🏠 总目录📚 本教程 03 · 数据契约
📑 本页目录(点开跳转)

03 · 数据契约

62 分钟 | ⭐ 把口头约定变成会失败的测试


🎯 一句话

「上游改字段没人告诉你」不是意外,是默认状态。 数据契约做的事只有一件: 把「我们说好了不改」变成「你一改,你的 CI 就红」 —— 因为只要代价不落在改的人身上,通知就永远不会发生。


🧨 一、先理解它为什么必然发生

不是上游的人不负责任。是三个结构性原因:

   ① 你不在他的 KPI 里
      上游团队的目标是"业务功能上线",不是"下游模型别挂"
      他甚至【不知道】有你这个消费者

   ② 下游是不可见的
      一张表可能有 40 个消费者:报表、风控、推荐、财务对账、外部接口
      上游打开表的元信息,看到的是:什么都没有

   ③ 改动成本严重不对称  ⭐⭐
      上游改一行 SQL / 一个枚举值:    5 分钟
      下游发现异常 → 定位 → 修复 → 回溯:3 天
      ——「省下 5 分钟」的收益归他,「花掉 3 天」的成本归你

🔑 所有数据治理方案,本质上都在解第 ③ 条:把成本挪回去。 契约不是文档,是一个把下游的痛,前置成上游的 CI 失败的装置。 一份没有强制执行的契约,等于一份措辞更正式的口头约定。


📄 二、一份契约该写什么:七个部分

# 部分 必须包含 少了会怎样
1 Schema 字段名、类型、可空、枚举取值全集、主键 类型变化静默通过
2 语义 ⭐⭐ 口径定义、单位、时区、边界情况、"不算什么" 第 2 章的曝光事故
3 质量 SLO 新鲜度、完整性、唯一性、取值域、行数区间 脏数据静默流入
4 所有权 负责人 + 值班群 + 兜底人 出事找不到人,平均多花 6 小时
5 变更流程 提前几天通知、双写窗口、灰度、回滚方案 改完你才知道
6 兼容性等级 哪些算破坏性变更、版本号规则 「我只是加了个枚举值」⚠️
7 消费者清单 谁在用、用在哪、联系人 上游根本不知道要通知谁

第 2 部分里「不算什么」是最容易漏、也最值钱的一行: 定义「曝光」时,写清楚「预加载不算、后台自动刷新不算、快速滑过不算」, 比写十行「曝光是指用户看到内容」有用得多。 口径的边界永远比口径的主体更容易出事。

一份最小可用契约长这样

dataset: dwd_feed_impression_di
owner: feed-client@team          # 4 所有权
oncall: "#feed-data-oncall"
consumers:                        # 7 消费者清单 ⭐
  - {team: rank-algo, usage: "CTR 模型负样本", contact: "@zhang"}
  - {team: bi,        usage: "日报曝光量",    contact: "@li"}

semantics:                        # 2 语义 ⭐⭐
  impression:
    definition: "item 进入视口且停留 >= 1000ms"
    not_counted: ["预加载", "后台刷新", "滑动中未停留"]
    unit: "次"
    timezone: "Asia/Shanghai"     # dt 分区按此时区切
    since_version: "v3"           # 该口径自哪一版生效

schema:                           # 1 Schema
  - {name: dt,         type: date,   nullable: false}
  - {name: user_id,    type: string, nullable: false}
  - {name: item_id,    type: string, nullable: false}
  - {name: channel,    type: string, nullable: false,
     enum: [app_home, app_search, h5, mini_program]}   # ⚠️ 枚举必须列全集
  - {name: dwell_ms,   type: int,    nullable: true, min: 0, max: 600000}
primary_key: [dt, user_id, item_id, ts]

slo:                              # 3 质量 SLO
  freshness: "T+1 07:00 前产出"
  row_count: {min: 3.0e8, max: 5.0e8}   # ⭐ 绝对量区间,比"波动 ±40%"严
  null_rate: {user_id: 0.0, dwell_ms: 0.05}
  uniqueness: {primary_key: 1.0}

change_policy:                    # 5 + 6
  breaking: ["删字段", "改类型", "改主键", "改口径定义", "缩小枚举"]
  non_breaking: ["加字段", "加枚举值(需下游确认)"]
  notice_days: 14
  dual_write_days: 7

