DDIA · 数据集成与哲学

Ch.13-14: 全书收束 — 一切皆派生: 让日志成为唯一事实源, 索引/缓存/视图都是它的投影; 及时性可以让步, 完整性寸步不让

数据库解绑 unbundling: 把库拆成"日志 + 投影" 全序变更日志 唯一事实源 (state machine) 搜索索引投影 缓存投影 物化视图投影 → 由流处理器 (Flink/Materialize) → 按 同一日志 确定性派生 → 各视图可删可重建 联邦 (Trino/FDW) 只统一读, 解决不了写同步 — 派生才解决 状态机复制 = 全序日志 × 确定性应用 所有输入按同一全序过日志, 每个副本/视图按序重放 → 状态必然一致 Ch.10 的全序广播在这里落地; Ch.11-12 的批/流是它的两种执行形态 lambda vs kappa: 两套代码 vs 只留流 lambda 架构 批层 (全量准确, 慢) + 加速层 (增量近似, 快) → 合并层拼结果 代价: 同一逻辑两套代码, 口径永远对不齐 Marz 自己也承认该被 kappa 取代 kappa 架构 只保留流: 一套逻辑, 需要重算 = 重放日志 (新版本消费者从 0 重跑) 批 = 对日志的 bounded 重放, 流 = 持续的 bounded 重放 前提: 日志保得够久 / 可重放 端到端论点 + timeliness vs integrity end-to-end argument: 底层兜不住业务正确性 TCP 去重 ≠ 幂等业务; 至少一次投递 ≠ 不重复扣款 解: 客户端生成 request ID 贯穿全链路, 端到端去重/校验 timeliness (暂时旧) vs integrity (永久错) 读到旧值: 暂时、自愈 (最终一致兜底) — 可接受 数据损坏: 永久、须人工修 — 不可接受 异步派生放弃 timeliness 换吞吐, 但 integrity 必须端到端保证 trust but verify: 对账/审计/scrubbing 是最后防线 确定性重放对账 + 全量 scrubbing 任务 — 假设所有组件都会说谎 跨服务一致性: saga 补偿链 + 伦理一页 下单 ✓ 扣款 ✓ 发货 ✗ 失败 补偿链反向执行 退款 ✓ 撤单 ✓ 每步必须可补偿/幂等; 比笛 2PC: 无全局阻塞, 但中间态可见 (业务要容忍) Ch.14 伦理: 数据系统的另一面 预测分析给人打分 → bias 偏见放大 → feedback loop 向下螺旋 consent 知情同意名存实亡 / 去标识化可被再识别 工程侧回应: 数据最小化 / crypto-shredding (Ch.3) / 可解释性

全书主旨: 一切皆派生

  • • 一个全序日志 = 唯一事实源
  • • 索引/缓存/视图 = 日志的确定性投影
  • • 投影可删可重建 (推倒重放)
  • • 批与流只是同一逻辑的两种跑法

集成的哲学

  • • kappa: 只留流, 重算=重放
  • • 端到端论点: request ID 贯穿全链
  • • timeliness 可让, integrity 不让
  • • trust but verify: 对账是最后防线

跨服务一致性

  • • saga 补偿链替代分布式事务
  • • 每步可补偿 + 幂等 + 可重试
  • • 协调避免: 异步派生保完整性
  • • 伦理: 最小化/可删除/防偏见

💡 一句话理解

全书在此收束成一句话: 数据库本来就是"日志 + 若干投影" — 索引是投影, 缓存是投影, 物化视图是投影, 那为什么不把它们拆出来, 让一个全序日志当唯一事实源, 流处理器当"通用投影机"? kappa是这个思想的最纯形态: 只留流, 要重算就重放日志。跨系统一致性则靠放下执念: 分布式事务不是必需品, saga 补偿链 + 端到端 request ID + 定期对账 (trust but verify) 足以守住完整性这条不可退让的底线 — 而读到一秒前的旧数据 (timeliness) 是可以买回来的代价。最后一章提醒我们: 数据系统的输出会决定人的贷款、招聘与自由, bias 会自我放大。

🧠 必知必会 必考 & 必会

