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 字节。这个数字是你判断成本的基础。

中间地带最常见的做法是混合:主链路走 CDC,历史回填和纠错走一次性批量重跑。第 05 章会讲回填怎么复用 ODS 里的顺序字段。

几处容易想反的地方

常见理解 实际机制 结论
CDC 就是「定时查增量表 / 查 updated_at」 那是轮询:两次查询之间被改回去的值看不到,DELETE 只能靠软删标记 要拿到完整变更序列,只能读日志
复制槽就是「订阅」,订阅了就自动收到 槽只负责记位点、留 WAL,没人来取就干等 槽是负债也是保险,必须有人消费
分析表 3 行 ≠ 处理了 3 条变更 分析表按主键 upsert,变更流水在 ODS 先想清楚你要的是「状态」还是「事件」
管道断了要重新导数据 槽保留了位点,重启后从断点续读 断点续传是槽的核心价值,前提是槽没被删

生产边界

动手

  1. 按 01 章 建一个复制槽,插入 5 万行,记录 pg_wal_lsn_diff(pg_current_wal_lsn(), restart_lsn) 的字节数,除以行数得到你机器上的「每行 WAL 成本」。用完立刻 pg_drop_replication_slot。
  2. 跑一遍 exp02_pipeline.sh(实验脚本在 lab/exp/stream-data-pipeline/),确认 orders_analytics 里恰好 3 行、cdc_applied 里 6 行。
  3. 断言题:把 JSONL 复制成两份喂给 sink,确认 duplicate_seq_skipped 等于输入条数,而 analytics_fingerprint_md5 与第一次完全相同。

自测

  1. 这条链路跨了几个进程、几块存储?其中哪些是「数据不能丢」的关键节点?
  2. 一个事务里改了 3 行,管道如果不保持事务边界会出什么问题?哪一段负责保住这个边界?
  3. 为什么分析表 6 条变更只落成 3 行?如果要保留 6 条,应该改哪里?
  4. 假设某业务每天产生 5000 万行更新,你凭什么数据判断该不该上 CDC?请给出算式。
  5. 「数据卡住了」这个现象,怎么用链路图把排查范围从 7 跳缩到 2 跳?

↓ 下一步:01 章 · CDC 的原理与实现 —— 把第 ③ 跳放大:复制槽到底是什么、不消费会怎样。

进入 keel 阅读