⚖️ 三、兼容性矩阵:哪些改动算「破坏」

这是契约里最容易吵架的一张表,因为「对机器兼容」和「对语义兼容」是两回事。

变更 对机器 对语义 判定 说明
加一个新字段 ✅ 兼容 非破坏 但要防列错位(第 4 章)
删字段 破坏 任务直接挂,反而是好事 ⭐
改字段类型 破坏 CAST 可能让它静默通过 ⚠️
字段改名 破坏 同上
枚举新增取值 ⚠️ 对下游是破坏 ⭐⭐ 见下面的事故
枚举删除取值 破坏 历史数据的含义也变了
放宽可空性 ⚠️ 破坏 下游可能没有缺失处理
改口径定义 最危险 💀 机器完全看不出来(第 2 章)
改单位 最危险 💀 数值 ×100,类型可能都不变
改分区时区 最危险 💀 每天的边界挪了 8 小时(第 4 章)

⭐⭐ 这张表的核心结论会让任务挂掉的变更,反而是安全的变更。 真正危险的是那三个 💀 —— 它们对机器 100% 兼容,任务照常成功,只是答案错了。 所以契约必须校验语义(口径版本号、单位、时区、取值分布), 而不只是 schema。


🚦 四、怎么强制:三道闸

   上游产出                落地入库                 你消费
   ────────                ────────                 ────────
   闸①:生产端 CI          闸②:到达即校验          闸③:训练入口断言
   ────────                ────────                 ────────
   在上游的 PR 里跑         数据一落盘就跑            train.py 第一行就跑
   schema diff             SLO 断言                 你自己的假设
   → 破坏性变更            → 失败则【隔离到          → 不满足就 raise
     直接 block 合并 ⭐      quarantine】而不是
                            让它往下流 ⭐⭐

   成本:改一次 CI 配置      成本:一个调度任务         成本:20 行代码
   收益:把问题挡在源头      收益:脏数据不污染下游      收益:兜底,永远该有

⭐⭐ 闸② 的「隔离而不是继续流」是最反直觉、也最重要的设计: 大多数团队的做法是「校验失败 → 发个告警 → 数据照常写入」。 结果告警没人看(第 8 章),脏数据已经进了下游 20 张表。 正确做法:写入 quarantine 分区,主表保持昨天的数据,宁可"旧"不可"错"。

4.1 校验失败该阻断还是告警?按后果分级

级别 例子 动作
P0 阻断 主键重复、必填字段全空、行数不到下限一半 停止写入主表,值班群电话
P1 阻断+可覆盖 行数超出 SLO 区间、枚举出现未知值 停止写入,值班人可手动放行 ⭐
P2 告警 缺失率超阈值 20%、分位数偏移 照常写入,但打标 quality_flag
P3 记录 新增了字段、注释变更 只进变更日志

⚠️ 不要把所有检查都设成 P0。 一个天天误报的阻断规则,两周内一定会被人加上 --skip-check, 然后所有检查一起失效宁可少设几条,也要保证设的每一条都值得半夜被叫醒。


🛑 读到这里可以停 —— 前半章讲完了(约 16 分钟)。 后半章还有:一个能直接用的契约校验器 · 事故复盘:一个枚举值,23% 的流量被喂了默认值 · 落地顺序:别一上来就搞平台 回来的时候不用重读,直接从下一节接着看就行。


🧑‍💻 五、一个能直接用的契约校验器

import pandas as pd

