引言
软件可观测性关注"服务是否健康",数据可观测性(Data Observability)则关注"数据是否可信"。数据管道运行成功并不等于数据正确:调度成功但分区延迟、行数骤降、字段类型悄悄变化,这些都不会触发"500 错误",却会悄然污染下游报表与模型。数据可观测性把软件工程中成熟的可观测性理念——指标、日志、追踪、告警、SLO——迁移到数据领域,让数据团队从"被动救火"走向"提前预警"。
数据宕机(Data Downtime)指数据不可用、不准确或不可信的时间段。它的成本与服务器宕机同样真实,只是更难被感知。
一、数据管道监控
1.1 管道级四类指标
管道监控关注"作业本身",按四类指标刻画健康度。
| 指标类 | 示例 | 采集方式 |
|---|---|---|
| Job 状态 | 成功/失败/重试次数 | 调度器事件 |
| 耗时 | 全链路耗时、各阶段耗时 | 运行时间戳 |
| 失败率 | 日失败率、失败间隔 | 历史统计 |
| 资源 | 并发、内存、Shuffle 量 | 引擎指标 |
1.2 Airflow DAG 级监控
在 Airflow 中为 DAG 配置 SLA、失败告警与运行明细,是管道可观测性的第一层。
# dag_orders_etl.py
from datetime import datetime, timedelta
from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.sensors.external_task import ExternalTaskSensor
with DAG(
dag_id="orders_etl",
schedule="0 */2 * * *",
start_date=datetime(2026, 9, 1),
catchup=False,
default_args={
"retries": 2,
"retry_delay": timedelta(minutes=5),
# 失败与 SLA 未达标时推送结构化告警
"on_failure_callback": notify_oncall,
"sla": timedelta(hours=2),
"email_on_retry": False,
},
max_active_runs=1,
) as dag:
wait_upstream = ExternalTaskSensor(
task_id="wait_orders_raw",
external_dag_id="orders_raw_ingest",
external_task_id="done",
timeout=3600,
soft_fail=True,
)
run_merge = PythonOperator(task_id="merge_orders", python_callable=merge_orders)
build_metric = PythonOperator(task_id="build_daily_revenue", python_callable=build_daily_revenue)
wait_upstream >> run_merge >> build_metric
1.3 失败率与趋势
用 Prometheus 记录每次作业的运行结果,形成失败率时间序列。
# pipeline_metrics.py
from prometheus_client import Counter, Histogram, Gauge
RUN_TOTAL = Counter("pipeline_run_total", "Pipeline runs", ["dag", "task"])
RUN_FAILED = Counter("pipeline_run_failed", "Pipeline failures", ["dag", "task"])
RUN_DURATION = Histogram("pipeline_run_duration_seconds", "Pipeline duration", ["dag"])
FRESHNESS = Gauge("data_freshness_seconds", "Data freshness lag", ["table"])
def record_run(dag: str, task: str, ok: bool, duration: float):
RUN_TOTAL.labels(dag, task).inc()
if not ok:
RUN_FAILED.labels(dag, task).inc()
RUN_DURATION.labels(dag).observe(duration)
二、数据新鲜度与数据宕机
2.1 新鲜度(Freshness)语义
新鲜度衡量"数据离上次更新的时间差"。它直接反映管道是否按时产出,是数据 SLO 中最常用的指标。
| 概念 | 定义 | 示例 |
|---|---|---|
| Freshness | 表最近更新距现在的时间 | 5 分钟 |
| 新鲜度 SLO | 允许的最大滞后 | < 30m |
| 迟到数据 | 事件时间远早于处理时间 | 补录的历史分区 |
2.2 数据宕机与业务影响
数据宕机(Data Downtime)比服务宕机更难发现,因为它通常没有显式报错。
服务宕机(显式): 500 错误 → 立即感知 → 快速响应
数据宕机(隐式): 行数骤降 → 无人知晓 → 报表失真 3 天
2.3 新鲜度检查落地
用 Soda 的 freshness 检查做表级新鲜度监控,超时即触发 SLO 违约。
# checks/freshness.yml
checks for ods_orders:
# 核心表新鲜度 SLO:滞后不超过 30 分钟
- freshness using created_at < 30m:
name: orders_freshness_slo
# 双保险:基于最大主键时间判断
- max(created_at) > ${{ yesterday_ts }}:
name: orders_max_ts_recent
2.4 新鲜度驱动的自动补偿
当新鲜度违约且检测到上游未产出时,可触发自动重跑或降级。
# freshness_controller.py
def handle_freshness_breach(table: str, lag_minutes: float):
if lag_minutes > 30:
# 触发上游重跑
trigger_dag("orders_raw_ingest")
# 通知订阅该表的团队
notify_subscribers(table, f"freshness {lag_minutes}m exceeded SLO")
if lag_minutes > 120:
# 超过 2 小时:标记数据宕机,接入 on-call
escalate_to_oncall("SEV2", table)
三、核心指标体系
3.1 四类数据健康指标
数据可观测性常以"四类指标 + 新鲜度"构建健康画像。
| 指标 | 关注点 | 典型异常 |
|---|---|---|
| Volume | 行数/字节数 | 行数骤降、突增 |
| Quality | 空值率、重复率、取值分布 | 空值激增 |
| Schema | 字段增删改、类型变化 | 上游加列导致下游断裂 |
| Freshness | 更新滞后 | 分区晚点 |
3.2 Volume 异常检测
对行数做基线比对,是最早能发现管道静默失败的手段。
# volume_anomaly.py
from statistics import mean, pstdev
def check_volume(history: list, current: int, z_threshold: float = 3.5) -> dict:
mu = mean(history)
sigma = pstdev(history) or 1.0
z = abs(current - mu) / sigma
return {
"z_score": round(z, 2),
"baseline": round(mu, 0),
"current": current,
"anomaly": z > z_threshold,
}
# 近 14 天订单量基线
history = [1_201_000, 1_210_000, 1_195_000, 1_208_000, 1_202_000,
1_196_000, 1_210_500, 1_199_000, 1_204_000, 1_207_000,
1_201_500, 1_198_500, 1_205_000, 1_203_000]
print(check_volume(history, current=1_205_000)) # 正常
print(check_volume(history, current=930_000)) # 异常
3.3 Schema 变更检测
Schema 漂移是最常见的数据事故来源,必须在消费前拦截。
-- 对比两日的表结构,检测新增/删除列
SELECT column_name, data_type,
COUNT(*) FILTER (WHERE snapshot_date = CURRENT_DATE) AS today_cnt,
COUNT(*) FILTER (WHERE snapshot_date = CURRENT_DATE - 1) AS yesterday_cnt
FROM information_schema.columns
CROSS JOIN generate_series(CURRENT_DATE - 1, CURRENT_DATE, INTERVAL '1 day') AS snapshot_date
WHERE table_name = 'ods_orders'
GROUP BY column_name, data_type
HAVING today_cnt <> yesterday_cnt;
四、工作流可观测性
4.1 编排器能力对比
Airflow、Dagster、Prefect 对"可观测性"的支持深度不同。
| 编排器 | 运行视图 | 资产视角 | 事件驱动 | 内置 DQ |
|---|---|---|---|---|
| Airflow | DAG 图 | 无(需插件) | 弱 | 无 |
| Dagster | 资产图 | ✅ 资产物化 | 强 | 插件 |
| Prefect | Flow/Deployment | 部分 | 强 | 无 |
4.2 Dagster 的资产可观测性
Dagster 以"资产(Asset)“为中心,天然把"数据产物的产出"作为可观测对象。
# assets_observability.py
from dagster import asset, AssetExecutionContext, MetadataValue
@asset(compute_kind="spark")
def daily_order_revenue(context: AssetExecutionContext):
df = run_spark_job("daily_revenue.sql")
# 把行数、质量分、血缘作为资产元数据暴露
context.log_event(
MetadataValue.int(len(df)),
MetadataValue.text("98.2"),
MetadataValue.url("lake.dws.daily_order_revenue"),
)
return df
# 资产未按时物化即触发告警
from dagster import asset_sensor, RunRequest
@asset_sensor(asset_key="daily_order_revenue")
def missing_revenue_sensor(context, event):
return RunRequest(run_key=f"rerun-{event.asset_key}")
4.3 调度级告警:SLA 未达成
Airflow 的 SLA 机制能把"整个 DAG 链路的端到端超时"转化为告警,而不是只看单个任务。
# sla_callback.py
def sla_miss_callback(dag, task_list, blocking_task_list, slas, *args):
"""DAG 级 SLA 违约回调:推送结构化事件到监控平台"""
event = {
"type": "sla_miss",
"dag": dag.dag_id,
"tasks": [t.task_id for t in task_list],
"blocking": [t.task_id for t in blocking_task_list],
"sla_date": slas[0].execution_time.isoformat(),
}
emit_metric("data_pipeline_sla_miss", 1, {"dag": dag.dag_id})
notify_channel("#data-oncall", event)
五、日志与链路追踪
5.1 数据链路的追踪
数据问题往往跨越多个作业,需要把"一次端到端数据流动"串成一条 trace。推荐用 OpenTelemetry 打点,贯通调度器与计算引擎。
# data_trace.py
from opentelemetry import trace
tracer = trace.get_tracer("data-pipeline")
with tracer.start_as_current_span("orders_freshness_check") as span:
span.set_attribute("table", "ods_orders")
span.set_attribute("expected_lag_min", 30)
lag = measure_freshness("ods_orders")
span.set_attribute("actual_lag_min", lag)
if lag > 30:
span.set_attribute("slo_breach", True)
5.2 结构化运行日志
每个作业应输出结构化日志(JSON),供按运行 ID 聚合检索。
{
"ts": "2026-09-27T02:00:00Z",
"dag": "orders_etl",
"run_id": "scheduled__2026-09-27T02:00:00+00:00",
"task": "merge_orders",
"level": "WARN",
"message": "partition dt=2026-09-26 late by 12m",
"metrics": {"rows_processed": 1200000, "late_partitions": 2}
}
5.3 日志检索与看板
将日志与指标接入统一平台(ELK / Loki + Grafana),实现"指标发现异常 → 日志定位根因"的闭环。
Grafana Data Observability
├── Pipeline: 成功率 / 耗时 P95 / 重试率
├── Data: Freshness / Volume / Quality score
├── Traces: 端到端数据流动时间线
└── Logs: 按 dag/run_id 检索结构化日志
六、告警分级与 on-call
6.1 告警分级体系
告警必须分级,否则 on-call 会被淹没在噪音中。
| 级别 | 含义 | 响应时限 | 通道 | 示例 |
|---|---|---|---|---|
| SEV1 | 核心数据不可用 | 15 分钟 | 电话 + 页面 | 财报表宕机 |
| SEV2 | 重要数据异常 | 1 小时 | 页面 + IM | 核心表新鲜度超 2h |
| SEV3 | 一般数据问题 | 当天 | IM | 非核心表空值率升高 |
| INFO | 仅记录 | - | 看板 | 规则变更通知 |
6.2 告警路由配置
把告警路由到正确的团队与 runbook,是降低 MTTR 的关键。
# alert_routing.yaml
routing:
- severity: SEV1
page: true
escalation:
- team: data-platform-oncall
delay_min: 0
- team: data-platform-leads
delay_min: 15
- severity: SEV2
page: false
channel: "#data-oncall"
notify_owner: true
- severity: SEV3
channel: "#data-quality"
business_hours_only: true
runbook_map:
freshness: "runbook://data/freshness-breach"
volume: "runbook://data/volume-drop"
schema: "runbook://data/schema-change"
6.3 on-call 保障机制
# oncall_handoff.py
def build_handoff(severity: str, incident: dict) -> dict:
return {
"title": f"[{severity}] {incident['table']} data incident",
"summary": incident["summary"],
"runbook": incident["runbook"],
"affected": incident.get("downstream", []),
"impact": incident.get("business_impact", "unknown"),
"handoff_notes": (
f"started_at={incident['started_at']}; "
f"baseline={incident['baseline']}; current={incident['current']}"
),
}
七、与 ELT 管道的集成
7.1 集成模式
现代 ELT(dbt + 数仓)的可观测性,本质上是在 ELT 的每个环节植入"检查点”。
| ELT 环节 | 可观测性动作 | 工具 |
|---|---|---|
| 数据同步 | 行数/新鲜度检查 | Airbyte/Fivetran 元数据 |
| 模型构建 | 每模型测试 + 运行明细 | dbt tests + elementary |
| 物化后 | 关键指标断言 | Great Expectations |
| 消费前 | 数据宕机检测 | Soda/Monte Carlo |
7.2 dbt + elementary 集成
elementary 把 dbt 的测试结果与运行元数据转化为可观测信号,并驱动告警。
#!/bin/bash
# run_dbt_with_observability.sh
set -euo pipefail
# 构建并测试
dbt build --select daily_order_revenue+ --target prod
# elementary 采集测试与运行元数据
edr monitor-report
# 若关键测试失败,退出非零码触发告警
dbt test --select tag:critical --target prod
# dbt_project.yml 中的测试声明
models:
daily_order_revenue:
columns:
- name: dt
tests: [not_null, unique]
- name: revenue
tests:
- not_null
- elementary.column_anomalies:
anomaly_direction: both
sensitivity: 1.5
7.3 管道级联动:质量阻断
当 ELT 管道某环节失败时,自动阻断下游并通知受影响方。
# elt_gate.py
def elt_quality_gate(model: str, check_result: dict):
if not check_result["success"]:
# 1. 阻断下游模型运行
block_downstream(model)
# 2. 标记数据宕机开始时间
mark_data_downtime(model, started_at=datetime.utcnow())
# 3. 通知订阅团队
notify_subscribers(model, check_result)
八、生产案例
8.1 案例:某电商的端到端数据可观测性
某电商公司从"报表不准才知道"升级到"分钟级预警",核心数据 SLA 达标率从 92% 提升到 99.5%。
| 阶段 | 动作 | 结果 |
|---|---|---|
| 指标 | 为 40 张核心表建立 freshness/volume/quality 基线 | 覆盖全部核心链路 |
| 集成 | Airflow SLA + dbt tests + elementary | 问题发现时间降至 10 分钟 |
| 告警 | SEV 分级 + runbook + 轮值 | MTTR 从 6h 降至 40min |
| 复盘 | 每月数据宕机回顾,沉淀 runbook | 重复事故减少 70% |
8.2 数据宕机成本度量
把数据宕机翻译成可量化的成本,才能支撑投入。
Data Downtime Cost Model
├── 影响时长 × 受影响用户数 × 单用户损失
├── 报表失真天数 × 错误决策概率
├── 数据团队处置工时 × 人力成本
└── 合规风险敞口(监管报表延迟)
九、常见问题与最佳实践
Q1: 数据可观测性与数据质量监控的区别?
数据质量监控回答"这张表是否符合规则",偏"检查";数据可观测性回答"整条数据链是否健康、为何异常、影响多大",偏"诊断与预警"。可观测性 = 指标 + 日志 + 追踪 + 告警 + SLO 的组合,质量监控只是其中一个输入。
Q2: 告警太多怎么办?
先分级、再收敛:SEV1/2 才允许页面与电话,SEV3 只发 IM;为每条规则设置 runbook 与"静默窗口";对连续 30 天零触发的规则自动降级。告警的目的是让 on-call 少收到"无意义警报",而不是更多。
Q3: 如何从"监控管道"演进到"监控数据"?
演进路径清晰:第一阶段管 Job(状态/耗时/失败率),第二阶段管 Freshness 与 Volume,第三阶段管 Quality 与 Schema,第四阶段做异常检测与业务影响分析。每一步都能独立产出价值,不必一步到位。
Q4: 需要自研还是买商业方案?
早期(<50 张核心表)完全可以用开源组合:Airflow + Prometheus + dbt/elementary + Soda + Grafana 就能跑通。当表规模上千、跨团队需要统一 SLO 与自动化根因分析时,再评估商业方案(如 Monte Carlo、Elementary 企业版),避免过早引入重平台。
总结
| 能力 | 工具/方法 | 关键指标 |
|---|---|---|
| 管道监控 | Airflow SLA + Prometheus | 失败率、耗时 |
| 数据健康 | Soda/elementary + 基线 | freshness、volume、quality |
| Schema 守护 | dbt tests + schema 检查 | schema 变更事件 |
| 追踪诊断 | OpenTelemetry + 结构化日志 | 端到端 trace |
| 告警与 on-call | SEV 分级 + runbook + 轮值 | MTTA / MTTR |
| 复盘改进 | 数据宕机回顾 | 重复事故率 |
数据可观测性的终极目标是让"数据宕机"像"服务宕机"一样被第一时间发现、被严肃对待、被系统化解决。它的落地不需要一次到位:从管道监控起步,逐步补上新鲜度、质量、Schema 与告警分级,最终形成"指标-日志-追踪-告警-复盘"的完整闭环,数据团队才能从救火队变成真正的可靠性工程师。
参考与延伸阅读
- Barr & Schell. Data Quality Fundamentals(数据宕机与可观测性章节)
- Monte Carlo Data 关于 Data Observability 的定义与指标框架
- Airflow 官方文档:SLA、on_failure_callback 与事件回调
- elementary 官方文档:dbt 测试结果与运行元数据采集
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。