data integration
多个专用系统 (库/索引/缓存/批) 如何保持同步 — Part III 的核心命题。答案不是"更好的 2PC", 而是让派生关系显式化。
# 反面: 5 个系统各自写自己的副本
# 正面: 1 个日志 → 4 个系统作为投影消费
state machine replication
所有输入经全序日志、各副本确定性地应用 → 状态必然一致。全序广播 (Ch.10) + 纯函数 fold = 它的实现。
state = {}
for ev in ordered_log():     # 全序保证
    state = fold(state, ev)   # 确定性: 同输入同输出
# 任何副本重放 → 同一状态 (呼应事件溯源)
unbundling the database
把数据库"解绑": 索引/物化视图/缓存都对变更日志派生, 用流处理器实现 — 相当于把 InnoDB 的分工变成一组可独立演化的组件。
# 库内: WAL → 二级索引 → 物化视图 (黑盒)
# 解绑: WAL → Kafka → Flink → ES/Redis/宽表 (白盒)
# 收益: 每个投影独立扩缩、独立技术栈
federated / polystore
联邦数据库 (Trino/FDW) 统一读 — 一个 SQL 查遍多源; 但它不搬数据、不解决写同步, 与派生是互补不是替代。
SELECT * FROM mysql.orders o
JOIN es.user_profile p ON o.user_id = p.id;  -- Trino
# 读爽了, 写还是各写各的 → 派生管道仍需要
lambda architecture
批层 (全量准)+加速层 (增量快)+合并层; 逻辑写两遍, 口径永远漂移 — 作者本人也承认被 kappa 取代。
# 同一个"月活"公式写两遍:
#   批: Java/Hive     流: Python/Storm
# 修 bug 要修两处, 还要处理两路结果的合并边界
kappa architecture
只保留流: 一套逻辑, 重算 = 重放日志 (新版本消费者从 0 跑); 批退化为"对日志的受限重放"。前提: 日志保存足够久、处理确定性。
# 口径变了?
#   新消费者 group 从 offset=0 重放 → 重建视图
#   旧视图切流下线 — 永远只有一份逻辑
end-to-end argument
底层 (TCP 去重、消息至少一次) 不能替代应用层检查: 客户端生成唯一 request ID 贯穿全链, 端到端去重 — 正确性放不下进任何单层。
# request_id 由客户端生成, 全链透传:
POST /pay {"req_id": "9f8e..."} → RPC 带 → MQ 带 → 库存表记录
# 任何一层重复投递, 终点都能识别并去重
saga / 补偿事务
跨服务一致性: 正向操作逐步执行, 失败时反向补偿 (撤单/退款); 每步必须幂等可重试 — 比笛 2PC 无全局阻塞, 但中间态对外可见。
# 订单→扣款→发货; 发货失败:
compensate: refund() → cancel_order()
# 每步落 saga log (断点续跑); 步骤语义 = "最终成功"
timeliness vs integrity
读到旧值: 暂时、自愈 (最终一致兜底) — 可让步换吞吐; 数据损坏: 永久、人工修 — 不可让步。异步派生的世界里, 这是核心权衡坐标。
# 让 timeliness: 关注页 2s 后看到新帖 (OK)
# 不让 integrity: 库存扣成负数 (永不可能 OK)
coordination avoidance
协调避免: 异步派生保 integrity、放弃 timeliness → 大多数约束 (唯一性/和不变量) 无需分布式事务也能守 — 冲突后再检测/补偿即可。
# 库存: 各副本先扣 (异步), 对账发现负数再补偿
# 比笛: 每次 zk 协调 (同步) → 吞吐低一个量级
# 适用: 短暂超卖可由补偿挽回的场景
trust but verify
信任但核实: 确定性重放对账、全量审计、scrubbing 任务 — 假设所有组件都会说谎, integrity 的最后防线。
# 每日对账:
sum(orders.amount) == sum(billing.amount) # ? 差额报警
# scrubbing: 全表扫描校验校验和 / 重算派生列
bias 与 feedback loop
算法从带偏见的数据学习并放大 (预测分析给贷款/招聘打分); 反馈循环: 信用差→难就业→更穷→分数更差 — 数据系统的输出反哺输入。
# 恶性闭环:
score low → 不给贷款 → 收入差 → 明年 score 更 low
# 工程回应: 剔除代理变量 + 人审出口 + 可解释性
consent / de-identification
GDPR 要求 freely-given 的知情同意, 但网络效应下"不同意就别用"使同意名存实亡; 去标识化常可被再识别 — 数据最小化才是根本解。
# 数据最小化: 存储成本 ≠ 全部成本
#   泄露/责任/合规/伦理 都在"存"这个动作里
# 配合 crypto-shredding 满足删除权 (Ch.3)