def check_contract(df: pd.DataFrame, contract: dict):
    """按契约校验一个 DataFrame。返回 (是否通过, 违规列表)。
    contract 结构见本章 YAML:schema / primary_key / slo。"""
    v = []
    sch = {c["name"]: c for c in contract.get("schema", [])}

    for name, spec in sch.items():
        if name not in df.columns:
            v.append(("P0", name, "字段缺失", ""))
            continue
        s = df[name]
        if not spec.get("nullable", True) and s.isna().any():
            v.append(("P0", name, "不可空字段出现空值", f"{s.isna().mean():.2%}"))
        if "enum" in spec:
            allowed = set(spec["enum"])
            unknown = set(s.dropna().unique()) - allowed
            # ⭐ 枚举出现未知值 = 上游加了新取值却没通知,是最高频的静默破坏
            if unknown:
                share = s.isin(unknown).mean()
                v.append(("P1", name, f"未知枚举值 {sorted(unknown)[:5]}", f"占比 {share:.2%}"))
        if "min" in spec and pd.api.types.is_numeric_dtype(s) and s.min() < spec["min"]:
            v.append(("P1", name, "低于取值下限", f"{s.min()}"))
        if "max" in spec and pd.api.types.is_numeric_dtype(s) and s.max() > spec["max"]:
            # ⭐ 上限断言能抓住"单位从元改成分"这类 ×100 事故
            v.append(("P1", name, "超出取值上限", f"{s.max()}"))

    pk = contract.get("primary_key")
    if pk and all(c in df.columns for c in pk):
        dup = df.duplicated(subset=pk).sum()
        if dup:
            v.append(("P0", ",".join(pk), "主键重复", f"{dup} 行"))

    slo = contract.get("slo", {})
    rc = slo.get("row_count")
    if rc:
        n = len(df)
        if n < rc["min"] or n > rc["max"]:
            level = "P0" if n < rc["min"] / 2 else "P1"
            v.append((level, "<table>", "行数越界", f"{n} 不在 [{rc['min']:.0f},{rc['max']:.0f}]"))
    for col, thr in (slo.get("null_rate") or {}).items():
        if col in df.columns and df[col].isna().mean() > thr:
            v.append(("P2", col, "缺失率超阈值", f"{df[col].isna().mean():.2%} > {thr:.2%}"))

    blocking = [x for x in v if x[0] in ("P0", "P1")]
    return len(blocking) == 0, pd.DataFrame(v, columns=["级别", "字段", "问题", "详情"])

# 用法:ok, report = check_contract(df, contract_dict)
#      if not ok: 写入 quarantine 分区,不要写主表  ⭐

配套:在 CI 里检测破坏性变更(上游 PR 时跑):

import sys

BREAKING = {"删字段", "改类型", "改主键", "缩小枚举", "改口径定义", "改单位", "改时区"}

def diff_schema(old: dict, new: dict):
    """对比两版契约,输出变更清单和是否含破坏性变更。"""
    changes = []
    o = {c["name"]: c for c in old["schema"]}
    n = {c["name"]: c for c in new["schema"]}
    for name in o.keys() - n.keys():
        changes.append(("删字段", name, o[name].get("type"), None))
    for name in n.keys() - o.keys():
        changes.append(("加字段", name, None, n[name].get("type")))
    for name in o.keys() & n.keys():
        if o[name].get("type") != n[name].get("type"):
            changes.append(("改类型", name, o[name].get("type"), n[name].get("type")))
        if o[name].get("nullable") is False and n[name].get("nullable") is True:
            changes.append(("放宽可空", name, False, True))
        oe, ne = set(o[name].get("enum", [])), set(n[name].get("enum", []))
        if oe - ne:
            changes.append(("缩小枚举", name, sorted(oe - ne), None))
        if ne - oe:
            # ⭐ 新增枚举对机器兼容、对下游可能是破坏:单列出来强制下游确认
            changes.append(("新增枚举", name, None, sorted(ne - oe)))
    if old["primary_key"] != new["primary_key"]:
        changes.append(("改主键", "<pk>", old["primary_key"], new["primary_key"]))
    # ⭐ 语义变更靠版本号发现:口径改了就必须 bump since_version
    for k, sem in (new.get("semantics") or {}).items():
        old_sem = (old.get("semantics") or {}).get(k, {})
        if old_sem and old_sem.get("definition") != sem.get("definition"):
            changes.append(("改口径定义", k, old_sem.get("definition"), sem.get("definition")))
        for key, label in (("unit", "改单位"), ("timezone", "改时区")):
            if old_sem and old_sem.get(key) != sem.get(key):
                changes.append((label, k, old_sem.get(key), sem.get(key)))
    has_breaking = any(c[0] in BREAKING for c in changes)
    return has_breaking, changes

