PostgreSQL 逻辑复制与 CDC 深度实战:从 WAL 解码到实时数据管道的工程边界
如果你所在的团队还在用「凌晨三点 SELECT * FROM orders WHERE updated_at > ? 全量拉一遍」的方式同步数据,那你大概率已经踩过这套方案的三个坑:删掉的行永远同步不过去、错过了中间态的一条记录更新、以及在源库上开了一个让 DBA 想拉你进小黑屋的长查询。
变更数据捕获(Change Data Capture,CDC)的价值就在于此:它不问「现在表里有什么」,而是问「从上次同步到现在,发生了什么」。而在 PostgreSQL 生态里,CDC 的地基几乎清一色建立在 逻辑复制(Logical Replication) 之上。
这篇文章不打算复述文档。我们聊工程:WAL 里到底躺着什么、逻辑解码插件怎么把二进制翻译成结构化变更、复制槽为什么是生产事故的头号制造者、以及什么场景下你其实不该上逻辑复制。
一、先搞清楚:物理复制和逻辑复制差在哪
PostgreSQL 的复制分两层,很多人含糊地统称为「主从」,但它们完全不是一回事。
| 维度 | 物理复制(Streaming Replication) | 逻辑复制(Logical Replication) |
|---|---|---|
| 复制单位 | WAL 字节流 | 行级逻辑变更(INSERT/UPDATE/DELETE) |
| 版本要求 | 主备大版本必须一致 | 跨大版本可用(PG10 → PG16 可行) |
| 目标库状态 | 只读备机 | 可读可写,可继续向下游传播 |
| 粒度 | 整个实例 | 按表(Publication)选择 |
| DDL | 自动同步 | 不同步,需人工处理 |
| 典型用途 | HA、灾备、读写分离 | CDC、数据集成、灰度迁移、双写对比 |
一句话总结:物理复制复制的是「磁盘块的结果」,逻辑复制复制的是「语义上的动作」。前者要求两端完全一致,后者给了你自由裁量的空间——也给了你踩坑的空间。
二、WAL 里到底装了什么
要理解逻辑解码,必须先理解 WAL(Write-Ahead Log)。它是 PostgreSQL 的 redo log:任何数据页修改落盘之前,对应的变更记录必须先落到 WAL。
每条 WAL record 大致长这样:
+----------+--------+---------+----------+-------------+
| xl_tot_ | xid | xl_prev | xl_info | xl_rmid | <- 头部 XLogRecord
| len | | | | |
+----------+--------+---------+----------+-------------+
| RelFileNode / BlockNumber / ... 块级定位信息 |
+----------+--------+---------+----------+-------------+
| 具体数据:insert 的 tuple、update 的 new tuple 等 |
+-------------------------------------------------------+
对物理复制而言,这些记录足够——备机按顺序回放,逐字节还原。但对 CDC 来说,这些信息是「只可意会」的:我们知道 BlockNumber=4712 的页面第 3 个 tuple 变了,但不知道它对应哪张表的哪一行、哪一列。
关键开关是 wal_level:
# postgresql.conf
wal_level = logical # minimal < replica < logical
max_replication_slots = 10
max_wal_senders = 10
只有 wal_level = logical 时,WAL 里才会额外记录足够的信息用于重建成完整的元组映像。minimal 甚至不足以支撑物理 archive。注意:改这个参数需要重启,不是 reload。
另一个常被忽略的参数是 REPLICA IDENTITY,它决定了 UPDATE/DELETE 事件里能不能拿到「变更前是什么」:
-- 默认值:只记录主键。UPDATE 的 old tuple 里非主键列为 NULL
ALTER TABLE orders REPLICA IDENTITY DEFAULT;
-- 记录所有列的旧值,能拿到完整 before image,但 WAL 体积膨胀
ALTER TABLE orders REPLICA IDENTITY FULL;
-- 用某个唯一索引作为行标识(推荐,性价比最高)
CREATE UNIQUE INDEX idx_orders_biz ON orders(order_no);
ALTER TABLE orders REPLICA IDENTITY USING INDEX idx_orders_biz;
这是 CDC 里最常见的「事故」之一:下游收到 UPDATE 事件却发现除主键外全是 null,于是以为数据没变化。真相是没有设置 REPLICA IDENTITY FULL(或合适的索引)。代价必须讲清楚:FULL 模式下 UPDATE 要在 WAL 里写一份完整的旧元组,写放大明显,高频更新表慎用。
三、复制槽:既是救命稻草,也是定时炸弹
逻辑解码通过 复制槽(Replication Slot) 维护消费位点。它记录三个关键 LSN:
SELECT slot_name, plugin, slot_type, database, active,
restart_lsn, confirmed_flush_lsn,
pg_size_pretty(pg_wal_lsn_diff(pg_current_wal_lsn(), restart_lsn)) AS retained_wal
FROM pg_replication_slots;
restart_lsn:解码器还需要从这里开始读confirmed_flush_lsn:消费者已确认处理到的位置- 两者之差以外的 WAL,主库永远不敢删
这就是那个经典生产事故:有人创建了复制槽,消费者进程挂了或者根本没启动,active = false,主库的 pg_wal 目录开始线性膨胀,直到磁盘 100%。数据库不是被流量打垮的,是被一个没人认领的槽撑死的。
必须配的监控告警:
-- 危险槽位:inactive 且保留了大量 WAL
SELECT slot_name,
pg_size_pretty(pg_wal_lsn_diff(pg_current_wal_lsn(), restart_lsn)) AS lag_bytes,
active
FROM pg_replication_slots
WHERE NOT active
AND pg_wal_lsn_diff(pg_current_wal_lsn(), restart_lsn) > 1073741824 -- 1GB
ORDER BY lag_bytes DESC;
配套参数建议设置 max_slot_wal_keep_size(PG13+),给槽位一个保留上限。超过后槽位会被标记为不可用——宁可让 CDC 链路断掉重建,也不能让主库磁盘写满。这是一个刻意的「舍卒保车」设计,务必启用。
四、手搓一条 CDC 管道:Publication 到 Kafka
4.1 源库侧:只暴露该暴露的
不要图省事 CREATE PUBLICATION p FOR ALL TABLES。它会同步后续所有新建的表,某天有人建了一张几十 GB 的日志表,你的 Kafka 就炸了。
-- 显式列白名单,并对 UPDATE DELETE 做列过滤
CREATE PUBLICATION cdc_pub FOR TABLE
orders (id, order_no, user_id, amount, status, updated_at),
order_items (id, order_id, sku_id, qty, price) WITH (publish_via_partition_root = true);
-- 只同步 INSERT 和 UPDATE,不同步 DELETE(很多数仓场景想保留墓碑由上层处理)
ALTER PUBLICATION cdc_pub SET (publish = 'insert,update');
publish_via_partition_root = true 这个参数对于分区表非常关键:默认 false 时,每个分区的变更会以分区自身的表名发出,下游会看到 orders_2026_09、orders_2026_10 一堆表;设为 true 后统一以父表名 orders 发布,下游 schema 才稳定。
4.2 Debezium 侧:波动更大的是 outbound 而不是 inbound
{
"name": "pg-orders-connector",
"config": {
"connector.class": "io.debezium.connector.postgresql.PostgresConnector",
"database.hostname": "pg-primary.internal",
"database.port": "5432",
"database.user": "cdc_user",
"database.password": "${file:/secrets/pg:cdc_pwd}",
"database.dbname": "shop",
"topic.prefix": "cdc.shop",
"publication.name": "cdc_pub",
"slot.name": "debezium_orders",
"plugin.name": "pgoutput",
"snapshot.mode": "no_data",
"tombstones.on.delete": "true",
"transforms": "unwrap",
"transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState",
"transforms.unwrap.delete.handling.mode": "rewrite",
"max.batch.size": 2048,
"poll.interval.ms": 100
}
}
几个容易被忽略的实战点:
plugin.name用pgoutput(PG10+ 内置),不推荐再用wal2json/decoderbufs。前者是官方支持的协议插件,兼容性最好且无需额外安装扩展。snapshot.mode选择很关键。initial会加锁并全量快照,大表上会阻塞(虽然 PG 用的是可重复读快照但不阻塞 DML,但快照期间 WAL 不能回收)。生产大表建议no_data+ 独立的离线回填任务,把回填和增量解耦。tombstones.on.delete = true让删除事件后面跟一个 null value 的消息,这是 Kafka compaction 能真正压缩掉旧值的前提。
4.3 直接用 Python 消费 raw 流(轻量场景)
不是所有场景都值得上 Kafka Connect。几十张表以下可以直连:
import select
import psycopg
conn = psycopg.connect(
"host=pg-primary dbname=shop user=cdc_user replication=true",
autocommit=True,
)
cur = conn.cursor()
cur.start_copy_expert(
"START_REPLICATION SLOT cdc_slot LOGICAL 0/1A2B3C8 "
"(proto_version '1', publication_names 'cdc_pub')"
)
keepalive = psycopg.replication.KeepAlive(conn)
while True:
msg = cur.read_message()
if msg is None:
# 没有数据时:回送 standby status update,推进位点并保活
now = psycopg.replication.now()
feedback = psycopg.replication.StandbyStatusUpdate(msg_lsn=now)
cur.send_feedback(feedback)
if select.select([conn], [], [], 10)[0]:
cur.consume_stream()
continue
payload = msg.payload.decode()
if payload[0] in ("I", "U", "D"):
handle_change(payload)
cur.send_feedback(flush_lsn=msg.data_start)
这里有个只要自己写客户端就一定会遇到的坑:必须周期性调用 send_feedback。逻辑复制流在 wal_sender_timeout(默认 60s)内没有收到消费者的位点反馈,主库会直接断开连接。很多手写客户端在消费慢或者有空闲期时静默掉线,就是因为漏了这段代码。
五、那些文档不会告诉你的边界
5.1 DDL 不会被复制
逻辑复制只复制 DML,DDL 一律不同步。加列、改类型、删列都要人工在两端执行,且顺序有讲究。推荐的操作顺序是:
- 下游先做向后兼容的 schema 变更(加列用 nullable,禁止加 NOT NULL 默认值)
- 上线新的上游写入逻辑
- 回填历史数据
- 下游切换读取新列
- 再考虑收紧约束
配合 Schema Registry(Avro/Protobuf)做兼容性校验,可以在 connector 层面就拒绝掉破坏性变更,而不是等到半夜下游 ETL 报错。
5.2 大事务是 decode 侧的内存杀手
逻辑解码必须按事务提交顺序输出。一个更新 500 万行的事务,会在解码器里攒够 500 万条的变更才一次性吐出。这意味着:
ReorderBuffer可能占用数 GB 内存,甚至 spill 到磁盘拖垮实例- 端到端延迟出现一个巨大的尖峰,而不是均匀分摊
PG14 引入了流式解码(streaming),可以边接收边下发未提交的事务块。连接参数加 streaming 'on':
START_REPLICATION SLOT cdc_slot LOGICAL 0/0 (proto_version '1', streaming 'on', publication_names 'cdc_pub')
同时在业务侧约束批量更新:把「一个巨型 UPDATE」拆成「按 id 分批的多个事务」。
5.3 TOAST 列:看不见的默认行为
被 TOAST 压缩或行外存储的大字段(比如一个 2MB 的 JSON),如果它在这次 UPDATE 中没有发生变化,逻辑解码输出时默认会用一个占位标记替代,而不是输出真实值。这会让下游拿到不完整的数据。
如果业务确实需要每次都拿到 TOAST 列的完整值,只能把该表的 REPLICA IDENTITY 设为 FULL——但要接受随之而来的写放大。
5.4 主备切换后复制槽会消失
这是 HA 场景下最痛的一点:逻辑复制槽不会同步到备库。一旦发生 failover,新主上没有原来的槽位,CDC 链路直接断裂,只能重新做快照。
解决方案有两个:
- 用
pg_failover_slots扩展,它会把槽位的创建和推进同步到备库,实现切换后衔接; - 或者接受重建策略:把 CDC 设计成幂等 upsert + 可重放,切换后触发一次有限窗口的重新快照。
我个人倾向后者配合前者:能自己恢复的链路,比依赖恢复的链路更可靠。
六、什么时候不要上逻辑复制
讲了这么多优点,也该说说边界:
- 宽表的全表刷新型分析:如果每天要重算整表,CDC 的增量收益被全量回填吃掉,反而多一层复杂度和故障面。
- 超高频更新表:每更新 10 次就要 10 条消息,下游存储压力巨大。这种场景更适合「状态快照 + 定期物化」,而不是「变更日志」。
- 极端低延迟要求:逻辑解码 → 网络 → Kafka → 消费者,链路上的每一跳都在加延迟。要求端到端毫秒级的场景,不如在应用层直接双写或用 outbox pattern。
一个务实的判断准则:当你关心的是「变更的过程」而不是「最终的状态」,并且下游的消费位点可以被可靠地管理时,逻辑复制才是对的工具。
七、小结
PostgreSQL 的逻辑复制是一套设计得相当克制的机制:它只做「把行级变更可靠地吐出来」这一件事,其余的工程复杂度——schema 演进、delivery 语义、故障恢复——都留给了使用者。
落到具体的生产实践,我的清单是:
- 槽位监控必须进核心告警,配
max_slot_wal_keep_size兜底; REPLICA IDENTITY按表逐个确认,别指望默认值能给你完整的 before image;- 手写消费客户端一定实现保活反馈;
- 上游禁 DDL 自动传播,走人工评审 + Schema Registry 校验;
- 拆分大事务,开启 streaming 解码;
- failover 预案提前演练,别等到真切换才发现槽位没了。
CDC 从来不是一个「配好了就不用管」的组件,它是数据栈里那条最细却最关键的血管。理解它的边界,比理解它的 API 重要得多。

发表评论 取消回复