跳转至

第14章 数据编排与质量


定时脚本“跑完不报错”,并不代表数据已经可以给 Agent 使用。上游分区可能没到、事实表可能为空、回填可能覆盖线上结果、质量规则也可能被临时关闭。DataAgent 最终看到的是数据产物,而不是任务退出码。

数据编排负责把资产按正确依赖生产出来,质量门禁负责判断这些资产何时可以被消费。 两者必须进入同一条发布链路,才能把脚本集合变成可运营的数据产品。

14.1 从调度脚本到数据产品:成功标准发生了什么变化

早期数据链路常用 cron 串起抽取、清洗和汇总。任务少时足够,但无法稳定回答:上游是否真的就绪、产出是否可信、失败后该重试还是回填、坏数据是否已经进入下游。

图14-1:编排不是孤立的调度器,质量也不是事后报表

图14-1:编排不是孤立的调度器,质量也不是事后报表。来源:本书自绘。Alt text:左侧"传统观念"中调度器与质量报表彼此分离,右侧"数据产品观念"中编排与质量门禁嵌在同一发布流程里,对比两种组织方式。

图14-2:脚本任务到数据产品的成功标准变化

图14-2:脚本任务到数据产品的成功标准变化。来源:本书自绘。Alt text:左列"脚本任务"成功标准是"跑完不报错",右列"数据产品"成功标准升级为契约满足、质量达标、可订阅、可追溯,对比成功定义的变化。

数据产品时代的成功标准至少包括:依赖满足、输入完整、质量达标、按 SLA 发布、失败可恢复、结果可追溯。

数据 DAG、资产依赖、业务工作流与 Agent Workflow 要分开

表14-1:数据 DAG、编排、质量门禁等概念的定义与区别。来源:本书整理。

概念 定义 与相邻概念的区别
数据 DAG 用有向无环图表达任务执行顺序 关注任务执行,不天然理解资产语义
资产依赖 表、视图、指标、特征的上下游关系 关注产物和影响范围
业务工作流 围绕审批、派单、补货等业务状态流转 不等于数据生产依赖
Agent Workflow Agent 调工具、检查结果、请求人工确认的路径 可以消费数据,但不替代数据编排
质量门禁 数据进入消费前的检查、阻断和降级 不只是监控报表

图14-3:边界清晰能降低系统耦合

图14-3:边界清晰能降低系统耦合。来源:本书自绘。Alt text:左侧职责混杂的任务相互交叉连线、耦合高,右侧编排、转换、质量、发布各司其职、连线清晰,对比边界清晰前后的耦合度。

数据 DAG 负责稳定生产,Agent 负责受控消费和解释。不要让 Agent 自由修复数据 DAG,也不要让调度器承担复杂业务审批。

14.2 触发与恢复:时间调度、数据集触发、回填不能混为一谈

常见触发方式包括时间、事件、数据集状态和人工回填。

表14-2:时间、事件、数据集触发与回填四种调度模型的优势与适用场景。来源:本书整理。

调度模型 触发条件 优势 代价 适用场景
时间调度 固定时间/周期 简单、可预测 上游未就绪时会空跑 日报、月报、稳定批处理
事件触发 文件或事件到达 响应快 要处理事件可靠性和去重 小时级增量、准实时
数据集触发 上游资产进入可用状态 更贴近数据依赖 需要资产状态平台 多团队共享数据产品
人工回填 指定历史区间重跑 可修复历史 容易污染当前结果 事故修复、口径重算

图14-4:调度触发模型与回填治理边界

图14-4:调度触发模型与回填治理边界。来源:本书自绘。Alt text:图中并列时间触发、事件触发、数据集触发三种模型,下方标出回填场景的特殊处理边界,说明常规调度与回填须分开治理。

回填尤其不能简单理解为“把历史日期再跑一次”。它需要影子分区、资源限额、幂等写入、质量复检、差异报告和下游通知,否则历史修复会直接改变正在被 Agent 使用的当前结果。

恢复动作也要分清:

  • 重试:处理网络、服务暂时不可用等瞬时故障;
  • 回填:补齐缺失的历史窗口;
  • 重算:处理代码、口径或业务逻辑错误。

图14-10:重试、回填与重算的恢复决策路径

图14-10:重试、回填与重算的恢复决策路径。来源:本书自绘。Alt text:决策树按"是瞬时失败、上游缺数据还是逻辑变更"分出重试、回填、重算三条恢复路径,帮助选择合适的恢复动作。

把所有失败都配置成自动重试,会掩盖 Schema、质量和权限类错误,而不是提高可靠性。

14.3 编排工具与质量工具:工具不能代替资产治理

Airflow、Dagster、Prefect、DolphinScheduler 的差异,主要在任务/资产建模、开发体验和组织协作方式。

表14-3:Airflow、Dagster、Prefect 等编排工具的优势、代价与适用场景。来源:本书整理。

方案 优势 代价 适用场景 本书建议
Airflow 生态成熟、调度稳定 资产语义需额外治理 传统批处理、已有数据平台 通用调度底座 + 资产/质量层
Dagster 资产建模强 需要接受资产开发范式 数据产品、强血缘 新建资产中心平台可重点评估
Prefect 开发体验轻、动态流程灵活 集中治理需补强 数据科学、应用自助流程 适合快速迭代
DolphinScheduler 可视化和任务类型丰富 资产语义和代码治理需补强 多团队协作、国产生态 适合重可视化组织

图14-5:工具选择应回到组织能力

图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:质量门禁的最小闭环

图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:把“写入临时分区”和“切换正式版本”拆开

图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:数据事故处理的止血与修复路径

图14-7:数据事故处理的止血与修复路径。来源:本书自绘。Alt text:事故流程分两段,先止血(下游降级、回退旧版本、告警),后修复(定位根因、回填重算、复盘),箭头表示先恢复可用性再追根因。

表14-7:上游未就绪、数据漂移等数据事故的检测方式与恢复策略。来源:本书整理。

失败模式 影响 检测 恢复策略
上游未就绪 空跑、旧数据 资产状态、行数、心跳 等待、重试、上一分区降级
字段漂移 转换失败或隐性错误 Schema/契约检查 阻断、变更评审
主键重复 指标翻倍、Join 膨胀 唯一性、对账 幂等重写、隔离重复
规则误报 无效阻断 告警确认率、历史分布 动态阈值、业务日历
规则漏报 坏数据进入 Agent 用户质疑、对账 复盘补规则
回填污染当前 指标前后不一致 发布审计、版本差异 影子回填后切换
告警无人处理 故障扩大 未确认时长 强制 Owner 和升级机制

14.6 质量状态必须进入 Agent 产品行为

质量状态不能只留在 Airflow/Dagster 页面。DataAgent 查询前需要知道资产是 publishedblocked、正在回填,还是只能使用上一版本。

可以把状态映射成四种产品行为:

  • 阻断:拒绝正式回答或报告;
  • 降级:使用上一版/受限数据,并明确截止时间;
  • 观察:允许分析,但提示软告警;
  • 正常:进入标准问数和报告链路。

用户不需要看到 DAG task id,但需要理解“为什么今天只能看到昨天数据”“哪些指标受影响”“什么时候可以恢复”。

用户质疑也应成为质量信号。DataAgent 的“这个数和看板不一样”“为什么区域口径错了”等反馈,要关联回资产、字段、指标和 Run,经 Owner 裁定后再沉淀为规则、元数据修正或评测样本。

图14-8:从事故复盘到规则沉淀的反馈链路

图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.