# CI 里:has_breaking, changes = diff_schema(old, new)
#        if has_breaking and not pr_has_label("approved-breaking"): sys.exit(1)

diff_schema 里最关键的一行是「改口径定义」那一段: 它只有在契约里把口径写成结构化字段时才可能被检出。 口径写在 wiki 里,机器永远查不出来;写进契约文件,才能进 CI。 这是「文档」和「契约」唯一的实质差别。


💀 六、事故复盘:一个枚举值,23% 的流量被喂了默认值

发生了什么

   系统:某平台的用户信用分模型,特征表里有一列 channel(注册渠道)
   训练时 channel 的取值:app_home / app_search / h5   (3 个)
   代码里的处理:
       CHANNEL_MAP = {"app_home": 0, "app_search": 1, "h5": 2}
       x = CHANNEL_MAP.get(row.channel, 3)     # ⚠️ 3 = "unknown"

   2 月 5 日:小程序端上线,埋点新增枚举值 mini_program
   上游的判断:"只是加了个取值,字段没动,向后兼容" —— 没有通知任何人

为什么没被发现

   ✅ 字段还在,类型没变,非空率 100%
   ✅ 行数在正常区间(小程序是增量,总量甚至涨了)
   ✅ 训练任务、推理服务全部正常,零报错
   ✅ 整体 AUC 从 0.781 微降到 0.774 —— 在正常波动范围内 💀

   💀 真相:mini_program 流量从 0% 一路涨到 23%,
      这 23% 的用户 channel 全部被映射成 "unknown"(3),
      而 "unknown" 在训练集里只占 0.4%,模型对它几乎没有学过
      → 这 23% 用户拿到的是一个近乎随机的分档

为什么整体 AUC 只降了 0.007:因为整体指标会稀释局部灾难。 23% 的流量彻底退化,但另外 77% 完全正常,加权平均后几乎看不出来。

   ⭐ 分渠道拆开之后(事后才做的):
      app_home / app_search / h5 :  AUC 0.784  (正常)
      mini_program              :  AUC 0.548  ⚠️ 接近随机

代价

   发现延迟:6 周(2/5 → 3/19)
   发现方式:小程序运营抱怨"通过率忽高忽低,看不出规律"

   影响:小程序渠道 41 万次授信决策使用了近乎随机的信用分
        其中约 1.7 万笔本该拒绝的通过了,事后逾期率是正常组的 3.4 倍
        直接坏账增量 ≈ 216 万元
   修复:重建映射 + 回溯 90 天特征 + 重训 = 8 人日

该补什么

缺失 补上之后
枚举没有全集约束 契约里 enum: 列全集,未知值触发 P1 阻断+可覆盖 ⭐⭐
.get(x, default) 静默兜底 改成显式:未知值 → 记录 + 计数 + 超阈值 raise ⭐
只看整体指标 关键维度(渠道/端/地域/新老客)必须分组看指标模型上线之后 14)⭐⭐
默认值占比没监控 unknown 占比从 0.4% → 23%,这条曲线本身就是完美的告警信号
上游不知道有消费者 契约的消费者清单 + 上游 CI 里的 diff_schema

💀 这个事故的两个通用教训① 「向后兼容」是上游视角的判断,不是下游视角的事实。 加一个枚举值对存储和类型完全兼容,对映射表是破坏性的。 ② 整体指标会稀释局部灾难。 23% 流量掉到随机水平,整体 AUC 只动了 0.007 —— 任何只看总体指标的监控,对"某一撮用户彻底坏掉"是全盲的。


🧠 七、落地顺序:别一上来就搞平台

   第 1 周:给你最重要的 3 张表,各写一份 30 行的契约(YAML 就行)⭐
           重点写【语义】和【消费者清单】——这两块光写下来就有价值
   第 2 周:把 check_contract 挂到训练入口(闸③),先只告警不阻断
   第 3–4 周:观察误报,收敛阈值;确认稳定后,把 P0/P1 改成真阻断
   第 2 个月:推动闸②(落地即校验 + quarantine 分区)
   第 3 个月:找上游谈闸①(CI 里跑 diff_schema)——
             ⭐ 这一步是政治问题不是技术问题,
                拿着前两个月攒下的【拦截记录】去谈,成功率高得多