🏭 生产实战 real world

场景 1 · lambda→kappa: 把两套月活口径合成一套

批层 Hive 与加速层 Storm 的"月活"永远差 0.3% — kappa 重放统一。

# 改造前: 月活 = batch(hive月表) ∪ merge(realtime storm)
# 差异来源: 两套去重逻辑 + 合并边界

# 改造后 (kappa):
#   唯一定义: MAU = Flink 消费 login-events, 按 event_time 日去重
#   历史重算: 新版本消费者从 offset(当年1月1日) 重放
#   批查询  = 对同一日志的 bounded 读取 (不再有独立批层代码)

口径差异从 0.3% 变 0 — 因为只剩一套逻辑。

场景 2 · 数据库解绑: 把 ES 索引变成"投影"运维

ES 索引坏了不再"手工补数" — 按投影协议重建, 与 CQRS 同一套纪律。

# 投影协议 (每个派生视图必须回答的三问):
#   1) 源日志: db.orders 的 CDC topic
#   2) fold 函数: 纯确定性 (代码评审重点)
#   3) 重建流程: create index v2 → 从 0 重放 → 别名切换
es_rebuild(alias="orders-idx", from_offset=0, fold=es_fold)
# 坏了 = 删了重放, 停机时间 = 重放时长 (分钟级)

配套监控: 每个投影记录 "重放水位 lag", 滞后过大自动告警。

场景 3 · 端到端 request ID: 终结"重复扣款"排查地狱

客户端重试 + MQ 重投 + RPC 超时重试三层叠加 — request ID 贯穿后一切可去重。

# 生成: 客户端 (用户点击时生成 uuid)
req_id = uuid4()
# 透传: HTTP header → RPC metadata → MQ header → 每层落库
dedup = INSERT IGNORE INTO seen(req_id, service)
if dedup.rowcount == 0:          # 见过 → 直接返回旧结果
    return lookup_result(req_id)
# 排查: req_id 一查到底, 全链路行为一目了然

这就是端到端论点的落地: 唯一性必须在"终点"由业务保证。

场景 4 · saga 落地: 订单-支付-库存补偿链

跨三个服务不能 2PC — saga 编排器逐步执行, 失败反向补偿。

steps = [
    (create_order, cancel_order),      # (正向, 补偿)
    (charge_wallet, refund_wallet),
    (deduct_stock, restore_stock),
]
for fwd, comp in steps:
    try:
        saga_log.append(fwd.__name__, "start")
        fwd(ctx)                        # 幂等 + 带 req_id
        saga_log.append(fwd.__name__, "done")
    except StepFailed:
        for prev in reversed(done_steps):
            prev.comp(ctx)              # 反向补偿
        raise OrderAborted

纪律: 中间态可见 (已扣款未发货) 要有对客解释; 补偿也可能失败 → 告警人工兜底。

场景 5 · 协调避免: 库存超卖的自愈设计

热点商品不再每次秒杀都 zk 协调 — 副本先扣, 对账发现超卖再补偿。

# 写路径 (无协调): 各分片本地原子扣减
UPDATE stock_local SET n = n - 1 WHERE sku=7 AND n > 0;

# 异步对账 (每 30s): 汇总各分片 n, 若总库存 < 0
if total < 0:
    compensate(latest_orders[-abs(total):])   # 撤销末尾订单+退款
# 业务前提: 极小概率的超卖可由补偿挽回 (发券道歉)

坐标自查: 换的是 timeliness (几秒内可能超卖), 保的是 integrity (终态不坏)。

场景 6 · trust but verify: 每日对账守护资金完整性

"我们的管道不丢数据"是信仰不是证据 — 每日三方对账是底线。

# 每日 03:00 对账任务 (确定性重放思路)
src  = sum(orders.amount, day=d)          # 源头 SoR
dest = sum(billing.amount, day=d)         # 派生计费
if src != dest:
    diff_rows = full_join_find_diff(d)    # 精确到行
    alert("INTEGRITY BREACH", diff_rows)  # page, 修复+根因

