29.2 CDC 与复制槽治理
内建 logical replication 的目标是 PostgreSQL table,CDC 的目标可能是 Kafka、 对象存储、搜索引擎、缓存、数据仓库或另一个业务服务。两者都从 logical decoding 出发,却在“谁保存位点、何时确认、如何处理重复、怎样表达 schema”上走向不同系统。
本节不绑定某一个 connector,而是建立任何 CDC 实现都绕不开的合同。
29.2.1 逻辑解码、插件与消费位点
从物理 WAL 到业务事件
PostgreSQL WAL 记录的是恢复数据库所需的物理/逻辑底层信息,并不是一个稳定的 JSON 业务接口。logical decoding 负责把一个 database 的持久表变化重组为按事务提交边界 可解释的 change stream:
output plugin 决定编码,不决定下游业务语义:
| plugin / path | 典型用途 | 重要边界 |
|---|---|---|
pgoutput |
内建 publication/subscription | 二进制 replication protocol,不是给人直接读的文本 |
test_decoding |
测试和理解 logical decoding | 测试插件,不是生产业务协议 |
| 第三方 JSON/connector plugin | 通用 CDC | 版本、扩展、schema envelope、升级兼容需独立负责 |
不要把“输出是 JSON”误认为 schema compatibility 已解决。事件仍需定义:
缺少 before image 时,DELETE 怎样定位?TOAST 列未变化时,connector 如何表示?列 rename 是新字段还是同一字段?这些都不是 JSON 格式本身能回答的。
SQL 解码 API 与 streaming protocol
教学或诊断可用:
peek 不消费,适合观察;get 会推进消费位置。对 pgoutput 应使用 binary API 或
replication protocol,而不是强行调用文本函数。
生产 connector 通常通过 replication protocol:
“客户端已经读到”“写进本地 buffer”“写到外部系统”“外部事务提交”“向 PostgreSQL 确认”是五个不同瞬间。connector 必须声明哪个瞬间触发 acknowledgement。
snapshot 与 stream 必须无缝衔接
新建 logical slot 时可导出一个 snapshot。这个 snapshot 表示:
于是通用 bootstrap 可以:
- 创建 slot 并取得 exported snapshot;
- 在另一个事务
SET TRANSACTION SNAPSHOT; - 全量导出 snapshot 中的基线;
- 从 slot 持续消费 B 之后的变化;
- 在下游把全量与增量合并。
内建 subscription 将这套协调封装在 table sync workers 中。自行开发 CDC 时,如果先
普通 SELECT 全表、很久以后才创建 slot,会在两者间漏变化;若先创建 slot却不消费也
不预算磁盘,会让 WAL 无界增长。bootstrap 顺序必须成为协议,不是运气。
四个位置不要混为一个“offset”
一个成熟 CDC 管道至少有:
它们应满足与具体协议一致的单调关系。最危险的情况是 connector 先确认 source,再把 buffer 异步写 sink:
反过来,sink 已提交而 checkpoint 尚未推进会造成 replay;这通常可以用幂等处理,永久 缺口却无法凭空修复。因此多数 CDC 更愿意接受 at-least-once,而不是冒险提前确认。
plugin、slot 与版本都要进 inventory
建议为每个 consumer 记录:
slot 本身不知道 consumer owner、SLO 或 sink 状态。没有外部 inventory,inactive slot 只能告诉你“现在没人连”,不能告诉你“可以删”。
29.2.2 至少一次、重复事件与下游幂等
PostgreSQL 已经明确允许 replay
logical slot 是 crash-safe 的,但它的当前位置只在 checkpoint 时持久化。服务器崩溃后, slot 可能回到较早 LSN,于是最近变化再次发送。网络断开、consumer 在 sink commit 后但 ack 前崩溃,也会产生相同结果。
因此正确假设是:
“测试十次没重复”不能升级为 exactly-once 保证。
exactly-once 是端到端属性
若 sink 也是 PostgreSQL,可以在同一个目标事务里同时写业务状态与 dedup ledger:
随后才向 source 确认。若事务已提交但 ack 丢失,replay 会命中唯一键,不重复副作用。
但如果副作用跨多个系统:
没有一个 PostgreSQL transaction 能原子覆盖全部。应使用:
- transactional outbox;
- sink 自身的幂等 key / conditional write;
- 去重 ledger;
- 可重建投影;
- saga/补偿;
- 对不可重复副作用的业务级 request identity。
不要用“处理成功后更新 offset”一句话掩盖多个 commit 之间的崩溃窗口。
事件身份怎么设计
只用 txid 不够:32-bit XID 会回卷,跨 cluster/database 也会重复。只用表主键也不够:
同一行可以变化很多次。一个实用 envelope 可包含:
具体 connector 能提供哪些字段取决于协议。核心要求是:
不要把消费时间戳或随机 UUID 当 event id;replay 时它们会变。
upsert 不自动等于幂等
下面的 sink 写法:
只有在“最后写入覆盖即可”且顺序严格时才近似幂等。它处理不了:
- 较旧事件晚到,覆盖较新状态;
amount = amount + delta被重复执行;- DELETE/tombstone 后旧 UPDATE 复活;
- 多个 source 同写一个 key;
- 事件 schema 变化导致部分字段保留旧值;
- 外部副作用已经发生。
更安全的投影常带 source version:
LSN 只在同一 source timeline/合同内有序;跨 source 合并仍需业务 version/vector 或冲突 规则。
transaction boundary 不能随意打散
源事务:
若 connector 把三条 row event 分别确认并让下游实时可见,中途失败可能暴露不平衡状态。 CDC envelope 应保留 begin/commit 或 transaction grouping;sink 要么原子应用整个事务, 要么明确只提供最终一致读模型并隐藏未完成 batch。
大事务还会带来另一组选择:
streaming = off:源端完整解码后发送,内存/延迟风险;streaming = on:未提交变化先写 subscriber 临时文件,commit 后应用;streaming = parallel:有 worker 时直接并行 apply,否则回退临时文件。
吞吐设置不能改变“只有源端 commit 后才把结果视为提交”的语义。
poison event 需要隔离,不是静默跳过
遇到无法解析或不满足 sink constraint 的事件:
都不是完整策略。应记录:
dead-letter queue 只是隔离区,不是数据正确性的垃圾桶。任何 approved skip 都要进入 reconciliation,并记录业务影响。
29.2.3 槽停滞、WAL 保留与磁盘风险
inactive 不等于无害
slot 与连接生命周期独立。consumer 下线后:
基础查询:
pg_wal_lsn_diff(current, restart_lsn) 是按当前时刻估算该 slot 的 WAL 保留距离,不等于
磁盘上所有 WAL 文件恰好这么大;checkpoint、archive、其他 slot 与 segment 粒度都会
影响实际 pg_wal。
把风险换成时间
若近期 WAL 产生率为 bytes/s,slot 当前保留 ,可用于增长的安全空间为 ,则最粗略的时间预算:
若配置 max_slot_wal_keep_size = M,距离 slot 可能在 checkpoint 后失去所需 WAL 的
预算:
实际告警应使用变化率和低水位:
只告警 active = false 会在计划维护时产生噪声;只告警磁盘使用率会在 WAL 已快速增长时
太晚。
上限保护的是 source,不保证 consumer 可恢复
max_slot_wal_keep_size 为非负值时,checkpoint 可以允许回收超出上限的 WAL,slot
可能进入 unreserved 乃至 lost。这避免一个遗忘 consumer 无限填满源盘,但代价是
consumer 必须 rebootstrap。
PostgreSQL 18 还提供 idle_replication_slot_timeout:slot inactive 超过时限后可在
checkpoint 被 invalidated。它同样不是“暂停后自动保存到对象存储”。启用前必须明确:
- 哪类 slot 允许因空闲失效;
- 谁收到即将失效告警;
- snapshot/rebootstrap 需要多久;
- 同步到 standby 的 slot 等豁免/特殊语义;
- checkpoint 周期带来的执行延迟。
不要为了保护磁盘把上限设得小于正常故障恢复窗口,然后把频繁 lost 当 consumer 问题。
本章停滞实验
正式实验先让 20,000 初始订单和 500/200/100 增量收敛,再执行:
确认目标 subenabled = false、worker 为零、源端 slot active = false 后,在源端插入
3,000 行有界 payload。
结果:
| 字段 | 写入前 | 写入后 |
|---|---|---|
active |
false | false |
confirmed_flush_lsn |
1/9FDBBF30 |
1/9FDBBF30 |
| retained bytes | 227,008 | 2,867,128 |
wal_status |
reserved |
reserved |
这里能确认因 consumer 停滞,确认位点没有推进且保留距离增长。不能用这 2.87 MB 推导 生产每 3,000 行的固定 WAL 成本:其他 database、full-page write、索引、payload、 checkpoint 和并发都会改变 WAL。
恢复:
验收同时要求:
“slot active 又变 true”仍不够。
停滞处置顺序
优先恢复 consumer,而不是先 drop slot。若 slot 已 lost,继续重试同一位点没有意义; 冻结下游写入,按 snapshot + new slot 的协议重建。若 consumer 已永久退役,保留审批与 最后消费位置后精确 drop。
与 Pigsty 观测对齐
Pigsty 的 PGSQL Replication / Persist / Instance / Alert 看板可把:
放在同一时间轴。原生 SQL 则确认 slot、subscription、table state 与 conflict identity。 平台看板用于发现趋势,不能替代 consumer owner 和业务 reconciliation。
进一步阅读:
- PostgreSQL 18:Logical Decoding Concepts
- PostgreSQL 18:Replication Settings
- PostgreSQL 18:
pg_replication_slots - PostgreSQL 18:Logical Replication Monitoring
- Pigsty:PGSQL Dashboards
上一节:逻辑复制原语 · 返回本章目录 · 下一节:批量装载与数据校验 · 查看全书目录 · 查看索引中心