第14章 数据编排与质量¶
定时脚本“跑完不报错”,并不代表数据已经可以给 Agent 使用。上游分区可能没到、事实表可能为空、回填可能覆盖线上结果、质量规则也可能被临时关闭。DataAgent 最终看到的是数据产物,而不是任务退出码。
数据编排负责把资产按正确依赖生产出来,质量门禁负责判断这些资产何时可以被消费。 两者必须进入同一条发布链路,才能把脚本集合变成可运营的数据产品。
14.1 从调度脚本到数据产品:成功标准发生了什么变化¶
早期数据链路常用 cron 串起抽取、清洗和汇总。任务少时足够,但无法稳定回答:上游是否真的就绪、产出是否可信、失败后该重试还是回填、坏数据是否已经进入下游。
图14-1:编排不是孤立的调度器,质量也不是事后报表。来源:本书自绘。Alt text:左侧"传统观念"中调度器与质量报表彼此分离,右侧"数据产品观念"中编排与质量门禁嵌在同一发布流程里,对比两种组织方式。
图14-2:脚本任务到数据产品的成功标准变化。来源:本书自绘。Alt text:左列"脚本任务"成功标准是"跑完不报错",右列"数据产品"成功标准升级为契约满足、质量达标、可订阅、可追溯,对比成功定义的变化。
数据产品时代的成功标准至少包括:依赖满足、输入完整、质量达标、按 SLA 发布、失败可恢复、结果可追溯。
数据 DAG、资产依赖、业务工作流与 Agent Workflow 要分开¶
表14-1:数据 DAG、编排、质量门禁等概念的定义与区别。来源:本书整理。
| 概念 | 定义 | 与相邻概念的区别 |
|---|---|---|
| 数据 DAG | 用有向无环图表达任务执行顺序 | 关注任务执行,不天然理解资产语义 |
| 资产依赖 | 表、视图、指标、特征的上下游关系 | 关注产物和影响范围 |
| 业务工作流 | 围绕审批、派单、补货等业务状态流转 | 不等于数据生产依赖 |
| Agent Workflow | Agent 调工具、检查结果、请求人工确认的路径 | 可以消费数据,但不替代数据编排 |
| 质量门禁 | 数据进入消费前的检查、阻断和降级 | 不只是监控报表 |
图14-3:边界清晰能降低系统耦合。来源:本书自绘。Alt text:左侧职责混杂的任务相互交叉连线、耦合高,右侧编排、转换、质量、发布各司其职、连线清晰,对比边界清晰前后的耦合度。
数据 DAG 负责稳定生产,Agent 负责受控消费和解释。不要让 Agent 自由修复数据 DAG,也不要让调度器承担复杂业务审批。
14.2 触发与恢复:时间调度、数据集触发、回填不能混为一谈¶
常见触发方式包括时间、事件、数据集状态和人工回填。
表14-2:时间、事件、数据集触发与回填四种调度模型的优势与适用场景。来源:本书整理。
| 调度模型 | 触发条件 | 优势 | 代价 | 适用场景 |
|---|---|---|---|---|
| 时间调度 | 固定时间/周期 | 简单、可预测 | 上游未就绪时会空跑 | 日报、月报、稳定批处理 |
| 事件触发 | 文件或事件到达 | 响应快 | 要处理事件可靠性和去重 | 小时级增量、准实时 |
| 数据集触发 | 上游资产进入可用状态 | 更贴近数据依赖 | 需要资产状态平台 | 多团队共享数据产品 |
| 人工回填 | 指定历史区间重跑 | 可修复历史 | 容易污染当前结果 | 事故修复、口径重算 |
图14-4:调度触发模型与回填治理边界。来源:本书自绘。Alt text:图中并列时间触发、事件触发、数据集触发三种模型,下方标出回填场景的特殊处理边界,说明常规调度与回填须分开治理。
回填尤其不能简单理解为“把历史日期再跑一次”。它需要影子分区、资源限额、幂等写入、质量复检、差异报告和下游通知,否则历史修复会直接改变正在被 Agent 使用的当前结果。
恢复动作也要分清:
- 重试:处理网络、服务暂时不可用等瞬时故障;
- 回填:补齐缺失的历史窗口;
- 重算:处理代码、口径或业务逻辑错误。
图14-10:重试、回填与重算的恢复决策路径。来源:本书自绘。Alt text:决策树按"是瞬时失败、上游缺数据还是逻辑变更"分出重试、回填、重算三条恢复路径,帮助选择合适的恢复动作。
把所有失败都配置成自动重试,会掩盖 Schema、质量和权限类错误,而不是提高可靠性。
14.3 编排工具与质量工具:工具不能代替资产治理¶
Airflow、Dagster、Prefect、DolphinScheduler 的差异,主要在任务/资产建模、开发体验和组织协作方式。
表14-3:Airflow、Dagster、Prefect 等编排工具的优势、代价与适用场景。来源:本书整理。
| 方案 | 优势 | 代价 | 适用场景 | 本书建议 |
|---|---|---|---|---|
| Airflow | 生态成熟、调度稳定 | 资产语义需额外治理 | 传统批处理、已有数据平台 | 通用调度底座 + 资产/质量层 |
| Dagster | 资产建模强 | 需要接受资产开发范式 | 数据产品、强血缘 | 新建资产中心平台可重点评估 |
| Prefect | 开发体验轻、动态流程灵活 | 集中治理需补强 | 数据科学、应用自助流程 | 适合快速迭代 |
| DolphinScheduler | 可视化和任务类型丰富 | 资产语义和代码治理需补强 | 多团队协作、国产生态 | 适合重可视化组织 |
图14-5:工具选择应回到组织能力。来源:本书自绘。Alt text:以"团队工程能力"和"资产治理需求"为轴的矩阵,把 Airflow、Dagster、Prefect 等工具落入不同象限,强调选型回到组织实际能力。
编排工具告诉你“任务有没有运行”,并不能自动证明数据正确。dbt tests 更贴近 SQL 模型;Great Expectations 和 Soda 更适合通用质量与持续巡检。
表14-4:dbt tests、Great Expectations 等转换与测试工具的优势与适用场景。来源:本书整理。
| 方案 | 优势 | 代价 | 适用场景 | 本书建议 |
|---|---|---|---|---|
| dbt tests | 与 SQL 模型贴近 | 复杂跨系统规则有限 | 事实表、维表、指标模型 | 模型内置测试起点 |
| Great Expectations | 规则丰富、文档化强 | 集成和规则维护有成本 | 跨源质量、契约验证 | 核心资产重点使用 |
| Soda | 配置简洁、适合持续监控 | 深度定制依赖集成 | 日常巡检、轻量门禁 | 标准化监控 |
| 自定义 SQL | 最灵活 | 易碎片化 | 特殊业务规则 | 允许存在,但统一登记 Owner |
14.4 质量门禁:规则必须绑定动作,而不是只产生红灯¶
数据质量至少覆盖完整性、唯一性、准确性、及时性、一致性和有效性。
表14-5:完整性、唯一性、准确性等数据质量维度的关注问题与示例规则。来源:本书整理。
| 质量维度 | 关注问题 | 示例规则 | 失败后的典型动作 |
|---|---|---|---|
| 完整性 | 是否缺失 | 行数处于合理历史范围 | 阻断或等待上游补齐 |
| 唯一性 | 主键是否重复 | order_id 分区内唯一 |
阻断并定位重复来源 |
| 准确性 | 数值是否符合业务事实 | 金额非负、时长不倒挂 | 隔离异常记录 |
| 及时性 | 是否在 SLA 前产出 | 08:00 前完成核心指标 | 告警、降级旧版本 |
| 一致性 | 跨表/系统是否对齐 | 订单和支付金额差异受控 | 暂停关键回答、触发对账 |
| 有效性 | 字段是否符合域约束 | 状态枚举合法 | 拒收或隔离坏数据 |
规则还要区分硬门禁、软告警和旁路隔离。
表14-6:硬门禁与软告警两种质量拦截策略的取舍。来源:本书整理。
| 方案 | 优势 | 代价 | 适用场景 | 本书建议 |
|---|---|---|---|---|
| 硬门禁 | 防止坏数据进入核心链路 | 可能造成不可用 | 主键重复、权限违规、金额非法 | 高风险资产 |
| 软告警 | 保持可用 | Agent 可能引用有风险结果 | 分布波动、轻微延迟 | 响应带质量状态 |
| 旁路隔离 | 主链路继续运行 | 需后续修复和对账 | 少量异常记录 | 隔离比例必须告警 |
图14-6:质量门禁的最小闭环。来源:本书自绘。Alt text:环形流程,产出数据、运行质量校验、通过则发布、不通过则阻断并告警,箭头表示失败样本回流修正规则,构成质量门禁闭环。
没有 Owner、严重级别和失败动作的质量规则,只是在制造告警噪声。
质量状态最好通过统一事件对外传播:
{
"asset_id": "ads.fulfillment_delay_daily",
"run_id": "run_20260611_020000",
"partition": "dt=2026-06-10",
"status": "blocked",
"severity": "critical",
"checks": [
{
"name": "order_id_unique",
"actual": "duplicate_count = 184",
"action": "block_publish"
}
],
"owner": "fulfillment-data-team",
"lineage": {
"downstream_consumers": ["DataAgent", "operations_dashboard"]
}
}
14.5 先写影子版本,再切正式版本:发布和事故恢复要连起来¶
质量检查应该发生在正式资产切换之前,而不是发布后做报表。
图14-9:把“写入临时分区”和“切换正式版本”拆开。来源:本书自绘。Alt text:发布分两步,先写入临时分区并校验,再原子切换正式版本指针,箭头表示校验通过才切换,避免半成品数据直接对外可见。
asset:
id: ads.fulfillment_delay_daily
owner: fulfillment-data-team
schedule: "0 6 * * *"
partition_key: dt
publish_mode: versioned_partition
dependencies:
- asset_id: dwd.orders_daily
freshness: 2h
quality_gates:
hard:
- name: order_id_unique
on_failure: block_publish
soft:
- name: row_count_anomaly
on_failure: publish_with_warning
fallback:
strategy: keep_previous_partition
max_age: 2d
发布伪代码只做一件关键事:坏数据不覆盖当前正式版本。
def publish_asset(run):
write_temp_partition(run.asset_id, run.partition, run.output)
quality_result = run_quality_checks(run.asset_id, run.partition)
if quality_result.has_blocking_failure:
mark_asset_state(run.asset_id, run.partition, "blocked", quality_result)
keep_previous_version(run.asset_id)
notify_owner(run.asset_id, quality_result)
return "blocked"
version = commit_versioned_partition(run.asset_id, run.partition)
mark_asset_state(run.asset_id, run.partition, "published", quality_result)
refresh_downstream_cache(run.asset_id, version)
return "published"
数据事故处理也应先止血、再修复:先冻结问题资产、降级到上一版本、阻断正式报告;再做血缘定位、回填、复检和复盘。
图14-7:数据事故处理的止血与修复路径。来源:本书自绘。Alt text:事故流程分两段,先止血(下游降级、回退旧版本、告警),后修复(定位根因、回填重算、复盘),箭头表示先恢复可用性再追根因。
表14-7:上游未就绪、数据漂移等数据事故的检测方式与恢复策略。来源:本书整理。
| 失败模式 | 影响 | 检测 | 恢复策略 |
|---|---|---|---|
| 上游未就绪 | 空跑、旧数据 | 资产状态、行数、心跳 | 等待、重试、上一分区降级 |
| 字段漂移 | 转换失败或隐性错误 | Schema/契约检查 | 阻断、变更评审 |
| 主键重复 | 指标翻倍、Join 膨胀 | 唯一性、对账 | 幂等重写、隔离重复 |
| 规则误报 | 无效阻断 | 告警确认率、历史分布 | 动态阈值、业务日历 |
| 规则漏报 | 坏数据进入 Agent | 用户质疑、对账 | 复盘补规则 |
| 回填污染当前 | 指标前后不一致 | 发布审计、版本差异 | 影子回填后切换 |
| 告警无人处理 | 故障扩大 | 未确认时长 | 强制 Owner 和升级机制 |
14.6 质量状态必须进入 Agent 产品行为¶
质量状态不能只留在 Airflow/Dagster 页面。DataAgent 查询前需要知道资产是 published、blocked、正在回填,还是只能使用上一版本。
可以把状态映射成四种产品行为:
- 阻断:拒绝正式回答或报告;
- 降级:使用上一版/受限数据,并明确截止时间;
- 观察:允许分析,但提示软告警;
- 正常:进入标准问数和报告链路。
用户不需要看到 DAG task id,但需要理解“为什么今天只能看到昨天数据”“哪些指标受影响”“什么时候可以恢复”。
用户质疑也应成为质量信号。DataAgent 的“这个数和看板不一样”“为什么区域口径错了”等反馈,要关联回资产、字段、指标和 Run,经 Owner 裁定后再沉淀为规则、元数据修正或评测样本。
图14-8:从事故复盘到规则沉淀的反馈链路。来源:本书自绘。Alt text:链路从一次数据事故出发,经复盘提炼出新的质量规则,沉淀进门禁,下次同类问题被提前拦截,箭头构成持续改进的反馈环。
在数据不可信时停止回答,是企业级 Agent 的保护能力,不是产品失败。
规则本身也要版本化。阈值、业务日历、空值策略和异常检测方法发生变化,会改变资产是否可用。高影响规则变更应回放核心问题,避免把“规则变松/变严”误解释成业务数据变化。
本章小结¶
数据编排的目标不是把脚本定时跑起来,而是让数据产品按依赖、质量、版本和 SLA 可靠发布。任务成功不等于资产可用,资产可用也不等于适合正式 Agent 回答。
DAG、资产依赖、业务工作流和 Agent Workflow 应分层;质量规则必须绑定 Owner、严重级别和失败动作;重试、回填、重算应使用不同恢复路径。
DataAgent 应只消费通过门禁或明确降级的数据资产,并把新鲜度和质量状态带到最终回答中。
参考文献¶
Apache Airflow. (n.d.). Documentation.
Dagster. (n.d.). Documentation.
Prefect. (n.d.). Documentation.
Great Expectations. (n.d.). Documentation.
Soda. (n.d.). Documentation.