最后那句是本章的实操核心: 上游不会因为「这是最佳实践」而接受一个会让他 CI 变红的东西。 他会因为「上季度这个检查拦下了 6 次事故,其中 2 次是你们改的」而接受。 先做闸③ 攒证据,再谈闸①。


🔗 这一章连到哪里

去哪 为什么
AI 基础设施 22 数据管线与存储 quarantine 分区、双写窗口、回溯补数在工程上怎么落
模型上线之后 17 版本回溯与可复现 契约版本号 + 数据版本号是「能复现」的前提
模型上线之后 14 辛普森悖论与七个陷阱 本章事故的第二个教训:整体指标会稀释局部灾难
智能体工程教程 14 评测 Evals 同一个思路的另一种应用:把「我觉得它变好了」变成会失败的测试

✅ 检查点

  1. 「上游改字段不通知」的三个结构性原因是什么?其中哪一条是所有治理方案真正在解的?
  2. 契约的七个部分是什么?为什么说「不算什么」那一行最值钱?
  3. 兼容性矩阵里,哪三种变更是「最危险」的?它们的共同特点是什么?为什么「会让任务挂掉的变更反而是安全的」?
  4. 三道闸分别在哪里、拦什么、成本多少?闸② 最反直觉的设计是什么、为什么?
  5. 校验失败分级里,P1 和 P2 的区别是什么?为什么不能把所有检查都设成 P0?
  6. check_contract 里哪两个断言分别能抓住什么事故?
  7. diff_schema 里「改口径定义」这段的前提条件是什么?「文档」和「契约」的实质差别是什么?
  8. 那个枚举事故:改了什么、23% 的流量发生了什么、为什么整体 AUC 只降 0.007、分渠道拆开后是多少、发现延迟和代价是多少?
  9. 落地顺序为什么是「先闸③ 再闸①」?谈闸① 时该拿什么去谈?
👀 答案
  1. 你不在他的 KPI 里(他甚至不知道有你这个消费者)②下游不可见(一张表可能有 40 个消费者,上游看到的是空白)③⭐⭐改动成本严重不对称(上游改一行 5 分钟,下游发现+定位+修复+回溯 3 天,省下的归他、花掉的归你)。所有治理方案本质都在解第 ③ 条:把成本挪回去。
  2. Schema / 语义 / 质量 SLO / 所有权 / 变更流程 / 兼容性等级 / 消费者清单。「不算什么」值钱是因为口径的边界永远比口径的主体更容易出事 —— 写清「预加载不算、后台刷新不算、快速滑过不算」比写十行定义有用得多。
  3. 改口径定义、改单位、改分区时区。共同点:对机器 100% 兼容,任务照常成功,只是答案错了。而删字段/改类型会让任务直接挂 —— 会让任务挂掉的变更反而是安全的变更,因为你立刻就知道了。
  4. 闸①生产端 CI(上游 PR 里跑 schema diff,破坏性变更 block 合并,成本=改一次 CI 配置);闸②落地即校验(数据一落盘跑 SLO 断言,成本=一个调度任务);闸③训练入口断言(train.py 第一行,成本=20 行代码)。闸② 最反直觉的是⭐⭐失败要隔离到 quarantine 分区、主表保持昨天的数据,而不是"告警但照常写入" —— 因为告警没人看,脏数据已经进了下游 20 张表。宁可旧,不可错。
  5. P1 阻断但值班人可手动放行(行数越界、未知枚举值);P2 只告警、照常写入但打 quality_flag(缺失率超阈值、分位数偏移)。不能全设 P0 是因为⚠️一个天天误报的阻断规则两周内一定会被加上 --skip-check,然后所有检查一起失效。宁可少设几条,保证每一条都值得半夜被叫醒。
  6. 枚举未知值检测抓「上游加了新取值没通知」(本章的 23% 事故);取值上限断言抓「单位从元改成分」这类 ×100 事故。
  7. 前提是把口径写成契约里的结构化字段definition / unit / timezone / since_version)。口径写在 wiki 里机器永远查不出来,写进契约文件才能进 CI —— 这是文档和契约唯一的实质差别。
  8. 小程序端上线,埋点新增枚举值 mini_program,上游认为「只是加个取值,向后兼容」没通知任何人。代码里 CHANNEL_MAP.get(row.channel, 3) 把它静默映射成 unknown,而 unknown 在训练集只占 0.4%。mini_program 流量从 0% 涨到 23%,这批用户拿到近乎随机的分档。整体 AUC 只从 0.781 降到 0.774,因为整体指标会稀释局部灾难(77% 正常);分渠道拆开:其他渠道 0.784,mini_program 只有 0.548,接近随机发现延迟 6 周,靠运营抱怨「通过率忽高忽低」;代价 41 万次授信用了随机分、约 1.7 万笔该拒的通过了、逾期率是正常组 3.4 倍、坏账增量约 216 万元,而修复只要 8 人日
  9. 因为闸① 是政治问题不是技术问题 —— 上游不会因为「这是最佳实践」接受一个会让他 CI 变红的东西。先做闸③(20 行代码、自己就能做)攒拦截记录,然后拿着「上季度这个检查拦下了 6 次事故,其中 2 次是你们改的」去谈,成功率高得多。

