数据可观测性:从管道监控到数据宕机的全方位保障

系统解析数据可观测性(Data Observability)的落地体系:管道 Job 状态/耗时/失败率监控、数据新鲜度与数据宕机概念、Volume/Quality/Schema 四类核心指标、Airflow/Dagster 工作流可观测性、日志与链路追踪、告警分级与 on-call 机制,以及与 dbt/ELT 管道的深度集成,附生产配置与案例。

引言

软件可观测性关注"服务是否健康",数据可观测性(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
AirflowDAG 图无(需插件)弱无
Dagster资产图✅ 资产物化强插件
PrefectFlow/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-callSEV 分级 + runbook + 轮值MTTA / MTTR
复盘改进数据宕机回顾重复事故率

数据可观测性的终极目标是让"数据宕机"像"服务宕机"一样被第一时间发现、被严肃对待、被系统化解决。它的落地不需要一次到位:从管道监控起步,逐步补上新鲜度、质量、Schema 与告警分级,最终形成"指标-日志-追踪-告警-复盘"的完整闭环,数据团队才能从救火队变成真正的可靠性工程师。


参考与延伸阅读

  • Barr & Schell. Data Quality Fundamentals(数据宕机与可观测性章节)
  • Monte Carlo Data 关于 Data Observability 的定义与指标框架
  • Airflow 官方文档:SLA、on_failure_callback 与事件回调
  • elementary 官方文档:dbt 测试结果与运行元数据采集

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「data-engineering」更多文章

  1. 特征存储(Feature Store)架构:从一致性到在线检索的完整实践
  2. 反向 ETL 与数据激活:让数据仓库的价值回到业务系统
  3. 湖仓一体架构:Iceberg、Delta Lake 与 Hudi 的统一数据底座