KEEL · 龙骨 · A CURRICULUM FOR THE AI ERA

01 · CDC 的原理与实现 — keel 龙骨

第 00 章说第 ③ 跳是「语义起点」。这一章把它放大:槽到底记了什么、不消费会怎样、UPDATE 的旧值为什么取不到。

第 00 章说第 ③ 跳是「语义起点」。这一章把它放大:槽到底记了什么、不消费会怎样、UPDATE 的旧值为什么取不到。

现场

上游 DBA 给了一个需求:把 orders 表的变更实时同步到下游。两种做法摆在桌上。

做法一,轮询:每 5 秒 select * from orders where updated_at > :last_ts。做法二,读日志:接一个逻辑复制槽,让数据库把变更解码后推给你。

轮询看着简单,但它有三个漏洞:两次查询之间被改两次的行只能看到最后一次;被物理 DELETE 的行从此消失,下游永远不知道;一个事务改了 100 行,下游看到的是「100 行都变了」,分不清它们属于同一个原子操作。这些漏洞不是实现问题,是「查表」这个动作本身的信息损失——表里只有当前状态,没有过程。

要过程,得去读数据库自己记录过程的地方。

全链路

flowchart TD
  subgraph DB["PostgreSQL 5433 / labcdc"]
    direction TB
    W["应用写入<br/>INSERT/UPDATE/DELETE public.orders"] --> HP["数据页<br/>orders (id PK)"]
    W --> WR["WAL 记录<br/>page-level 改动"]
    WR --> DEC["逻辑解码<br/>插件 test_decoding"]
    DEC --> S["复制槽 cdc_slot<br/>记 restart_lsn / confirmed_flush_lsn"]
  end

  subgraph OUT["消费侧"]
    S -->|"pg_logical_slot_get_changes"| Q["SQL 函数取变更<br/>一次性、取完即确认"]
    S -->|"pg_recvlogical --start"| R["流式接收<br/>持续、按事务推"]
  end

  subgraph RISK["不消费时的行为"]
    S -. "无人确认:WAL 不能回收" .-> G1["restart_lsn 停在原地<br/>积压随写入增长"]
    S -. "长期残留:磁盘被写满" .-> G2["主库停止接受写入<br/>(100k 行 ≈ 16.25 MB)"]
  end

  style W fill:#e8f0fe,stroke:#3b5bdb,color:#1a1a2e
  style HP fill:#e8f0fe,stroke:#3b5bdb,color:#1a1a2e
  style WR fill:#fff4e6,stroke:#e8590c,color:#1a1a2e
  style DEC fill:#fff4e6,stroke:#e8590c,color:#1a1a2e
  style S fill:#fff4e6,stroke:#e8590c,color:#1a1a2e
  style Q fill:#e6fcf5,stroke:#0ca678,color:#1a1a2e
  style R fill:#e6fcf5,stroke:#0ca678,color:#1a1a2e
  style G1 fill:#ffe3e3,stroke:#c92a2a,color:#1a1a2e
  style G2 fill:#ffe3e3,stroke:#c92a2a,color:#1a1a2e

图上三条边要从左看到右:正常路径(实线)是「写 → WAL → 解码 → 消费」;两条虚线是同一个机制失效后的两种表现。

精确定义

复制槽(replication slot) 不是订阅,也不是缓冲区。它只在主库上持久化两件事:这个下游最多还需要从哪里开始的 WAL(restart_lsn),以及下游已经确认消费到哪里(confirmed_flush_lsn)。主库据此决定哪些 WAL 段还不能删。

confirmed_flush_lsn 由消费方推进。谁来取一次变更、确认一次,它就往前走。

restart_lsn 才是决定磁盘的那个值——它是最老的一段「因为可能有下游要用,所以不能删」的 WAL。它不由消费动作直接推进,这条在本机跑出来一个很反直觉的结果,下面单独说。

REPLICA IDENTITY 是表级属性,决定 UPDATE/DELETE 时日志里带多少「行原来的样子」。默认是 DEFAULT,只带主键。

一次完整运行:解码原文与事务边界

原始输出见 lab/evidence/stream-data-pipeline/01-cdc-internals.txt。先看一个显式事务解出来是什么样:

    lsn     |  xid  | data
------------+-------+--------------------------------------------------
 0/4E66D878 | 21878 | BEGIN 21878
 0/4E66D878 | 21878 | table public.orders: INSERT: order_id[bigint]:1 customer_id[bigint]:100 status[text]:'created' ...
 0/4E66D978 | 21878 | table public.orders: INSERT: order_id[bigint]:2 ...
 0/4E66DA18 | 21878 | table public.orders: UPDATE: order_id[bigint]:1 ... status[text]:'paid' ...
 0/4E66DA88 | 21878 | table public.orders: DELETE: order_id[bigint]:2
 0/4E66DB00 | 21878 | COMMIT 21878

一个事务里的 4 条变更被 BEGIN / COMMIT 包在一起,xid 都是 21878。这个边界是解码器给的,不是数据库表给的——表里你永远看不出这 4 条改动是同一个原子操作。

换成两条独立语句:

 0/4E66DB38 | 21879 | BEGIN 21879
 0/4E66DB38 | 21879 | table public.orders: INSERT: order_id[bigint]:3 ...
 0/4E66DC08 | 21879 | COMMIT 21879
 0/4E66DC08 | 21880 | BEGIN 21880
 0/4E66DC08 | 21880 | table public.orders: UPDATE: order_id[bigint]:3 ...
 0/4E66DCA8 | 21880 | COMMIT 21880

