Ch.13-14: 全书收束 — 一切皆派生: 让日志成为唯一事实源, 索引/缓存/视图都是它的投影; 及时性可以让步, 完整性寸步不让
全书在此收束成一句话: 数据库本来就是"日志 + 若干投影" — 索引是投影, 缓存是投影, 物化视图是投影, 那为什么不把它们拆出来, 让一个全序日志当唯一事实源, 流处理器当"通用投影机"? kappa是这个思想的最纯形态: 只留流, 要重算就重放日志。跨系统一致性则靠放下执念: 分布式事务不是必需品, saga 补偿链 + 端到端 request ID + 定期对账 (trust but verify) 足以守住完整性这条不可退让的底线 — 而读到一秒前的旧数据 (timeliness) 是可以买回来的代价。最后一章提醒我们: 数据系统的输出会决定人的贷款、招聘与自由, bias 会自我放大。
# 反面: 5 个系统各自写自己的副本 # 正面: 1 个日志 → 4 个系统作为投影消费
state = {}
for ev in ordered_log(): # 全序保证
state = fold(state, ev) # 确定性: 同输入同输出
# 任何副本重放 → 同一状态 (呼应事件溯源)# 库内: WAL → 二级索引 → 物化视图 (黑盒) # 解绑: WAL → Kafka → Flink → ES/Redis/宽表 (白盒) # 收益: 每个投影独立扩缩、独立技术栈
SELECT * FROM mysql.orders o JOIN es.user_profile p ON o.user_id = p.id; -- Trino # 读爽了, 写还是各写各的 → 派生管道仍需要
# 同一个"月活"公式写两遍: # 批: Java/Hive 流: Python/Storm # 修 bug 要修两处, 还要处理两路结果的合并边界
# 口径变了? # 新消费者 group 从 offset=0 重放 → 重建视图 # 旧视图切流下线 — 永远只有一份逻辑
# request_id 由客户端生成, 全链透传: POST /pay {"req_id": "9f8e..."} → RPC 带 → MQ 带 → 库存表记录 # 任何一层重复投递, 终点都能识别并去重
# 订单→扣款→发货; 发货失败: compensate: refund() → cancel_order() # 每步落 saga log (断点续跑); 步骤语义 = "最终成功"
# 让 timeliness: 关注页 2s 后看到新帖 (OK) # 不让 integrity: 库存扣成负数 (永不可能 OK)
# 库存: 各副本先扣 (异步), 对账发现负数再补偿 # 比笛: 每次 zk 协调 (同步) → 吞吐低一个量级 # 适用: 短暂超卖可由补偿挽回的场景
# 每日对账: sum(orders.amount) == sum(billing.amount) # ? 差额报警 # scrubbing: 全表扫描校验校验和 / 重算派生列
# 恶性闭环: score low → 不给贷款 → 收入差 → 明年 score 更 low # 工程回应: 剔除代理变量 + 人审出口 + 可解释性
# 数据最小化: 存储成本 ≠ 全部成本 # 泄露/责任/合规/伦理 都在"存"这个动作里 # 配合 crypto-shredding 满足删除权 (Ch.3)
批层 Hive 与加速层 Storm 的"月活"永远差 0.3% — kappa 重放统一。
# 改造前: 月活 = batch(hive月表) ∪ merge(realtime storm) # 差异来源: 两套去重逻辑 + 合并边界 # 改造后 (kappa): # 唯一定义: MAU = Flink 消费 login-events, 按 event_time 日去重 # 历史重算: 新版本消费者从 offset(当年1月1日) 重放 # 批查询 = 对同一日志的 bounded 读取 (不再有独立批层代码)
口径差异从 0.3% 变 0 — 因为只剩一套逻辑。
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", 滞后过大自动告警。
客户端重试 + 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 一查到底, 全链路行为一目了然
这就是端到端论点的落地: 唯一性必须在"终点"由业务保证。
跨三个服务不能 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
纪律: 中间态可见 (已扣款未发货) 要有对客解释; 补偿也可能失败 → 告警人工兜底。
热点商品不再每次秒杀都 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 (终态不坏)。
"我们的管道不丢数据"是信仰不是证据 — 每日三方对账是底线。
# 每日 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 — 不是"有没有对账"。
"先全存了再说"让一次泄露变成公司事故 — 收集时就裁剪。
# 需求: 推荐系统需要"兴趣画像", 不需要身份证号 # 错: 全字段进特征宽表 (含手机号/身份证) # 对: 最小化 schema + 密钥隔离 features = {"interest_tags": tags, "activity_bucket": bucket} # PII (手机号) 只在账号服务, 派生侧永远拿不到 # 需要 join 时用不可逆 user_hash
记账方式: 每存一份 PII = 多一份泄露面 + 一条 GDPR 删除义务。
模型上线前测分组指标, 防止"数据偏见"变成产品歧视。
# 上线门禁 (与性能回归门禁并列): for g in (gender, region, age_band): pass_rate[g] = model_approval_rate(group=g) # 断言: 最大组间差异 < 5%, 超标 → 拦截 + 特征审计 # 审计清单: 有没有代理变量 (邮编≈收入≈种族)?
呼应 feedback loop: 模型决定谁被看见, 被看见的数据又训练模型 — 切断闭环要靠人审出口。
用户注销后, 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 报警, 分位数组合读 → 长尾问题 # ② 事务页: 写偏斜? 事务内 RPC? → 定位到双写索引脏数据 # ③ 批流页: dual-write 交错 → 索引与库分叉 # ④ 集成页: 改造为 outbox + CDC, 索引 = 投影 # ⑤ 自愈: 索引投影从日志头重放 → 数据自动收敛 # ⑥ 韧性页: 重试预算/熔断防复发; 复盘进 blameless 文档
这就是 320+ 术语的最终用法: 它们不是名词表, 是一条能走通的排障路径。
# 错: "上了 Trino 就不用同步数据了" # 对: 联邦只管读, 写走日志派生
# 错: Hive 口径 vs Storm 口径并存 # 对: 同一函数, 日志重放重建
# 错: fold(state, ev): state["t"]=now() # 对: state["t"]=ev.event_ts
# 错: 直接 UPDATE ES 文档修数 # 对: 源头修正 → 日志重放
# 错: "Kafka 至少一次, 不会丢" # 对: 至少一次=会重复 → 终点幂等
# 错: 发短信 → 扣款 → 扣款失败 # 对: 扣款 → 发短信 (短信放最后)
# 错: 卡在"已扣款"三天没人管 # 对: 中间态超时 → 自动补偿/人工队列
# 错: refund(req_id) 不去重 # 对: refund ON CONFLICT(req_id)
# 错: 余额先扣后对账 # 对: 余额强约束; 库存类可异步+补偿
# 错: 上线三年没对过账 # 对: 每日对账任务 + 差异 page
# 错: 为了低延迟跳过校验 # 对: 异步校验 + 坏数据隔离区
# 错: retention=7d, 口径要重算 1 年 # 对: 快照 + 日志续接 (compacted 基线)
# 错: 每层服务自己 uuid4 # 对: 客户端生成, header 透传
# 错: features += [zipcode, gender] # 对: 代理变量审计 + 差异 <5% 门禁
# 错: score<60 全自动拒绝 # 对: 边界带人审 + 定期重训纠偏
# 错: DELETE FROM users 完事 # 对: delete-user 事件 → 各投影清理
# 错: 宽表 select * 全带 # 对: user_hash + 特征白名单
# 错: topic 12 分区, 各副本各自读 # 对: 需要全序的状态进单分区/共识日志
# 错: 只监控进程存活 # 对: projection_lag > 5min 告警
# 错: "谁配错的" → 记过 # 对: 变更防呆 + 校验自动化