KEEL · 龙骨 · A CURRICULUM FOR THE AI ERA
00 · 一条变更的旅程 — keel 龙骨
这一章先把整条路的形状画出来。后面六章都在放大这条路上的某一段,但在你脑子里有这张图之前,每一段的细节都是孤立的。
这一章先把整条路的形状画出来。后面六章都在放大这条路上的某一段,但在你脑子里有这张图之前,每一段的细节都是孤立的。
现场
客服改了一个订单的状态,从 created 改成 paid。数据库里只是一次 UPDATE。半小时后运营问:「为什么报表里这单还是未支付?」
如果报表是每晚跑批算出来的,答案是「它明天才会变」。如果报表是流式的,那这次 UPDATE 应该已经走进了分析表,只是没有——这正是要排查的现场:数据卡在了链路的哪一段。
先看清楚这条路上到底有几个进程、几块存储。
全链路
flowchart TD
subgraph PG["业务库 PostgreSQL 5433 (labcdc)"]
APP["应用 / 人工操作<br/>UPDATE orders"] --> TBL["表 public.orders<br/>order_id PK"]
TBL --> WAL["pg_wal<br/>物理 WAL 记录"]
WAL --> SLOT["逻辑复制槽 cdc_slot<br/>test_decoding 插件"]
end
subgraph PIPE["管道进程(独立于数据库)"]
SLOT -->|①解码后的变更原文| PROD["cdc_pipeline.py<br/>读槽 + 解析"]
PROD --> JSONL["下游A:本地 JSONL<br/>02-cdc-orders.jsonl"]
PROD -->|②JSON 信封| KAFKA["下游B:Kafka 集群B<br/>topic cdc-orders-01"]
end
subgraph SINK["落库进程"]
JSONL --> CONS["cdc_sink.py<br/>幂等 upsert"]
KAFKA -. "另一条消费路径(本题未落库)" .-> CONS
CONS --> ANA["分析表<br/>labcdc.orders_analytics"]
end
SLOT -. "没人消费:WAL 一直被压住" .-> WALGROW["pg_wal 持续增长<br/>约 160 字节/行"]
PROD -. "进程崩了:槽保留位点,数据不丢" .-> SLOT
CONS -. "重放同一批:靠幂等键去重" .-> ANA
style APP fill:#e8f0fe,stroke:#3b5bdb,color:#1a1a2e
style TBL fill:#e8f0fe,stroke:#3b5bdb,color:#1a1a2e
style WAL fill:#fff4e6,stroke:#e8590c,color:#1a1a2e
style SLOT fill:#fff4e6,stroke:#e8590c,color:#1a1a2e
style PROD fill:#e6fcf5,stroke:#0ca678,color:#1a1a2e
style JSONL fill:#e6fcf5,stroke:#0ca678,color:#1a1a2e
style KAFKA fill:#e6fcf5,stroke:#0ca678,color:#1a1a2e
style CONS fill:#f3f0ff,stroke:#7048e8,color:#1a1a2e
style ANA fill:#f3f0ff,stroke:#7048e8,color:#1a1a2e
style WALGROW fill:#ffe3e3,stroke:#c92a2a,color:#1a1a2e
图上共 7 跳,跨了 4 个存储边界(表、WAL、Kafka、分析表)和 3 个进程边界(数据库、管道进程、落库进程)。
分步走一遍
① 应用写一行。 UPDATE orders SET status='paid' WHERE order_id=1001。这条语句先改数据页,再写 WAL。注意顺序:先写 WAL,再改页(Write-Ahead)。这一点决定了后面所有事——数据库里已经有一份「发生过什么」的完整记录,只是它以物理格式存在。
② WAL 里的物理记录。 表数据页会被改,WAL 里留下对应的 Heap 页改动。这时候的记录还认不出「表、列、旧值」——它记的是页和偏移。
③ 逻辑复制槽解码。 cdc_slot 挂在 WAL 上,用 test_decoding 插件把物理记录翻译成人类可读的变更:
table public.orders: UPDATE: order_id[bigint]:1001 customer_id[bigint]:1 status[text]:'paid' ...
这一步是整个管道的语义起点:从这里往下,数据不再是「页怎么变」,而是「哪张表的哪一行的哪个字段变成了什么」。
④ 管道进程解析。 cdc_pipeline.py 从槽里取变更原文,解析成结构化记录(op、key、after、事务号、LSN),给每条赋一个递增的序号 seq。序号很关键,它是后面所有幂等判断的锚。
⑤ 分发到两个下游。 同一批变更写成两种形态:本地 JSONL 文件(下游 A,最容易验证),和 Kafka topic(下游 B,最能扛量、能多消费组并行)。这一步在图上是个分叉,两条路的可靠性保证完全不同。
⑥ 消费者落库。 cdc_sink.py 读 JSONL,把每条变更 upsert 进 orders_analytics。这里有两道闸门:cdc_applied(seq) 挡住重复投递的序号,orders_analytics 上的 src_seq 比较挡住乱序到达的旧版本。
⑦ 分析表可查。 运营在报表里看到的是这一步的结果。
一次完整运行
原始输出见 lab/evidence/stream-data-pipeline/02-pipeline.txt。造一个显式事务(2 条 INSERT + 1 条 UPDATE)外加 2 条独立语句和 1 条 DELETE,解码得到 6 条变更:
{"seq": 0, "op": "INSERT", "table": "orders", "key": {"order_id": "1001"}, "after": {"status": "created", ...}, "txn_xid": "21938", "txn_size": 3}
{"seq": 1, "op": "INSERT", ... "txn_xid": "21938", "txn_size": 3}
{"seq": 2, "op": "UPDATE", ... "txn_xid": "21938", "txn_size": 3}
{"seq": 3, "op": "INSERT", ... "txn_xid": "21939", "txn_size": 1}
{"seq": 4, "op": "UPDATE", ... "txn_xid": "21940", "txn_size": 1}
{"seq": 5, "op": "DELETE", "key": {"order_id": "1002"}, "after": null, "txn_xid": "21941", "txn_size": 1}
txn_xid=21938 的三条记录共享同一个 txn_size=3,对应数据库里那个 begin ... commit。管道这边按序号顺序处理,就能保证「一个事务里的多条变更一起交出去」——它是链路上唯一还保留事务边界的地方,再往下(JSONL、Kafka 单条消息)这个边界就散了。
落库结果:
first_run applied=6 duplicate_seq_skipped=0 orders_analytics_rows=3
duplicate_replay applied=0 duplicate_seq_skipped=12 orders_analytics_rows=3
orders_analytics_rows=3 而不是 6,因为 1001 和 1003 各被改过、1002 被删了——分析表存的是主键下的最新状态,不是变更流水。想留流水,ODS 层才是它的位置(第 05 章)。
什么时候不该上管道
这不是道德题,是算术题。本机实测:灌 10 万行,逻辑槽未消费时积压 16,254,472 字节 WAL(lab/evidence/stream-data-pipeline/01-cdc-internals.txt),摊到每行约 160 字节。这个数字是你判断成本的基础。
- 每天 1 亿行更新 × 160 字节 ≈ 16 GB/天的额外 WAL 留存——只要槽没有被及时消费,这部分就一直占着主库磁盘。如果你的业务是「夜间批量刷新报表」,把 CDC 全套架上去,换来的是运维复杂度,不是业务价值。
- 反过来,如果下游要按分钟级做风控、库存扣减、推荐特征,或者业务方明确要求「删除也要能被感知」(轮询看不到 DELETE 之外的中间态),那管道就是唯一答案。
中间地带最常见的做法是混合:主链路走 CDC,历史回填和纠错走一次性批量重跑。第 05 章会讲回填怎么复用 ODS 里的顺序字段。
几处容易想反的地方
| 常见理解 | 实际机制 | 结论 |
|---|---|---|
CDC 就是「定时查增量表 / 查 updated_at」 |
那是轮询:两次查询之间被改回去的值看不到,DELETE 只能靠软删标记 | 要拿到完整变更序列,只能读日志 |
| 复制槽就是「订阅」,订阅了就自动收到 | 槽只负责记位点、留 WAL,没人来取就干等 | 槽是负债也是保险,必须有人消费 |
| 分析表 3 行 ≠ 处理了 3 条变更 | 分析表按主键 upsert,变更流水在 ODS | 先想清楚你要的是「状态」还是「事件」 |
| 管道断了要重新导数据 | 槽保留了位点,重启后从断点续读 | 断点续传是槽的核心价值,前提是槽没被删 |
生产边界
- 教学替身 vs 真实依赖:本机的「管道进程」是一个一次性脚本,跑完就退出;真实系统里它是常驻服务,要有健康检查、指标上报和自动重启。真实环境里 Kafka 至少 3 副本,本机集群 B 是单节点单副本,它挂掉等于下游 B 直接不可用。
- 要盯的指标:复制槽的
confirmed_flush_lsn与当前 LSN 的距离、pg_wal目录大小、管道进程的消费速率与槽产出速率的差、分析表的写入延迟。 - 失败策略:槽积压超阈值要能告警并自动暂停上游大批量写入;落库失败要停在当前
seq而不是跳过——跳过会在分析表里留下一个永远不会被修正的旧值。
动手
- 按 01 章 建一个复制槽,插入 5 万行,记录
pg_wal_lsn_diff(pg_current_wal_lsn(), restart_lsn)的字节数,除以行数得到你机器上的「每行 WAL 成本」。用完立刻pg_drop_replication_slot。 - 跑一遍
exp02_pipeline.sh(实验脚本在lab/exp/stream-data-pipeline/),确认orders_analytics里恰好 3 行、cdc_applied里 6 行。 - 断言题:把 JSONL 复制成两份喂给 sink,确认
duplicate_seq_skipped等于输入条数,而analytics_fingerprint_md5与第一次完全相同。
自测
- 这条链路跨了几个进程、几块存储?其中哪些是「数据不能丢」的关键节点?
- 一个事务里改了 3 行,管道如果不保持事务边界会出什么问题?哪一段负责保住这个边界?
- 为什么分析表 6 条变更只落成 3 行?如果要保留 6 条,应该改哪里?
- 假设某业务每天产生 5000 万行更新,你凭什么数据判断该不该上 CDC?请给出算式。
- 「数据卡住了」这个现象,怎么用链路图把排查范围从 7 跳缩到 2 跳?
↓ 下一步:01 章 · CDC 的原理与实现 —— 把第 ③ 跳放大:复制槽到底是什么、不消费会怎样。