🛑 可以停在这里

走神救援

契约做的事只有一件:把「我们说好了不改」变成「你一改,你的 CI 就红」。 上游不通知有三个结构性原因:①你不在他的 KPI 里(他甚至不知道有你)②下游不可见(一张表 40 个消费者,他打开元信息什么都看不到)③⭐⭐改动成本严重不对称——上游改一行 5 分钟、下游定位修复回溯 3 天,省下的归他、花掉的归你;所有治理方案本质都在解第③条。⚠️没有强制执行的契约 = 措辞更正式的口头约定。 七个部分:Schema / 语义 / 质量 SLO / 所有权 / 变更流程 / 兼容性等级 / 消费者清单;语义里⭐「不算什么」那一行最值钱——写清「预加载不算、后台刷新不算、快速滑过不算」比写十行定义有用,口径的边界永远比主体更容易出事。⭐⭐兼容性矩阵的核心结论:会让任务挂掉的变更反而是安全的(删字段、改类型立刻暴露);真正危险的是三个💀——改口径定义、改单位、改时区,它们对机器 100% 兼容、任务照常成功、只是答案错了枚举新增取值对机器兼容、对下游映射表是破坏三道闸:①上游 CI 跑 diff_schema,破坏性变更 block 合并;②落地即校验,⭐⭐失败要写进 quarantine 分区、主表保留昨天的数据,而不是"告警但照常写入"(告警没人看,脏数据已经进了下游 20 张表)——宁可旧不可错;③训练入口断言,20 行代码,永远该有。分级:P0 阻断(主键重复、必填全空)、P1 阻断可覆盖(行数越界、未知枚举)、P2 只告警打 quality_flag、P3 记录;⚠️别全设 P0——天天误报的规则两周内一定被加 --skip-check,然后所有检查一起失效。代码里两个关键断言:未知枚举检测抓上游偷加取值,取值上限断言抓「元改成分」的 ×100。⭐口径写进契约的结构化字段才进得了 CI,写在 wiki 里机器永远查不出来——这是文档和契约唯一的实质差别。 💀事故:小程序上线,埋点新增枚举 mini_program,上游认为「向后兼容」没通知;代码 CHANNEL_MAP.get(x, 3) 把它静默映射成 unknown(训练集里只占 0.4%)。这批流量从 0% 涨到 23%,拿到近乎随机的信用分,而整体 AUC 只从 0.781 掉到 0.774——因为整体指标会稀释局部灾难;分渠道拆开才看到:其他渠道 0.784、mini_program 只有 0.548发现延迟 6 周,靠运营抱怨「通过率没规律」;41 万次授信受影响、约 1.7 万笔该拒的通过、逾期率是正常组 3.4 倍、坏账增量约 216 万元,修复只要 8 人日。两条通用教训:「向后兼容」是上游视角的判断不是下游视角的事实任何只看总体指标的监控,对「某一撮用户彻底坏掉」是全盲的。落地顺序:先闸③ 攒拦截记录,再拿记录去谈闸①——那是政治问题不是技术问题。

下一节 👉 04-脏数据的十种形态.md

打卡记录保存在你的浏览器里,首页能看到总进度