PostgreSQL 逻辑复制与 CDC:发布订阅、WAL 解码、零停机迁移与 Debezium/Kafka 集成

系统讲解 PostgreSQL 逻辑复制与变更数据捕获(CDC):逻辑复制原理与 WAL 解码(pgoutput/wal2json/decoderbufs)、发布订阅(CREATE PUBLICATION/SUBSCRIBE)配置、逻辑复制做故障切换与零停机迁移、初始同步机制、DDL 限制与数据冲突处理、与 Debezium/Kafka 的数据流集成、复制延迟监控与性能优化。

逻辑复制(Logical Replication)是 PostgreSQL 10+ 引入的生产级数据同步方案。与基于物理块复制的流复制不同,逻辑复制在逻辑层面按"行变更"传播数据,这让它成为数据分发、CDC(变更数据捕获)、跨版本迁移和多活架构的核心工具。

核心认知:物理复制复制"数据库文件的变化",逻辑复制复制"一行行的 INSERT/UPDATE/DELETE 语义"。前者适合高可用,后者适合数据集成与迁移。


一、逻辑复制原理与 WAL 解码

1.1 逻辑复制 vs 物理复制

维度物理复制(流复制)逻辑复制
复制单位数据块(block)行级变更(tuple)
跨版本要求主从大版本一致可跨大版本(协议兼容)
跨数据库/实例不行可以
数据过滤不能可按表/行/列过滤
是否复制 DDL是否(需手动处理)
用途高可用、PITR数据分发、迁移、CDC

1.2 WAL 解码机制

逻辑复制依赖 WAL(Write-Ahead Log)解码:PostgreSQL 将每个事务的变更以逻辑格式从 WAL 中解码出来,再通过网络发送给订阅端。

-- 开启逻辑复制所需的配置
ALTER SYSTEM SET wal_level = 'logical';
ALTER SYSTEM SET max_replication_slots = 10;
ALTER SYSTEM SET max_wal_senders = 10;
SELECT pg_reload_conf();

1.3 解码插件

插件输出格式典型用途
pgoutputPostgreSQL 原生协议内置发布订阅(默认)
wal2jsonJSON轻量 CDC、调试、外部消费
decoderbufsProtobufDebezium 默认(配合 pgoutput 亦可)
test_decoding文本测试/学习
-- 直接观察 WAL 解码输出(wal2json)
SELECT * FROM pg_logical_slot_get_changes('slot_demo', NULL, NULL, 'format_version', '2');
-- 输出示例:{"change":[{"kind":"insert","schema":"public","table":"orders","columnnames":["id","amount"],...}]}

1.4 逻辑复制架构

Publisher (主库)
  └── WAL 写入 → 逻辑解码 → 复制槽 (slot)
                         │  wal_sender 进程
                         ▼
                网络 (TCP)
                         │
                         ▼
Subscriber (订阅库)
  └── 接收变更 → 应用到表

二、发布订阅配置

2.1 创建发布(Publication)

-- 在源库(Publisher)上创建发布
CREATE PUBLICATION mypub FOR TABLE orders, users;

-- 发布全部表
CREATE PUBLICATION all_tables FOR ALL TABLES;

-- 只发布特定操作
CREATE PUBLICATION orders_insert_only
  FOR TABLE orders WITH (publish = 'insert');
发布选项说明
publishinsert/update/delete/truncate,默认全部
publish_via_partition_root分区表通过根表发布
FOR ALL TABLES发布全库所有表
FOR TABLE ... WHERE ...行级过滤(PG 15+)

2.2 创建订阅(Subscription)

-- 在目标库(Subscriber)上创建订阅
CREATE SUBSCRIPTION mysub
  CONNECTION 'host=primary.example.com port=5432 dbname=mydb user=repl password=xxx'
  PUBLICATION mypub
  WITH (create_slot = true, enabled = true);

关键参数:

参数说明
enabled订阅是否立即启动
create_slot是否自动在主库创建复制槽
slot_name复制槽名(默认订阅名)
synchronous_commit可设 off 加速复制消费
copy_data是否先做初始全量同步
failover订阅在故障切换后是否自动接续(PG 17+)

