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