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 放大。
生产边界
- 教学替身 vs 真实依赖:本机的消费端是
pg_logical_slot_get_changes(一次性、取完即确认)和pg_recvlogical(流式)。真实系统里通常是 Debezium 一类工具,它替你管槽、做断点、做 schema 演化——但槽的物理约束一条都不会少。 - 要盯的指标:每个槽的
pg_wal_lsn_diff(pg_current_wal_lsn(), restart_lsn)、active是否长期为f、pg_wal目录大小、以及max_replication_slots的剩余额度。confirmed_flush_lsn落后不代表磁盘紧张,restart_lsn落后才是。 - 失败策略:槽长时间
active=f且积压上升,应该先让上游降速或停掉大批量写入,再去修下游;直接pg_drop_replication_slot虽然立刻释放磁盘,但会让下游从断点彻底接不上,只能全量重来。 - DDL 不会自动同步:加列、改类型都要在两端手工协调,顺序错了解码会直接中断。
动手
- 建槽
cdc_slot,跑一遍[3a]那个多语句事务,用pg_logical_slot_get_changes打印原文,数一数BEGIN/COMMIT之间有几条。 - 灌 5 万行、不消费,记录
pg_wal_lsn_diff(pg_current_wal_lsn(), restart_lsn);再取走变更、做一次checkpoint,看这个值降了多少。实验结束务必select pg_drop_replication_slot('cdc_slot');。 alter table orders replica identity full;前后各做一次非主键列的 UPDATE,比较解码输出里多了什么字段,并解释它为什么会变大。
自测
- 轮询查表会丢掉哪三类信息?分别对应数据库里的什么机制?
confirmed_flush_lsn和restart_lsn各自由谁推进、各自决定什么?本机实验里哪一个不会随消费立刻前进?- 为什么一个「建了但没人消费」的槽是纯负债?它是怎么把一个只读的日志变成主库故障的?
REPLICA IDENTITY FULL换来了旧值,代价是什么?什么场景下必须开,什么场景下开了是浪费?- 你在监控里看到某个槽的
confirmed_flush_lsn正常前进,但pg_wal还在涨。可能的原因是什么?
↓ 下一步:02 章 · 管道拓扑与 Kafka Connect —— 把第 ④ 跳放大:变更交给谁、状态记在哪、断了怎么续。