2.3 管理与维护

-- 查看发布与订阅
SELECT * FROM pg_publication;
SELECT * FROM pg_subscription;
SELECT * FROM pg_stat_subscription;

-- 暂停/恢复/删除
ALTER SUBSCRIPTION mysub DISABLE;
ALTER SUBSCRIPTION mysub ENABLE;
DROP SUBSCRIPTION mysub;

-- 修改订阅连接参数
ALTER SUBSCRIPTION mysub CONNECTION 'host=newhost ...';

三、初始同步

3.1 初始同步流程

创建订阅时(copy_data = true 默认),PostgreSQL 自动执行:

1. 在主库创建复制槽
2. 对发布表执行基础快照(pg_basebackup 类似,但按表)
3. 将快照数据 COPY 到订阅端
4. 快照之后的新变更通过 WAL 解码持续同步

3.2 手动初始化(已有数据的大表)

对于超大表,推荐先建发布与订阅但禁用 copy,再手工同步:

-- 目标库:创建订阅但不自动复制
CREATE SUBSCRIPTION mysub
  CONNECTION '...' PUBLICATION mypub
  WITH (copy_data = false);

-- 用 pg_dump 只导数据(不导 schema,避免覆盖订阅表)
pg_dump -h primary -d mydb --data-only -t orders -t users > data.sql
psql -d target -f data.sql

-- 重新启用订阅开始接续变更
ALTER SUBSCRIPTION mysub ENABLE;

3.3 初始同步监控

-- 观察复制状态
SELECT subname, received_lsn, latest_end_lsn, last_msg_send_time
FROM pg_stat_subscription;
-- received_lsn 滞后于 latest_end_lsn 说明还有积压

提示:初始同步期间,订阅表上的触发器默认不执行(逻辑复制默认不触发触发器,除非 REPLICA IDENTITY 与触发器配置正确)。需要 ALTER TABLE ... ENABLE TRIGGER ALL。


四、DDL 限制与冲突处理

4.1 逻辑复制的限制

限制说明应对
DDL 不复制表结构变更需手动在两端执行借助 pglogical 或变更管理流程
序列不复制SERIAL 自增值两端各自独立主键冲突时手动调序列
大对象不复制lo 类型不支持避免使用或另做同步
无主键表UPDATE/DELETE 无法定位必须设 REPLICA IDENTITY FULL
非复制表未发布表不复制发布时明确表清单
触发器默认不执行订阅端触发器不触发按需 ENABLE TRIGGER

4.2 REPLICA IDENTITY 与冲突

订阅端表中必须能定位被更新/删除的行。默认用主键,若表无主键:

-- 无主键表:让副本记录完整旧行(开销较大)
ALTER TABLE orders REPLICA IDENTITY FULL;

-- 有唯一索引但无主键:用唯一索引
ALTER TABLE orders REPLICA IDENTITY USING INDEX orders_uidx;

4.3 冲突处理策略

逻辑复制是"单向写入传播",若订阅端也允许写入,可能冲突:

冲突类型表现建议
主键冲突插入相同主键报 23505订阅端只读,或提前改 key
行不存在UPDATE/DELETE 0 行用 REPLICA IDENTITY FULL 辅助
数据不一致静默覆盖定期校验(数据比对)
-- 数据一致性校验(简单比对行数)
SELECT count(*) FROM orders ON primary;
SELECT count(*) FROM orders ON subscriber;
-- 更严格:用 checksum 按 id 分片比对

4.4 DDL 变更规范流程

1. 在订阅端先执行兼容 DDL(加列可空等)
2. 在主库执行相同 DDL
3. 若 DDL 改变行定位(删主键等),先暂停订阅
ALTER SUBSCRIPTION mysub DISABLE;
-- 两端执行 DDL
ALTER SUBSCRIPTION mysub ENABLE;

五、逻辑复制做零停机迁移

5.1 为什么逻辑复制适合迁移

  • 支持跨大版本(如 PG 13 → PG 16)
  • 支持跨平台(Linux → 云托管实例)
  • 迁移期间源库持续可写,无需停机
  • 可灰度、可回滚

5.2 零停机迁移步骤

