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 一律不同步。加列、改类型、删列都要人工在两端执行,且顺序有讲究。推荐的操作顺序是:

  1. 下游先做向后兼容的 schema 变更(加列用 nullable,禁止加 NOT NULL 默认值)
  2. 上线新的上游写入逻辑
  3. 回填历史数据
  4. 下游切换读取新列
  5. 再考虑收紧约束

配合 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 语义、故障恢复——都留给了使用者。

落到具体的生产实践,我的清单是:

  1. 槽位监控必须进核心告警,配 max_slot_wal_keep_size 兜底;
  2. REPLICA IDENTITY 按表逐个确认,别指望默认值能给你完整的 before image;
  3. 手写消费客户端一定实现保活反馈;
  4. 上游禁 DDL 自动传播,走人工评审 + Schema Registry 校验;
  5. 拆分大事务,开启 streaming 解码;
  6. failover 预案提前演练,别等到真切换才发现槽位没了。

CDC 从来不是一个「配好了就不用管」的组件,它是数据栈里那条最细却最关键的血管。理解它的边界,比理解它的 API 重要得多。

点赞(0) 打赏

评论列表 共有 0 条评论

暂无评论
立即
投稿

微信公众账号

微信扫一扫加关注

发表
评论
返回
顶部