# 加固: 关键表行数/校验和日环比, 异动即查

原则: 对账的差异数与定位时长都是 SLI — 不是"有没有对账"。

场景 7 · 数据最小化: 别把 PII 存进宽表

"先全存了再说"让一次泄露变成公司事故 — 收集时就裁剪。

# 需求: 推荐系统需要"兴趣画像", 不需要身份证号
# 错: 全字段进特征宽表 (含手机号/身份证)

# 对: 最小化 schema + 密钥隔离
features = {"interest_tags": tags, "activity_bucket": bucket}
# PII (手机号) 只在账号服务, 派生侧永远拿不到
# 需要 join 时用不可逆 user_hash

记账方式: 每存一份 PII = 多一份泄露面 + 一条 GDPR 删除义务。

场景 8 · bias 审计: 给推荐模型加"公平性回归"

模型上线前测分组指标, 防止"数据偏见"变成产品歧视。

# 上线门禁 (与性能回归门禁并列):
for g in (gender, region, age_band):
    pass_rate[g] = model_approval_rate(group=g)
# 断言: 最大组间差异 < 5%, 超标 → 拦截 + 特征审计
# 审计清单: 有没有代理变量 (邮编≈收入≈种族)?

呼应 feedback loop: 模型决定谁被看见, 被看见的数据又训练模型 — 切断闭环要靠人审出口。

场景 9 · 可删除架构: 删除权从设计开始

用户注销后, 20 个派生系统里的数据怎么删 — 密钥粉碎 + 投影重建双保险。

# SoR 侧: crypto-shredding (Ch.3)
kms.destroy_key(user_123)              # PII 永不可读

# 派生侧: 投影按"删除事件"清理
# CDC: {"op":"delete-user","uid":123} → 各投影消费
#   ES: delete_by_query(uid=123)
#   缓存: del user:123:*
#   数仓: 次日分区重写 (无该用户)
# 审计: 删除完成率报表 — 合规可验证

设计期问题: "这个字段进了几个投影?" — 回答不了就是还没准备好删除权。

场景 10 · 全书终局演练: 从 P99 报警到"日志重放"的完整闭环

把 10 页串成一个事故: 报警 → 定位 → 修复 → 用"重放"自愈。

# ① 总纲页: P99 报警, 分位数组合读 → 长尾问题
# ② 事务页: 写偏斜? 事务内 RPC? → 定位到双写索引脏数据
# ③ 批流页: dual-write 交错 → 索引与库分叉
# ④ 集成页: 改造为 outbox + CDC, 索引 = 投影
# ⑤ 自愈: 索引投影从日志头重放 → 数据自动收敛
# ⑥ 韧性页: 重试预算/熔断防复发; 复盘进 blameless 文档

这就是 320+ 术语的最终用法: 它们不是名词表, 是一条能走通的排障路径。

⚠️ 编码注意与常见坑 pitfalls