-- 1. 在目标新库建好 schema(schema 不自动复制)
-- 2. 源库发布
CREATE PUBLICATION mig_pub FOR ALL TABLES;

-- 3. 目标库订阅(自动初始同步 + 持续变更)
CREATE SUBSCRIPTION mig_sub
  CONNECTION 'host=old-db ...' PUBLICATION mig_pub;

-- 4. 观察追赶进度,直到 received_lsn 追上主库 latest_end_lsn
SELECT received_lsn, latest_end_lsn,
       pg_wal_lsn_diff(latest_end_lsn, received_lsn) AS lag_bytes
FROM pg_stat_subscription;

5.3 切换(Cutover)流程

1. 应用改为"双写"或先切读流量到新库
2. 暂停源库写入窗口(可选)
3. 等待复制延迟归零
4. 停止订阅,新库完全接管
ALTER SUBSCRIPTION mig_sub DISABLE;
-- 修改应用连接串指向新库
5. 旧库降级为归档/只读

5.4 回滚保障

若新库有问题,可在旧库执行 ALTER SUBSCRIPTION mig_sub DISABLE; 停止订阅,并把应用连接切回旧库;修复后在低峰期重新启用订阅反向追数。零停机迁移的成败关键在于复制延迟监控和主键/唯一键无冲突,切换前务必对账。


六、与 Debezium/Kafka 集成

6.1 架构

PostgreSQL 主库
   │ WAL
   ▼
Debezium PostgreSQL Connector
   │ 解码变更事件
   ▼
Kafka Topic (orders.cdc)
   │
   ├─ 数据湖 / 数仓(实时入湖)
   ├─ 事件驱动微服务
   └─ 实时搜索/分析(Elasticsearch)

6.2 Debezium Connector 配置

# debezium postgres connector 配置片段
{
  "name": "orders-connector",
  "config": {
    "connector.class": "io.debezium.connector.postgresql.PostgresConnector",
    "database.hostname": "primary.example.com",
    "database.port": "5432",
    "database.user": "debezium",
    "database.password": "secret",
    "database.dbname": "mydb",
    "database.server.name": "pg-primary",
    "plugin.name": "pgoutput",
    "slot.name": "debezium_slot",
    "publication.name": "dbz_publication",
    "table.include.list": "public.orders,public.users",
    "heartbeat.interval.ms": "5000",
    "topic.prefix": "cdc"
  }
}

6.3 消费变更事件

# 消费 Kafka 中的 CDC 事件(orders 表的 UPDATE)
kafka-console-consumer.sh \
  --bootstrap-server kafka:9092 \
  --topic cdc.public.orders \
  --from-beginning

# 事件 payload 示例(wal2json 风格)
# {"op":"u","ts_ms":1720000000000,"before":{"id":1,"amount":90},"after":{"id":1,"amount":80}}

6.4 Debezium 配置注意事项

配置项建议
plugin.namePG 10+ 用 pgoutput,旧版本用 decoderbufs/wal2json
slot.name自定义复制槽,避免被自动清理
publication.autocreate.modeall_tables / filtered,控制发布范围
heartbeat.interval.ms设置心跳避免 connector 过期
database.history需 Kafka topic 记录 schema 历史

七、监控与排障

7.1 复制状态监控

-- 订阅端:查看复制进度与延迟
SELECT subname, pid, received_lsn,
       latest_end_lsn,
       pg_wal_lsn_diff(latest_end_lsn, received_lsn) AS lag_bytes,
       last_msg_send_time
FROM pg_stat_subscription;

-- 主库:查看复制槽消费位置
SELECT slot_name, database, active,
       restart_lsn,
       pg_size_pretty(pg_wal_lsn_diff(pg_current_wal_lsn(), restart_lsn)) AS wal_retained
FROM pg_replication_slots;

7.2 复制槽 WAL 堆积

逻辑复制槽会保留未被消费的 WAL。若订阅端长时间宕机,主库 WAL 会无限堆积撑爆磁盘:

-- 找出 WAL 堆积最多的复制槽
SELECT slot_name, active,
       pg_size_pretty(pg_wal_lsn_diff(pg_current_wal_lsn(), restart_lsn)) AS retained