两条语句各自成事务,各自 BEGIN/COMMIT。管道要保住事务语义,就得自己按 xid 分组——第 00 章的 txn_xid / txn_size 字段干的就是这件事。

不消费会怎样:WAL 堆积量化

建好 cdc_slot 后不取任何变更,直接灌 10 万行:

写之前: restart_lag_bytes=1184     confirmed_lag_bytes=0
写之后: restart_lag_bytes=16255656 confirmed_lag_bytes=16254472

16,254,472 字节 ≈ 16.25 MB,摊到 10 万行是每行约 162 字节。这就是 CDC 的隐藏成本:不是网络,不是 CPU,是主库上被迫留住的 WAL。一个被遗忘的槽会让 pg_wal 一直涨,直到撑满磁盘,然后主库拒绝写入——这是生产上排得上号的故障类型。

restart_lsn 什么时候才前进

把变更取走之后,confirmed_flush_lsn 立刻追平当前 LSN(confirmed_lag_bytes 归零)。但真正决定磁盘的 restart_lsn 不动,这是本机最有价值的一次实测:

基线(刚建槽)        restart_lag=56
灌 1 万行、不消费      restart_lag=1545288
取走变更             confirmed_lag 归零,restart_lag 仍是 1545288

随后做两组对照。第一组,取走变更后、每轮之间继续产生新 WAL 再 checkpoint:

round 1: restart_lsn=0/595868F8  confirmed=0/596FFE08  restart_lag=1545664
round 2: restart_lsn=0/596FFEB8  confirmed=0/5970A210  restart_lag=41992
round 3: restart_lsn=0/5970A210  confirmed=0/5970BED0  restart_lag=7536
round 4: restart_lsn=0/5970BED0  confirmed=0/5970DBF8  restart_lag=7640
round 5: restart_lsn=0/5970F970  confirmed=0/5970F9A8  restart_lag=232

第二轮 restart_lsn 恰好等于第一轮的 confirmed_flush_lsn——它的推进滞后一个 checkpoint 周期。第二组,取走变更后不再产生新 WAL、连跑三次 checkpoint:

idle-checkpoint 1: restart_lsn=0/5970FA58 restart_lag=1553192
idle-checkpoint 2: restart_lsn=0/5970FA58 restart_lag=1553368
idle-checkpoint 3: restart_lsn=0/5970FA58 restart_lag=1553544

restart_lsn 纹丝不动,1.55 MB 一直压着。两组对照的差别说明:「消费了」只是必要条件,WAL 真正被放开还要等后续 checkpoint 把它往前推。对运维的直接含义是——一条停止消费的管道,即使你立刻把积压读完,磁盘也不会马上就还给你;反过来,一个没人管的槽,永远不会自己好。

UPDATE 的 old key 与 REPLICA IDENTITY

默认属性下(relreplident='d'),一次只改非主键列的 UPDATE:

table public.orders: UPDATE: order_id[bigint]:1 customer_id[bigint]:100 status[text]:'paid' ...

只有新值,没有旧值。下游拿到的是一条「当前状态快照」,不是「状态迁移」。想做「按状态变化做增量计数」,你只能记住上一版自己算差值。

改成 REPLICA IDENTITY FULL(relreplident='f')再改一次:

table public.orders: UPDATE: old-key: order_id[bigint]:1 customer_id[bigint]:100 status[text]:'paid' ... new-tuple: order_id[bigint]:1 customer_id[bigint]:999 status[text]:'refunded' ...

old-key: 和 new-tuple: 都在了。代价是每行 UPDATE 的 WAL 体积变大(旧值也要落盘)。默认只带主键的取舍很清楚:主键不变时,下游靠主键就能定位到行,旧值只是冗余;但一旦你需要「这一行是从什么变成什么的」,就得显式开 FULL,并接受 WAL 放大。

生产边界

动手

  1. 建槽 cdc_slot,跑一遍 [3a] 那个多语句事务,用 pg_logical_slot_get_changes 打印原文,数一数 BEGIN/COMMIT 之间有几条。
  2. 灌 5 万行、不消费,记录 pg_wal_lsn_diff(pg_current_wal_lsn(), restart_lsn);再取走变更、做一次 checkpoint,看这个值降了多少。实验结束务必 select pg_drop_replication_slot('cdc_slot');。
  3. alter table orders replica identity full; 前后各做一次非主键列的 UPDATE,比较解码输出里多了什么字段,并解释它为什么会变大。

自测

  1. 轮询查表会丢掉哪三类信息?分别对应数据库里的什么机制?
  2. confirmed_flush_lsn 和 restart_lsn 各自由谁推进、各自决定什么?本机实验里哪一个不会随消费立刻前进?
  3. 为什么一个「建了但没人消费」的槽是纯负债?它是怎么把一个只读的日志变成主库故障的?
  4. REPLICA IDENTITY FULL 换来了旧值,代价是什么?什么场景下必须开,什么场景下开了是浪费?
  5. 你在监控里看到某个槽的 confirmed_flush_lsn 正常前进,但 pg_wal 还在涨。可能的原因是什么?

↓ 下一步:02 章 · 管道拓扑与 Kafka Connect —— 把第 ④ 跳放大:变更交给谁、状态记在哪、断了怎么续。

进入 keel 阅读