坑 1 · 联邦查询当集成 — Trino 读得爽, 写同步烂摊子还在. 原因: 混淆读/写解耦. 正解: 联邦+派生管道互补。
# 错: "上了 Trino 就不用同步数据了"
# 对: 联邦只管读, 写走日志派生
坑 2 · lambda 双代码路径 — 批/流两套逻辑口径漂移. 原因: 历史惯性. 正解: kappa — 一套逻辑 + 重放。
# 错: Hive 口径 vs Storm 口径并存
# 对: 同一函数, 日志重放重建
坑 3 · 重放依赖不确定代码 — fold 里有 now()/rand(), 重放结果不同. 原因: 状态机复制前提被破坏. 正解: 投影纯函数化, 时间取事件时间。
# 错: fold(state, ev): state["t"]=now()
# 对: state["t"]=ev.event_ts
坑 4 · 派生视图当 SoR — 手工改 ES 里的数据, 投影再也不会收敛. 原因: 图省事. 正解: 只改源日志, 视图一律重放。
# 错: 直接 UPDATE ES 文档修数
# 对: 源头修正 → 日志重放
坑 5 · 底层去重当业务幂等 — TCP 不丢≠不会重复扣款. 原因: 端到端论点缺席. 正解: request ID 端到端 + 终点去重。
# 错: "Kafka 至少一次, 不会丢"
# 对: 至少一次=会重复 → 终点幂等
坑 6 · saga 步骤不可补偿 — "发短信"没法撤回, 补偿链断. 原因: 步骤设计没考虑撤销. 正解: 不可补偿动作放最后 / 预授权模式。
# 错: 发短信 → 扣款 → 扣款失败
# 对: 扣款 → 发短信 (短信放最后)
坑 7 · saga 中间态裸奔 — 已扣款未发货对用户不可解释. 原因: 只想成功路径. 正解: 状态机可见 + 超时补偿告警。
# 错: 卡在"已扣款"三天没人管
# 对: 中间态超时 → 自动补偿/人工队列
坑 8 · 补偿不幂等 — 补偿重试 = 退款两次. 原因: 只测正向. 正解: 补偿同样带 req_id 幂等。
# 错: refund(req_id) 不去重
# 对: refund ON CONFLICT(req_id)
坑 9 · coordination avoidance 用错场景 — 资金账户"先扣后补", 超卖现金. 原因: 无条件乐观. 正解: 只对"可补偿挽回"的约束避免协调, 资金类仍强约束。
# 错: 余额先扣后对账
# 对: 余额强约束; 库存类可异步+补偿
坑 10 · 对账缺失 — 派生错了几个月才发现. 原因: "管道很稳". 正解: 日级对账 + 差异数当 SLI。
# 错: 上线三年没对过账
# 对: 每日对账任务 + 差异 page
坑 11 · integrity 让位于 timeliness — "实时优先"把校验全关了, 数据悄悄坏. 原因: 权衡坐标颠倒. 正解: 校验异步化但不省略 (先落可疑区再核)。
# 错: 为了低延迟跳过校验
# 对: 异步校验 + 坏数据隔离区
坑 12 · kappa 日志保不住 — 要重放时日志已 TTL 清掉. 原因: 没按重算窗口规划保留期. 正解: 日志保留 ≥ 最长重算窗口, 或定期快照+日志续接。
# 错: retention=7d, 口径要重算 1 年
# 对: 快照 + 日志续接 (compacted 基线)
坑 13 · 多副本无唯一 request ID — 各处自生成, 去重失效. 原因: ID 在中层生成. 正解: 客户端生成, 透传不重写。
# 错: 每层服务自己 uuid4
# 对: 客户端生成, header 透传
坑 14 · 伦理字段裸进特征 — 邮编/性别直接当特征, 代理歧视. 原因: 只看 AUC. 正解: 特征审计 + 分组公平性门禁。
# 错: features += [zipcode, gender]
# 对: 代理变量审计 + 差异 <5% 门禁
坑 15 · 偏见反馈环无出口 — 模型决定一切, 闭环锁死用户. 原因: 全自动化. 正解: 低分用户人审出口 + 抽样反事实测试。
# 错: score<60 全自动拒绝
# 对: 边界带人审 + 定期重训纠偏
坑 16 · 注销只删主库 — 20 个派生系统里 PII 还活着. 原因: 没有删除事件流. 正解: 删除事件进日志, 投影各自消费清理 + 完成率审计。
# 错: DELETE FROM users 完事
# 对: delete-user 事件 → 各投影清理
坑 17 · 全量 PII 进宽表 — 特征表带手机号, 泄露面×N. 原因: 收集无最小化. 正解: 最小化 schema + 不可逆 hash join。
# 错: 宽表 select * 全带
# 对: user_hash + 特征白名单
坑 18 · 状态机复制依赖可变排序 — 日志分区乱序/多写入者, 重放分叉. 原因: 全序没保障 (Ch.10). 正解: 单分区/共识定序, 事件位序即序。
# 错: topic 12 分区, 各副本各自读
# 对: 需要全序的状态进单分区/共识日志
坑 19 · 投影无 lag 监控 — 视图落后 6 小时没人知道, 业务决策用旧数. 原因: 只看管道"在跑". 正解: 重放水位 lag + 停滞告警。
# 错: 只监控进程存活
# 对: projection_lag > 5min 告警
坑 20 · 复盘追责而非改系统 — 事故根因停在"人失误". 原因: blame 文化 (Ch.2). 正解: blameless 复盘, 改进项落在系统/流程。
# 错: "谁配错的" → 记过
# 对: 变更防呆 + 校验自动化