FROM pg_replication_slots
ORDER BY pg_wal_lsn_diff(pg_current_wal_lsn(), restart_lsn) DESC;

-- 确认无用的槽及时删除
SELECT pg_drop_replication_slot('stale_slot');

7.3 常见故障与排查

故障可能原因排查命令
订阅长时间无进展网络/磁盘/大事务pg_stat_subscription
复制槽 WAL 暴涨订阅端离线、主库写入大pg_replication_slots
主键冲突报错订阅端有写入查看日志 23505
表结构不一致DDL 未同步比对两端 \d
复制中断不重连网络抖动ALTER SUBSCRIPTION mysub DISABLE; ENABLE;

7.4 告警规则

# Prometheus 告警:复制延迟
groups:
  - name: pg-logical-replication
    rules:
      - alert: ReplicationLagHigh
        expr: pg_stat_subscription_lag_bytes > 104857600   # 100MB
        for: 5m
        labels: { severity: warning }
      - alert: ReplicationSlotWALRetention
        expr: pg_replication_slots_retained_wal_bytes > 10737418240  # 10GB
        for: 10m
        labels: { severity: critical }

八、性能与最佳实践

8.1 提升复制吞吐

-- 订阅端:关闭同步提交加速 apply
ALTER SUBSCRIPTION mysub SET (synchronous_commit = off);

-- 主库:必要时提升 max_wal_senders / max_replication_slots
-- 调整 WAL 级参数需要重启

8.2 复制性能瓶颈分析

瓶颈特征对策
WAL 生成速率主库写入量大垂直扩容、分区、批量写入
网络带宽延迟随数据量上升只发布必要表/行过滤
订阅端 apply 慢目标库负载高调 maintenance_work_mem、并行
无主键表UPDATE/DELETE 全表定位加主键或 REPLICA IDENTITY

8.3 最佳实践要点

□ 发布表尽量都有主键(否则设 REPLICA IDENTITY FULL)
□ 只发布需要的表,避免全库复制
□ 监控复制槽 WAL 保留量,防磁盘爆满
□ DDL 变更走统一流程,两端同步执行
□ 订阅端保持只读,避免数据冲突
□ 使用 sslmode=require 加密复制链路
-- 加密复制链路示例
CREATE SUBSCRIPTION mysub
  CONNECTION 'host=primary.example.com port=5432 dbname=mydb user=repl sslmode=require'
  PUBLICATION mypub;

常见问题(FAQ)

逻辑复制和流复制选哪个?

高可用主从切换用物理流复制(延迟更低、可自动 failover);数据分发、CDC、跨版本迁移用逻辑复制。两者可以共存:一个实例同时做物理从库和逻辑发布,都是常见架构。

为什么订阅端收不到 DDL?

逻辑复制只复制 DML,不复制 DDL,这是设计如此。表结构变更必须在订阅端手动执行(或借助变更管理工具),且顺序上要保证主键/唯一约束两端一致,否则后续 DML 会失败。

主键冲突报错 23505 怎么办?

说明订阅端也存在写入,和复制过来的数据撞了主键。规范做法是订阅端只读;如果必须双写,要保证应用写入的主键与源库不重叠,或用唯一键做幂等合并。

逻辑复制有延迟吗?怎么降低?

有。延迟来源包括 WAL 生成、网络传输、订阅端 apply 单行执行。降低手段:关闭订阅端 synchronous_commit、只发布必要表、为无主键表加主键、订阅端提升硬件/并行。持续监控 pg_stat_subscription 的 lag_bytes。

复制槽会带来磁盘风险吗?

会。逻辑复制槽会保留未消费的 WAL,订阅端宕机越久,主库 WAL 积压越多,最终可能撑爆主库磁盘。必须监控 pg_replication_slots 的 retained WAL 量,并配置告警;长期无用的槽及时 pg_drop_replication_slot。


相关阅读

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「database」更多文章

  1. 数据库安全加固与审计实战:权限最小化、加密、脱敏与合规
  2. 数据库容量规划与资源治理:从评估、监控到扩展路径
  3. 数据库字符集、排序规则与乱码实战:utf8mb4、Collation 选择与排查