引言
数据质量(Data Quality)是数据工程中最容易被忽视、却最致命的环节。上游一张表的字段语义漂移、一次调度延迟、一个重复的订单主键,都会顺着血缘链向下游传导,最终让 BI 报表失真、模型特征退化、业务决策失误。与软件工程中的单元测试一样,数据质量监控的目标不是"消灭所有问题",而是"让问题在影响业务之前被尽早发现"。本文将围绕数据质量的六大维度、两大开源框架(Great Expectations 与 Soda Core)、质量门禁与异常检测四个层面,给出可直接在生产环境复用的监控体系。
数据质量是数据产品的 SLO,而不是可选项。没有监控的数据管道,本质上是一个尚未爆炸的定时炸弹。
一、数据质量维度体系
1.1 六大质量维度
业界通常用以下六个维度刻画数据集是否"可信"。它们并非互斥,而是一个数据集在不同消费场景下需要同时满足的多维约束。
| 维度 | 英文 | 典型问题 | 监控手段 |
|---|---|---|---|
| 完整性 | Completeness | 关键字段为空、整行缺失 | not_null、row_count 下限 |
| 唯一性 | Uniqueness | 主键重复、去重后数量异常 | unique、duplicate_count |
| 一致性 | Consistency | 币种、状态码取值越界 | in_set、regex、跨表对比 |
| 时效性 | Timeliness | 数据迟到、新鲜度超过 SLA | freshness、事件时间水位 |
| 准确性 | Accuracy | 均值漂移、分布突变、异常值 | mean/stdev 区间、异常检测 |
| 有效性 | Validity | 格式错误、类型不匹配 | schema 校验、valid 格式 |
1.2 质量分级与 SLA
并非所有表都需要 100% 质量,质量门槛应与其业务重要程度绑定。
# quality_tier.yaml
tiers:
- name: critical # 直接支撑财报 / 监管 / 在线推荐
completeness: 100.0
freshness_minutes: 5
action: BLOCK_AND_PAGE
- name: important # 支撑日常运营报表
completeness: 99.99
freshness_minutes: 30
action: BLOCK_AND_ALERT
- name: nice_to_have # 探索性分析
completeness: 95.0
freshness_minutes: 720
action: ALERT_ONLY
1.3 监控架构分层
一套完整的数据质量监控体系应当覆盖"表级-字段级-业务指标级"三个层次:每条数据管道内联 Checkpoint 逐层验证,运行结果写入质量元数据库(含历史基线),再由质量看板、告警与门禁判断统一消费。这样既能做到"越早发现越好",又能让每一次质量事件都有历史上下文可供分析。
二、Great Expectations 框架核心概念
2.1 Expectation / Suite / Checkpoint
Great Expectations(简称 GX)的核心心智模型只有三个概念:Expectation 是对数据"应当如何"的单条断言;Expectation Suite 是这些断言的集合(对应一张表);Checkpoint 是"在什么数据上、用哪个 Suite、在何时"执行验证的编排单元。
| 概念 | 英文 | 类比 | 生命周期 |
|---|---|---|---|
| 期望 | Expectation | 单条测试用例 | 声明一次 |
| 期望套件 | Expectation Suite | 一张表的测试套件 | 随表演进 |
| 验证 | Validation | 测试执行 | 每次运行 |
| 检查点 | Checkpoint | 测试流水线 | 定时/事件触发 |
| 数据文档 | Data Docs | 测试报告 | 每次运行生成 |
2.2 初始化 Data Context
GX 1.x 使用 gx 命名空间,先创建项目级 Data Context,再挂载数据源。
# init_gx.py
import great_expectations as gx
context = gx.get_context(mode="file") # 或 gx.get_context(project_config=...)
datasource = context.sources.add_or_update_sql_datasource(
name="orders_db",
connection_string="postgresql+psycopg2://user:pass@dw-host:5432/warehouse",
)
# 为 orders 表创建 asset
table_asset = datasource.add_table_asset(
name="orders_asset",
table_name="ods_orders",
schema_name="ods",
)
2.3 期望套件示例
下面用 GX 的 Expectation Suite 定义一张订单贴源表的核心质量契约。
# suite_orders.py
from great_expectations import gx
context = gx.get_context(mode="file")
suite = context.add_expectation_suite("ods_orders.suite")
# 主键唯一性:order_id 必须唯一
suite.expectation_context.add_expectation(
expect_column_values_to_be_unique(column="order_id")
)
# 完整性:核心字段不允许为空
suite.expectation_context.add_expectation(
expect_column_values_to_not_be_null(column="order_id")
)
suite.expectation_context.add_expectation(
expect_column_values_to_not_be_null(column="user_id", mostly=0.9999)
)
# 准确性/业务规则:金额必须为正
suite.expectation_context.add_expectation(
expect_column_values_to_be_between(column="amount", min_value=0, max_value=1_000_000)
)
# 有效性:状态码在合法枚举内
suite.expectation_context.add_expectation(
expect_column_values_to_be_in_set(
column="status",
value_set=["CREATED", "PAID", "SHIPPED", "CANCELLED", "REFUNDED"],
)
)
# 时效性兜底:当日数据量下限
suite.expectation_context.add_expectation(
expect_table_row_count_to_be_between(min_value=100_000, max_value=50_000_000)
)
2.4 Checkpoint 与验证报告
Checkpoint 将 Suite 与数据资产绑定,并在运行后产出机器可读的验证结果。
# run_checkpoint.py
import great_expectations as gx
context = gx.get_context(mode="file")
checkpoint = context.add_or_update_checkpoint(
name="ods_orders.daily_check",
validations=[
{
"batch_request": {
"datasource_name": "orders_db",
"data_asset_name": "orders_asset",
},
"expectation_suite_name": "ods_orders.suite",
}
],
result_format="COMPLETE",
)
result = context.run_checkpoint(checkpoint_name="ods_orders.daily_check")
# 结果判定:全部通过则 success=True
print("SUCCESS" if result.success else "FAILED", {
e["expectation_config"]["expectation_type"]: e["success"]
for e in result.list_validation_results()[0]["results"]
})
三、Soda Core 与 SodaCL 实战
3.1 Soda 的定位差异
Soda 与 GX 最大的差异在于:Soda 采用 声明式配置优先(SodaCL 在 YAML 中写检查),开箱即含 freshness、anomaly score 等监控型检查,更适合持续监控场景;GX 以 Python 代码为主,更适合在 CI 中做"单元测试式"门禁。两者可以并存:CI 中用 GX,运行时用 Soda。
3.2 配置数据源
# configuration.yml
data_source orders_dw:
type: postgres
host: dw-host
port: 5432
username: soda_monitor
password: ${SODA_PG_PASSWORD}
database: warehouse
schema: ods
3.3 SodaCL 检查文件
SodaCL 检查文件放在 checks/ 目录,一套 YAML 覆盖多类检查。
# checks/ods_orders.yml
checks for ods_orders:
# 时效性:数据必须每 5 分钟内更新
- freshness using column created_at < 5m
# 完整性:行数不能低于下限
- row_count > 100000
- missing_count(order_id) = 0
- missing_count(user_id) < 100
# 唯一性:重复数量为零
- duplicate_count(order_id) = 0
# 有效性:金额取值区间
- max(amount) < 1000000
- min(amount) >= 0
# Schema:不允许出现未声明的列
- schema:
warn:
when unknown columns: 0
3.4 运行 Soda Scan
#!/bin/bash
# run_soda.sh
export SODA_PG_PASSWORD=$(aws secretsmanager get-secret-value \
--secret-id dw/monitor --query SecretString --output text)
soda scan -d orders_dw -c configuration.yml checks/ods_orders.yml -s scan_results.json
# 退出码:0=通过,非 0=失败,可直接用于 CI 门禁
echo "Soda scan exit code: $?"
四、数据质量规则定义与代码
4.1 规则即代码的工程约束
质量规则应当纳入版本控制、走代码评审,并且与表结构一起演进。推荐在仓库中划分 suites/(GX 期望套件)、checks/(SodaCL YAML)、gate/(门禁脚本)与 ci/(CI 编排)四个目录;同结构多张表的规则可用 Jinja 模板参数化,避免复制粘贴导致的规则漂移。
4.2 规则元数据化
每一条规则都应当可寻址、可溯源、可灰度。下面用一个 JSON 描述规则的生命周期状态机。
{
"rule_id": "dq-rule-orders-0001",
"expectation_type": "expect_column_values_to_not_be_null",
"column": "order_id",
"tier": "critical",
"status": "active",
"owner": "order-squad",
"introduced_in": "dq-repo@v1.4.0",
"alert_channels": ["#data-oncall", "email:dw@example.com"],
"mttr_reference": "runbook://data-quality/null-order-id"
}
五、质量门禁进 CI/CD
5.1 门禁语义
质量门禁(Quality Gate)遵循 fail fast 与 fail loud 原则:下游模型发布前必须通过上游质量检查;一旦失败,阻断发布并触发告警。
| 门禁阶段 | 检查内容 | 失败动作 |
|---|---|---|
| PR 阶段 | 规则本身可编译、元数据完整 | 阻断合并 |
| 预发布 | 对抽样数据跑 Suite | 阻断部署 |
| 发布后 | 全量数据定时扫描 | 告警 + 自动回滚开关 |
| 运行期 | freshness / anomaly | 降级流量 + 通知 |
5.2 GitHub Actions 集成
在数据管道仓库中,将 GX 门禁作为 CI Job 的一部分执行。
# .github/workflows/dq-gate.yml
name: Data Quality Gate
on:
push:
branches: [main]
paths: ["models/**", "suites/**"]
schedule:
- cron: "0 * * * *"
jobs:
run-gx-checkpoint:
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@v4
- uses: actions/setup-python@v5
with:
python-version: "3.11"
- name: Install deps
run: pip install "great_expectations>=1.0" snowflake-connector-python
- name: Run validation
env:
SNOWFLAKE_PASSWORD: ${{ secrets.SNOWFLAKE_PASSWORD }}
run: |
python -m great_expectations checkpoint run ods_orders.daily_check
# 失败则非零退出码阻断 CI
- name: Upload data docs
if: always()
uses: actions/upload-artifact@v4
with:
name: data-docs
path: great_expectations/uncommitted/data_docs/
5.3 数据回填的版本化门禁
历史回填(backfill)同样是质量事故高发区:旧逻辑与新逻辑混跑、分区重复、口径漂移。回填任务必须携带版本号,并在完成后与存量数据做一致性 diff(行数与关键指标漂移阈值内才允许切换读路径),以此保证回填不会污染既有质量基线。
六、异常检测:基数与均值漂移
6.1 静态阈值 vs 动态基线
静态区间(between 0 and 1000000)只能拦截"粗野"错误,无法发现分布渐变。现代 DQ 平台普遍引入统计式异常检测:基于历史窗口构建基线,用 z-score / 分位数判断当前值是否显著偏离。
| 指标 | 检测目标 | 算法示例 |
|---|---|---|
| 基数 | 去重数骤降(上游去重逻辑被改动) | 环比变化率 |
| 均值 | 均值漂移(汇率/折扣逻辑变更) | 滚动 z-score |
| 标准差 | 波动放大(脏数据混入) | 双样本方差检验 |
| 分布 | 取值分布突变 | KS 检验 / 直方图散度 |
6.2 基数(Cardinality)异常检测
基数变化通常意味着数据语义的静默变更——例如订单去重键被改、上游加了新的 join 导致维度膨胀。
# cardinality_anomaly.py
from collections import deque
class CardinalityDetector:
"""基于历史基数的 EWMA 基线检测"""
def __init__(self, window: int = 28, z_threshold: float = 4.0):
self.history: deque[float] = deque(maxlen=window)
self.z_threshold = z_threshold
def observe(self, distinct_count: float) -> dict:
if len(self.history) < 10:
self.history.append(distinct_count)
return {"anomaly": False, "reason": "cold_start"}
# EWMA 均值与波动
alpha = 0.3
mu, var = self.history[0], 0.0
for x in self.history:
mu = alpha * x + (1 - alpha) * mu
var = sum((x - mu) ** 2 for x in self.history) / len(self.history)
sigma = var ** 0.5 or 1.0
z = abs(distinct_count - mu) / sigma
self.history.append(distinct_count)
return {
"anomaly": z > self.z_threshold,
"z_score": round(z, 2),
"baseline_mean": round(mu, 2),
"current": distinct_count,
}
detector = CardinalityDetector()
for d in [1_200_000, 1_210_000, 1_195_000, 1_205_000, 1_198_000, 980_000]:
print(d, detector.observe(d))
6.3 均值漂移与质量反馈
均值漂移往往先于报表投诉数小时出现,是质量监控中最有"预警价值"的信号。将均值检查与 GX 结合:
# mean_drift.py
import great_expectations as gx
context = gx.get_context(mode="file")
suite = context.add_expectation_suite("ods_orders.mean_drift")
# 基于近 28 天基线计算出的置信区间
suite.expectation_context.add_expectation(
expect_column_mean_to_be_between(column="amount", min_value=118.0, max_value=132.0)
)
suite.expectation_context.add_expectation(
expect_column_stdev_to_be_between(column="amount", min_value=50.0, max_value=70.0)
)
6.4 根因定位与通知
一旦触发异常,应将上下文(数据集、指标、观测值、基线、z-score、窗口)打包成结构化告警,附上可能原因与 runbook 链接,供 on-call 快速决策——告警质量决定了 MTTR。
七、质量报告与血缘联动
7.1 质量评分模型
将多维检查结果汇总为 0-100 的质量分,作为数据产品的"健康度"对外暴露。
| 组件 | 权重 | 说明 |
|---|---|---|
| 完整性 | 25% | 空值率越界扣分 |
| 唯一性 | 20% | 主键重复扣分 |
| 准确性 | 25% | 均值/基数异常扣分 |
| 时效性 | 15% | 新鲜度超标扣分 |
| Schema 稳定性 | 15% | 未声明列变更扣分 |
7.2 质量分计算
# quality_score.py
def compute_quality_score(checks: dict) -> dict:
weights = {
"completeness": 0.25,
"uniqueness": 0.20,
"accuracy": 0.25,
"timeliness": 0.15,
"schema": 0.15,
}
score = 0.0
detail = {}
for name, w in weights.items():
ok, total = checks[name]["passed"], checks[name]["total"]
rate = ok / total if total else 1.0
detail[name] = round(rate, 3)
score += w * rate
grade = "A" if score >= 0.98 else "B" if score >= 0.95 else "C" if score >= 0.9 else "D"
return {"score": round(score * 100, 1), "grade": grade, "detail": detail}
print(compute_quality_score({
"completeness": {"passed": 8, "total": 9},
"uniqueness": {"passed": 4, "total": 4},
"accuracy": {"passed": 6, "total": 7},
"timeliness": {"passed": 2, "total": 3},
"schema": {"passed": 1, "total": 1},
}))
7.3 血缘联动:从表到业务影响的追溯
质量结果应写入 OpenLineage 图谱,使"某张表质量下降"能自动映射到"受影响的看板与下游模型"。
# openlineage_emit.yaml
job:
namespace: dq
name: ods_orders.daily_check
facets:
qualityCheck:
producer: great_expectations/1.0
score: 87.5
grade: B
failed_expectations:
- expect_column_values_to_be_in_set.status
- freshness
run:
status: FAILED
failed_at: "2026-09-27T02:00:00Z"
outputs:
- namespace: postgres.warehouse
name: ods.ods_orders
facets:
dataQuality:
pass_rate: 0.875
7.4 质量看板
质量分与异常事件应进入统一看板(如 Grafana / Tableau),展示全局质量健康度、质量最差数据集 TOP10、异常事件时间线、新鲜度达标率与质量分趋势,供数据团队与业务方共同观测。
八、生产落地案例
8.1 案例:电商订单域 DQ 体系
某电商公司将订单域作为 DQ 先行试点,两周内将"数据问题平均发现时间"从 6 小时缩短到 12 分钟。
| 阶段 | 动作 | 效果 |
|---|---|---|
| 盘点 | 圈定 23 张核心表,标注 Tier | 明确优先级 |
| 建规则 | 每表 8-15 条 Expectation + SodaCL 检查 | 规则覆盖 90% 维度 |
| 入 CI | 发布管道串入 GX 门禁 | 预发布问题拦截率 70% |
| 上线监控 | Soda 定时扫描 + 异常检测 | MTTR 降至 40 分钟 |
| 血缘联动 | 质量事件写入 OpenLineage | 业务影响自动通知 |
8.2 案例关键配置
该体系的持续扫描通过 Kubernetes CronJob 每 15 分钟运行一次 Soda Scan(sodadata/soda-core 镜像,凭证由 Secret 注入),异常事件自动路由到告警平台。
九、常见问题与最佳实践
Q1: GX 与 Soda 应该选哪个?
二者并非替代关系。GX 更擅长"代码化、可版本化的规则 + 丰富的 Expectation 类型",适合放在 CI 中作为发布门禁;Soda 更擅长"持续扫描 + freshness + anomaly score",适合放在运行时做在线监控。成熟团队通常两者并用:GX 管"能不能上线",Soda 管"上线后是否仍然健康"。
Q2: 规则过多导致告警疲劳怎么办?
告警疲劳的根因是"没有分级"。把规则按 Tier 分级,critical 规则失败才触发页面告警(page),important 仅发 IM,nice_to_have 只写看板。同时为每条规则配备 runbook 链接,降低 on-call 的决策成本。若某规则连续 30 天零触发,应下线或降级。
Q3: 如何让业务方信任质量分?
质量分必须透明可解释:不仅给一个总分,还要给出"哪条规则失败、影响哪张表、何时开始、影响哪些下游看板"。将质量看板直接嵌入业务方常用的 BI 入口,让质量状态成为业务决策的一部分,而不是数据团队内部指标。
Q4: 历史数据回填要不要过 DQ?
要。回填是质量事故的高发区——旧逻辑与新逻辑混跑、分区重复、口径漂移。回填任务必须带版本号,回填完成后自动执行一次与存量数据的 diff 检查(行数、关键指标漂移阈值),通过后才允许切换读路径。
总结
| 维度 | 推荐工具 | 触发方式 | 失败响应 |
|---|---|---|---|
| 完整性/唯一性/有效性 | GX Expectation Suite | CI 门禁 + 定时 | 阻断/告警 |
| 时效性 | Soda freshness | 每 15 分钟扫描 | 按 Tier 分级 |
| 准确性/漂移 | Soda anomaly + z-score | 每小时基线比对 | 告警 + runbook |
| 全链路影响 | OpenLineage 血缘联动 | 事件驱动 | 自动通知受影响方 |
数据质量监控没有银弹,但有一条清晰路径:用六大维度框定范围,用 GX 管好发布门禁,用 Soda 守住运行健康,用异常检测捕捉渐变退化,最后用血缘把质量问题翻译成业务影响。 先让 20% 的核心表达到监控闭环,再逐步扩展,是性价比最高的落地方式。
参考与延伸阅读
- Great Expectations 官方文档:Expectation / Suite / Checkpoint 概念
- Soda Core 官方文档:SodaCL 检查语言与
soda scanCLI - OpenLineage 规范:质量事件与血缘图谱的元数据模型
- Data Quality Fundamentals(O’Reilly,Barr & Schell 著)中的维度框架
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。