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:
WAL records
-> decode relation and tuple changes
-> reorder by transaction
-> output plugin
-> logical replication protocol or SQL decoding API
-> consumeroutput plugin 决定编码,不决定下游业务语义:
| plugin / path | 典型用途 | 重要边界 |
|---|---|---|
pgoutput | 内建 publication/subscription | 二进制 replication protocol,不是给人直接读的文本 |
test_decoding | 测试和理解 logical decoding | 测试插件,不是生产业务协议 |
| 第三方 JSON/connector plugin | 通用 CDC | 版本、扩展、schema envelope、升级兼容需独立负责 |
不要把“输出是 JSON”误认为 schema compatibility 已解决。事件仍需定义:
source identity
database / schema / table
transaction and commit position
operation
replica identity before image
new column values
schema version
event id
producer version缺少 before image 时,DELETE 怎样定位?TOAST 列未变化时,connector 如何表示?列 rename 是新字段还是同一字段?这些都不是 JSON 格式本身能回答的。
SQL 解码 API 与 streaming protocol
教学或诊断可用:
SELECT *
FROM pg_create_logical_replication_slot(
'demo_slot',
'test_decoding'
);
SELECT *
FROM pg_logical_slot_peek_changes(
'demo_slot',
NULL,
100
);
SELECT *
FROM pg_logical_slot_get_changes(
'demo_slot',
NULL,
100
);peek 不消费,适合观察;get 会推进消费位置。对 pgoutput 应使用 binary API 或
replication protocol,而不是强行调用文本函数。
生产 connector 通常通过 replication protocol:
START_REPLICATION SLOT ...
-> server sends changes
-> client sends standby status updates
-> slot confirmed_flush_lsn advances“客户端已经读到”“写进本地 buffer”“写到外部系统”“外部事务提交”“向 PostgreSQL 确认”是五个不同瞬间。connector 必须声明哪个瞬间触发 acknowledgement。
snapshot 与 stream 必须无缝衔接
新建 logical slot 时可导出一个 snapshot。这个 snapshot 表示:
snapshot sees database state at boundary B
slot emits all committed changes after boundary B于是通用 bootstrap 可以:
- 创建 slot 并取得 exported snapshot;
- 在另一个事务
SET TRANSACTION SNAPSHOT; - 全量导出 snapshot 中的基线;
- 从 slot 持续消费 B 之后的变化;
- 在下游把全量与增量合并。
内建 subscription 将这套协调封装在 table sync workers 中。自行开发 CDC 时,如果先
普通 SELECT 全表、很久以后才创建 slot,会在两者间漏变化;若先创建 slot却不消费也
不预算磁盘,会让 WAL 无界增长。bootstrap 顺序必须成为协议,不是运气。
四个位置不要混为一个“offset”
一个成熟 CDC 管道至少有:
restart_lsn
publisher still needs WAL from here
confirmed_flush_lsn
source slot believes consumer confirmed through here
consumer durable checkpoint
connector can restart from here
sink commit / business watermark
downstream effects are durable through here它们应满足与具体协议一致的单调关系。最危险的情况是 connector 先确认 source,再把 buffer 异步写 sink:
confirmed_flush_lsn advances
-> connector crashes before sink commit
-> source is allowed to recycle old changes
-> permanent downstream gap反过来,sink 已提交而 checkpoint 尚未推进会造成 replay;这通常可以用幂等处理,永久 缺口却无法凭空修复。因此多数 CDC 更愿意接受 at-least-once,而不是冒险提前确认。
plugin、slot 与版本都要进 inventory
建议为每个 consumer 记录:
consumer_id: search-orders-v3
source_cluster: pg-prod
database: shop
slot: cdc_search_orders_v3
plugin: pgoutput
publication: cdc_search_orders
owner: search-platform
schema_contract: order-event-v7
checkpoint_store: kafka-connect-offsets
ack_after: sink_transaction_commit
max_replay_window: 15m
max_slot_retained_wal: 80GiB
rebootstrap_method: snapshot-plus-stream
retirement_ticket: nullslot 本身不知道 consumer owner、SLO 或 sink 状态。没有外部 inventory,inactive slot 只能告诉你“现在没人连”,不能告诉你“可以删”。
29.2.2 至少一次、重复事件与下游幂等
PostgreSQL 已经明确允许 replay
logical slot 是 crash-safe 的,但它的当前位置只在 checkpoint 时持久化。服务器崩溃后, slot 可能回到较早 LSN,于是最近变化再次发送。网络断开、consumer 在 sink commit 后但 ack 前崩溃,也会产生相同结果。
因此正确假设是:
event may be delivered more than once
transaction order is meaningful within a stream
acknowledged history can no longer be requested from that slot“测试十次没重复”不能升级为 exactly-once 保证。
exactly-once 是端到端属性
若 sink 也是 PostgreSQL,可以在同一个目标事务里同时写业务状态与 dedup ledger:
BEGIN;
INSERT INTO cdc_applied(event_id, source_lsn, applied_at)
VALUES (:event_id, :source_lsn, clock_timestamp())
ON CONFLICT (event_id) DO NOTHING;
-- 只有上一条确实插入时,才应用业务变化。
UPDATE search_projection
SET ...
WHERE ...
AND :new_event_was_inserted;
COMMIT;随后才向 source 确认。若事务已提交但 ack 丢失,replay 会命中唯一键,不重复副作用。
但如果副作用跨多个系统:
write database
send email
charge payment
publish Kafka
advance offset没有一个 PostgreSQL transaction 能原子覆盖全部。应使用:
- transactional outbox;
- sink 自身的幂等 key / conditional write;
- 去重 ledger;
- 可重建投影;
- saga/补偿;
- 对不可重复副作用的业务级 request identity。
不要用“处理成功后更新 offset”一句话掩盖多个 commit 之间的崩溃窗口。
事件身份怎么设计
只用 txid 不够:32-bit XID 会回卷,跨 cluster/database 也会重复。只用表主键也不够:
同一行可以变化很多次。一个实用 envelope 可包含:
source system identifier
database identity
slot / publication contract
commit LSN
transaction identity
change ordinal within transaction
schema contract version具体 connector 能提供哪些字段取决于协议。核心要求是:
same logical change -> same event id on replay
different logical changes -> different event ids不要把消费时间戳或随机 UUID 当 event id;replay 时它们会变。
upsert 不自动等于幂等
下面的 sink 写法:
INSERT ... ON CONFLICT (id) DO UPDATE ...只有在“最后写入覆盖即可”且顺序严格时才近似幂等。它处理不了:
- 较旧事件晚到,覆盖较新状态;
amount = amount + delta被重复执行;- DELETE/tombstone 后旧 UPDATE 复活;
- 多个 source 同写一个 key;
- 事件 schema 变化导致部分字段保留旧值;
- 外部副作用已经发生。
更安全的投影常带 source version:
INSERT INTO order_projection (
order_id, status, source_commit_lsn, payload
)
VALUES (...)
ON CONFLICT (order_id) DO UPDATE
SET status = excluded.status,
source_commit_lsn = excluded.source_commit_lsn,
payload = excluded.payload
WHERE order_projection.source_commit_lsn
< excluded.source_commit_lsn;LSN 只在同一 source timeline/合同内有序;跨 source 合并仍需业务 version/vector 或冲突 规则。
transaction boundary 不能随意打散
源事务:
debit account A
credit account B
append ledger若 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 的事件:
stop entire stream forever
skip and lose data silently
retry at full speed forever都不是完整策略。应记录:
event_id: ...
source_position: ...
schema_version: ...
error_class: conversion | constraint | permission | code
first_seen: ...
attempts: ...
raw_payload_ref: encrypted-private-location
owner: ...
decision: repair-and-replay | compensate | approved-skipdead-letter queue 只是隔离区,不是数据正确性的垃圾桶。任何 approved skip 都要进入 reconciliation,并记录业务影响。
29.2.3 槽停滞、WAL 保留与磁盘风险
inactive 不等于无害
slot 与连接生命周期独立。consumer 下线后:
active = false
confirmed_flush_lsn stops
source continues writing WAL
restart_lsn remains old
pg_wal retained bytes grow
catalog_xmin may also hold catalog tuples基础查询:
SELECT slot_name,
database,
active,
active_pid,
inactive_since,
xmin,
catalog_xmin,
restart_lsn,
confirmed_flush_lsn,
pg_size_pretty(
pg_wal_lsn_diff(pg_current_wal_lsn(), restart_lsn)
) AS retained,
wal_status,
pg_size_pretty(safe_wal_size) AS safe_wal,
invalidation_reason
FROM pg_replication_slots
WHERE slot_type = 'logical'
ORDER BY restart_lsn NULLS FIRST;pg_wal_lsn_diff(current, restart_lsn) 是按当前时刻估算该 slot 的 WAL 保留距离,不等于
磁盘上所有 WAL 文件恰好这么大;checkpoint、archive、其他 slot 与 segment 粒度都会
影响实际 pg_wal。
把风险换成时间
若近期 WAL 产生率为 (r) bytes/s,slot 当前保留 (L),可用于增长的安全空间为 (F),则最粗略的时间预算:
$$ T_{\text{disk}} \approx \frac{F}{r} $$
若配置 max_slot_wal_keep_size = M,距离 slot 可能在 checkpoint 后失去所需 WAL 的
预算:
$$ T_{\text{slot}} \approx \frac{M - L}{r} $$
实际告警应使用变化率和低水位:
retained bytes
safe_wal_size
wal_status / invalidation_reason
pg_wal filesystem free
archive health
consumer lag and error rate
catalog_xmin / XID age只告警 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 增量收敛,再执行:
ALTER SUBSCRIPTION pg36_shop_sub DISABLE;确认目标 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。
恢复:
ALTER SUBSCRIPTION pg36_shop_sub ENABLE;验收同时要求:
confirmed_flush_lsn >= source marker LSN
both relation states = r
source logical manifest = target logical manifest“slot active 又变 true”仍不够。
停滞处置顺序
1. identify exact slot and owner
2. confirm source database / plugin / consumer contract
3. capture restart, confirmed, xmin, wal_status, safe_wal_size
4. inspect consumer, network, auth, schema and apply conflicts
5. calculate disk and slot-loss deadlines
6. choose resume, repair, rebootstrap or retire
7. verify acknowledgement plus semantic convergence
8. only then close or drop the slot优先恢复 consumer,而不是先 drop slot。若 slot 已 lost,继续重试同一位点没有意义; 冻结下游写入,按 snapshot + new slot 的协议重建。若 consumer 已永久退役,保留审批与 最后消费位置后精确 drop。
与 Pigsty 观测对齐
Pigsty 的 PGSQL Replication / Persist / Instance / Alert 看板可把:
slot retention
WAL production
archive
disk free
physical replica lag
logical pub/sub
host I/O放在同一时间轴。原生 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
上一节:逻辑复制原语 · 返回本章目录 · 下一节:批量装载与数据校验 · 查看全书目录 · 查看索引中心