序言
数据库里的一条记录变了,影响面往往比想象中大得多:
- 商品或内容改了,Elasticsearch 索引得跟上;
- 用户资料或配置更新了,Redis 缓存和派生视图要刷新;
- 业务数据源源不断地流入 Kafka、数据仓库或数据湖;
- 账户、权限这些关键数据的修改,不光要留审计记录,还得触发后续流程。
下游系统五花八门,但它们依赖的能力其实是一样的:持续捕获数据库里的 INSERT、UPDATE 和 DELETE,并可靠地把这些变化送出去。这就是 CDC(Change Data Capture,变更数据捕获)。
不过,真正让这件事变得棘手的,不是“读到一次变化”,而是当消费者重启、网络中断、全量扫描、下游拥塞或者主备切换发生时,仍然要保证:
- 已经提交的数据不能丢;
- 同一行的变化顺序要可控;
- 消费者能从正确的位置恢复;
- 存量数据和后续增量能无缝衔接。
这篇文章就从这些问题出发,一步步拆解 PostgreSQL WAL CDC 的原理和工程实现,顺便和更熟悉的 MySQL Binlog CDC 做个对比。

1. 为什么选择日志型 CDC
假设要把数据库的变化同步到搜索引擎、缓存或消息队列,最直接的办法似乎是让应用在写数据库之后,再写一次下游。但一放到故障场景里,问题就来了:数据库已经提交,消息发送却失败了,该以谁为准?
常见的 CDC 方案大概有四类:
| 方案 | 工作方式 | 优点 | 主要问题 |
|---|---|---|---|
| 应用层投递 | 应用写数据库后主动投递消息 | 能携带业务语义 | 存在双写一致性问题,也覆盖不了批量 SQL 和绕过应用的改库 |
| 查询轮询 | 定期按 updated_at 查询 | 简单、权限要求低 | 有延迟,增加主库压力;为了感知删除,还得引入 deleted_at 之类的软删除标记 |
| 数据库触发器 | 在业务表上创建触发器,行数据变化时同步写入变更表 | 实时,且与事务同步 | 侵入业务,产生写放大,数据库额外承担了太多工作 |
| 日志解析 | 从数据库事务日志还原变化 | 低侵入、低延迟、覆盖完整 | 依赖数据库日志机制,消费端实现更复杂 |
至于为什么日志解析成了现代 CDC 的主流,原因其实很简单:数据库已经替我们完成了最困难的那一步——把所有已提交的变化,按可恢复的顺序,老老实实地记录了下来。
在 MySQL 里,这份日志是 binlog;在 PostgreSQL 里,则是 WAL(Write-Ahead Log,预写日志)。两者目标相似,但底层的语义和工程责任并不完全一样。
2. PostgreSQL 与 MySQL CDC 全景对比
如果已经熟悉 MySQL Binlog CDC,那换个角度来理解 PostgreSQL 会更容易。
| 对比维度 | MySQL Binlog CDC | PostgreSQL WAL CDC | 工程影响 |
|---|---|---|---|
| 日志基础 | binlog 原本主要服务于复制和增量恢复 | WAL 原本主要服务于崩溃恢复和物理复制 | PostgreSQL 需要先做逻辑解码,MySQL ROW event 已经有较强的逻辑语义 |
CDC 前置配置 | log_bin=ON、binlog_format=ROW | wal_level=logical | 都需要实例级配置,部分参数变更可能得重启 |
| 解码方式 | 解析 binlog event | Logical Decoding + Output Plugin | PostgreSQL 可以通过插件选择输出协议 |
| 内置输出格式 | ROW event | pgoutput 二进制协议 | 两者通常都交给成熟客户端库去解析 |
| 数据范围 | 通常在消费端按库表过滤 | Publication 在服务端声明表和操作范围;PG 15+ 支持行过滤与列列表 | PostgreSQL 能在源端缩小发布范围 |
| 消费位置 | binlog file/position 或 GTID | LSN | 都要持久化并用于断点恢复,但语义并不完全等价 |
| 进度保存 | 通常由消费者保存位点 | Replication Slot 在服务端保存确认水平 | PostgreSQL 对消费者更友好,但把资源保留责任放到了数据库 |
| 日志保留 | binlog 通常按时间或空间过期 | Slot 会保留消费者仍然需要的 WAL | PG 需要关注 Slot 对 WAL 保留量的影响;MySQL 需要关注日志过期后的续传问题 |
| UPDATE/DELETE 旧值 | binlog_row_image 控制行镜像 | 逐表设置 REPLICA IDENTITY | PostgreSQL 上线前要逐表检查旧值需求 |
| 全量与增量衔接 | 一致性快照 + binlog 位点/GTID | 一致点 LSN + 导出快照 | 原理相同:先确定日志边界,再读取同一时刻的存量 |
| DDL | binlog 通常包含 DDL Query Event | 逻辑复制默认不传播 DDL | PostgreSQL 需要单独设计 Schema 演进机制 |
| 大字段 | 受 row image 配置影响 | 需要识别 unchanged TOAST | PG 消费端不能把“未变化”误认为 NULL |
| 故障切换 | 关注 GTID 连续性和拓扑切换 | 关注 timeline、Slot 是否延续及新主库起始 LSN | 两者都不能只靠重连地址判断是否能安全续传 |

最值得记住的区别是:PostgreSQL 通过 Slot 在服务端保存消费者的确认水平,这使得断点恢复的逻辑更直接。但代价是,当消费者长时间没推进时,Slot 所需的 WAL 保留量就会增加,因此需要配合容量规划和监控来管理。
3. 一张图看懂 PostgreSQL CDC
先不看那些复杂的协议细节,一套 PostgreSQL CDC 链路可以概括成这样:

一次事务从数据库执行到变更投递,大致会经过下面这些环节:
- PostgreSQL 在修改数据页前,先把相关日志写入
WAL,提交时再写入提交记录; - 事务确认提交后,Logical Decoding 把底层变化还原成行级语义;
- Output Plugin(通常是
pgoutput)把变化编码成协议消息; CDC消费者通过 Replication Slot 持续读取;- 消费者完成转换、路由和投递;
- 下游确认成功后,消费者才向 PostgreSQL 上报安全 LSN;
- PostgreSQL 推进 Slot 水平,并回收不再需要的
WAL。
这是一条闭环。只读不确认,Slot 需要保留的 WAL 就会持续增加;还没投递成功就提前确认,则可能永久丢数据。
4. WAL 如何变成行级事件
WAL 是 PostgreSQL 保证持久性和崩溃恢复的基础。数据页落盘之前,相关的修改必须先写入 WAL,这也是“Write-Ahead”这个名字的由来。
原始的 WAL 是面向数据页和内部操作的物理/物理逻辑日志,它主要服务于:
- 崩溃恢复:重启后重放
WAL,让数据恢复到一致状态; - 物理流复制:把
WAL传给备库,得到主库的数据页副本。
举个例子,原始的 WAL 记录可能只告诉你“某个数据页发生了变化”,但没办法直接告诉消费者“public.products 表里主键为 42 的商品价格变成了 99”。前者是底层存储的变化,后者才是 CDC 需要的行级业务语义。
为了把这种底层变化转换成 CDC 能消费的行级事件,PostgreSQL 从 9.4 开始引入了 Logical Decoding(逻辑解码)。服务端读取 WAL,结合事务、表结构等信息,把它还原成带有表和行语义的变化。
要让 WAL 携带逻辑解码需要的信息,实例参数 wal_level 必须设为 logical。
逻辑解码只负责还原变化,最终输出什么格式,取决于 Output Plugin:
| 插件 | 格式 | 提供方式 | 适合场景 |
|---|---|---|---|
pgoutput | PostgreSQL 二进制逻辑复制协议 | PG 10+ 内置 | 生产系统首选,标准、性能好、无需安装扩展 |
wal2json | JSON | 第三方插件 | 调试友好,消费端容易接入 |
decoderbufs | Protobuf | 第三方插件 | 二进制输出,Debezium 早期方案 |
test_decoding | 文本 | 随 PostgreSQL contrib 提供,取决于安装包 | 学习和测试,不面向生产 |
新系统通常优先选用 pgoutput。它的缺点是人眼读不了,但这属于客户端库该解决的问题,没必要为了调试方便,长期承担第三方插件的运维成本。
5. 四个核心概念
开始消费逻辑复制流之前,最好先理解四个彼此关联的概念:
| 概念 | 可以理解为 | 回答的问题 |
|---|---|---|
| Publication | 发布清单 | 哪些表和操作需要输出? |
| Replication Slot | 消费者书签 | 这个消费者已经读到哪里了? |
| LSN | WAL 坐标 | 某次变化在日志的哪个位置? |
| REPLICA IDENTITY | 旧记录识别规则 | UPDATE/DELETE 时,怎么定位变化前的记录? |
它们共同描述了一次逻辑复制:Publication 确定捕获范围,LSN 标记每个变化的位置,Replication Slot 用 LSN 保存消费者进度,REPLICA IDENTITY 决定 UPDATE 和 DELETE 能携带哪些旧值。

Publication:捕获哪些变化
Publication 是 PostgreSQL 对“哪些表参与逻辑复制”的服务端声明:
从早期版本起,Publication 就能选择表和操作类型,PG 15+ 还进一步支持了行过滤与列列表:
复制代码-- 发布所有表
CREATE PUBLICATION cdc_pub FOR ALL TABLES;-- 只发布指定表
CREATE PUBLICATION cdc_pub
FOR TABLE public.users, public.products;-- 只发布 INSERT 和 UPDATE
CREATE PUBLICATION cdc_pub
FOR TABLE public.users
WITH (publish = 'insert, update');-- PG 15+:行过滤用于 UPDATE/DELETE 时,过滤列必须包含在 Replica Identity 中
ALTER TABLE public.orders REPLICA IDENTITY FULL;CREATE PUBLICATION paid_orders_pub
FOR TABLE public.orders WHERE (status = 'paid');-- PG 15+:只发布指定列
CREATE PUBLICATION product_price_pub
FOR TABLE public.products (id, name, price);
和 MySQL 常见的消费端过滤相比,Publication 可以在源端就缩小范围。不过要注意,它并不是权限系统:发布范围、用户的读取权限和复制权限,还是需要分别配置的。
Replication Slot:消费者读到哪里了
设想一下,数据仓库维护两小时,CDC 消费者暂时离线了。它恢复后应该从哪里继续?离线期间的日志又该由谁来保留?
Replication Slot 就是 PostgreSQL 为消费者维护的服务端书签。逻辑 Slot 主要关注两个水平:
restart_lsn:该 Slot 可能还需要的最早WAL位置;confirmed_flush_lsn:消费者已经明确确认安全持久化的位置。
在没超过 max_slot_wal_keep_size 等保留限制时,PostgreSQL 会为 Slot 保留仍然需要的 WAL,所以消费者能够断点续传。但如果超过了限制,Slot 可能因为所需的 WAL 已经被移除而失效;消费者长时间不推进,也需要关注 WAL 保留量的变化。
这一点也体现了 PostgreSQL 和 MySQL 在日志保留机制上的差异:PostgreSQL 把“保留哪些日志”的责任更多地交给了服务端,而 MySQL 则依赖消费者自己去管理日志的过期和续传风险。
LSN:变化发生在哪里
LSN(Log Sequence Number)是 WAL 里的位置坐标,文本形式看起来像 16/B374D848。在同一集群的正常 WAL 序列中,它会随着 WAL 写入不断向前推进。
CDC 系统会同时遇到好几个 LSN:
- 服务端当前写到哪里了;
- 复制流已经发送到哪里了;
- 消费者已经处理到哪里了;
- 下游已经安全确认到哪里了。
真正允许上报给 Slot 的,只能是最后一个。
REPLICA IDENTITY:旧值能看到多少
DELETE 之后已经没有新行了,UPDATE 也可能需要知道修改前的分区键。PostgreSQL 必须决定在 WAL 里保留哪些旧值,这由逐表属性 REPLICA IDENTITY 控制。
| 设置 | UPDATE/DELETE 可用的旧值 | 典型用途 |
|---|---|---|
DEFAULT | DELETE 携带旧主键;UPDATE 仅在主键变化时携带旧键 | 下游按主键定位 |
USING INDEX | DELETE 携带索引旧值;UPDATE 仅在索引键变化时携带旧键 | 没有主键但有合适的唯一键 |
FULL | 所有列 | 审计、计算差异、旧分片数据清理 |
NOTHING | 不提供旧键 | 发布 UPDATE/DELETE 时无法满足复制要求,相关操作会报错 |
复制代码ALTER TABLE public.accounts REPLICA IDENTITY FULL;
MySQL 主要通过 binlog_row_image 控制行镜像,而 PostgreSQL 是逐表来设置的。不要机械地给所有表都设为 FULL:它能提供更完整的 before image,但也会增加 UPDATE/DELETE 的 WAL 体积。具体怎么选,应该根据下游的实际需求来定。
6. 逻辑复制的流式协议
CDC 消费者使用复制连接连上 PostgreSQL,然后通过基于 COPY 的流式协议持续接收数据。

协议里的关键动作包括:
IDENTIFY_SYSTEM:获取集群标识、timeline 和当前的WAL位置;START_REPLICATION SLOT ... LOGICAL:从指定的 Slot 和 LSN 开始消费;XLogData:承载逻辑解码后的数据;- Primary Keepalive:服务端发送的心跳,可能要求消费者立即回复;
- Standby Status Update:消费者上报自己已接收、已刷盘和已应用的位置。
使用 pgoutput 时,还需要解析 Begin、Commit、Relation、Insert、Update、Delete、Truncate 这些逻辑消息。一个典型的事务可以理解为:
复制代码Begin → Relation(必要时)→ Insert/Update/Delete... → Commit
默认模式下,逻辑复制流是按事务提交顺序输出的,Begin 和 Commit 之间的内容属于同一个事务。PG 14+ 开启大事务流式解码后,一个事务可能被拆成多个 Stream Start/Stop 片段,最后收到 Stream Commit 或 Stream Abort;消费端应该按 xid 归组,确认提交后再应用。生产环境里,还是应该用成熟的客户端库来处理 CopyData 和插件协议,不要假设一次网络读取就恰好对应一条业务变更。
需要特别澄清的是:正常的 PostgreSQL 主备切换通常仍然属于同一个数据库集群,systemid 不一定变化。恢复时除了检查 systemid,还必须确认 timeline、逻辑 Slot 是否已经同步到新主库,以及请求的 LSN 是否仍然有效。
7. 服务端配置与最小示例
在继续讨论快照和工程实现之前,先用一组明确的示例对象,把最小链路跑通。后续命令统一使用以下名称:
| 对象 | 示例值 |
|---|---|
| 数据库 | appdb |
| Schema | public |
| 业务表 | products |
| 复制用户 | cdc_user |
| Publication | cdc_pub |
| Replication Slot | demo_slot |
| 客户端网段 | 10.0.0.0/8 |
配置 PostgreSQL 实例
复制代码# postgresql.conf
wal_level = logical
max_replication_slots = 10
max_wal_senders = 10
这些值至少应该覆盖实际的 Slot 和复制连接数量,并预留一些运维空间。修改 wal_level 这类启动参数后,需要重启实例才能生效。
准备数据库、示例表和复制用户
先创建示例数据库,然后用 psql 的 connect 命令切到该数据库。以下示例里的所有对象都位于 appdb:
复制代码CREATE DATABASE appdb;connect appdbCREATE TABLE public.products (
id bigint PRIMARY KEY,
name text NOT NULL,
price numeric(12, 2) NOT NULL,
attributes jsonb,
updated_at timestamptz NOT NULL DEFAULT now()
);-- 后续消息示例需要展示完整 before image
ALTER TABLE public.products REPLICA IDENTITY FULL;CREATE USER cdc_user
WITH REPLICATION LOGIN PASSWORD 'replace-with-a-secret';GRANT CONNECT ON DATABASE appdb TO cdc_user;
GRANT USAGE ON SCHEMA public TO cdc_user;
GRANT SELECT ON TABLE public.products TO cdc_user;
其中 REPLICATION 权限用于建立复制连接,SELECT 权限用于读取已发布表和执行全量快照。
逻辑复制连接需要指定实际数据库,所以 pg_hba.conf 应该放行示例数据库 appdb,而不是物理复制连接里用的特殊数据库关键字 replication:
复制代码# TYPE DATABASE USER ADDRESS METHOD
host appdb cdc_user 10.0.0.0/8 scram-sha-256
修改后执行 SELECT pg_reload_conf(); 或 reload 配置。
生产链路:创建 Publication
复制代码CREATE PUBLICATION cdc_pub FOR TABLE public.products;
用 pgoutput 构建生产链路时,这条命令声明了需要捕获 public.products。生产消费者通过复制协议订阅 cdc_pub,Publication 决定了进入逻辑复制流的表和操作范围。
本地验证链路:使用 test_decoding 观察变化
为了让变化能直接看得到,下面单独用文本输出插件 test_decoding。这条本地验证链路不会读取上面的 cdc_pub,和用 pgoutput + Publication 的生产链路不是一回事。
复制代码SELECT *
FROM pg_create_logical_replication_slot('demo_slot', 'test_decoding');INSERT INTO public.products (id, name, price)
VALUES (1, 'Mechanical Keyboard', 699.00);UPDATE public.products
SET price = 649.00, updated_at = now()
WHERE id = 1;DELETE FROM public.products WHERE id = 1;
上面三条语句分别产生了一条 INSERT、UPDATE 和 DELETE 变化。现在读取 Slot 里已经解码的内容:
复制代码SELECT lsn, xid, data
FROM pg_logical_slot_get_changes('demo_slot', NULL, NULL);
pg_logical_slot_get_changes 会消费变化并推进位置;pg_logical_slot_peek_changes 只查看、不消费,更适合反复调试。
用完后记得及时清理测试 Slot:
复制代码SELECT pg_drop_replication_slot('demo_slot');
这个例子用 test_decoding 和 SQL 函数,是为了让变化肉眼可见。生产系统通常会用 pgoutput、Publication 和复制协议来持续消费,而不是轮询上面的 SQL 函数。
8. 存量快照与增量如何衔接
逻辑复制只能提供 Slot 保留范围内的增量变化,没办法自动还原表里早就存在的全部数据。所以,第一次启动 CDC 时,需要完成两件事:
- 读取当前已有的数据,给下游建立一个完整的基线;
- 从一个确定的 LSN 开始,消费此后的增量变化。
真正的难点不在于分别完成全量扫描和增量消费,而在于让两者共享同一个边界。可以把这个边界想象成一次切分:
比如,把 appdb 里的 public.products 全量导入搜索引擎,可能需要一小时,但扫描期间商品价格可能还在变。如果等扫描结束才临时确定增量起点,中间提交的变化就可能漏掉;如果没控制好应用顺序,较旧的快照数据还可能覆盖掉较新的增量结果。

PostgreSQL 通过两个相互对应的值来建立这条边界:
| 返回值 | 作用 |
|---|---|
consistent_point | 逻辑复制流的安全起始 LSN,边界之后提交的变化可以从这里开始读 |
snapshot_name | 和该边界对应的 MVCC 快照,让一个或多个扫描连接能看到同一时刻的数据 |
完整的衔接过程
- 建立边界:通过复制协议创建逻辑 Slot,并要求导出快照,获得
consistent_point和snapshot_name。 - 导入快照:用一个或多个普通数据库连接开启只读的
REPEATABLE READ事务,执行查询前导入同一个snapshot_name。 - 扫描存量:各个连接并发读取
public.products,把快照数据写入下游;在所有快照数据确认完成之前,不应用更晚的增量事件。 - 接续增量:从
consistent_point启动逻辑复制,读取边界建立之后提交的变化。
快照扫描连接的示例:
复制代码BEGIN ISOLATION LEVEL REPEATABLE READ READ ONLY;
SET TRANSACTION SNAPSHOT '00000003-0000001A-1';SELECT id, name, price, updated_at
FROM public.products
WHERE id >= 1 AND id < 100000;COMMIT;
导出快照是有生命周期限制的。在所有扫描连接成功执行 SET TRANSACTION SNAPSHOT 之前,应该保持创建快照的复制连接打开,并避免在该连接上继续执行其他复制命令。具体的客户端库通常会封装创建 Slot 和导出快照的协议细节,但调用方仍然需要保证这个时序。
大表如何并发扫描
所有 worker 都导入同一个快照后,可以用不同的条件来拆分扫描任务,同时保持一致的数据视图。常见的分片方式包括:
- 按主键范围:
WHERE id >= ? AND id < ?,简单,但受主键分布影响; - 按业务分区:按日期、租户或原生分区表并发扫描;
- 按 CTID 页范围:更接近物理顺序,适合缺少均匀主键的大表。
CTID 会因 VACUUM FULL、CLUSTER 等表重写操作而发生变化。如果用 CTID 分片,快照期间应该避免这些操作,并在实际版本和表结构上验证扫描计划。
应用顺序与幂等
最容易验证的实现是:先等所有快照数据写入并确认完毕,再从 consistent_point 顺序应用增量。这样,扫描期间发生的价格修改会暂存在 Slot 所保留的 WAL 里,等快照完成后依次追平。
如果为了降低延迟而同时处理快照和增量,就必须保证旧的快照行不会覆盖更新的增量状态,比如设置阶段屏障;需要并发应用时,可以把 consistent_point 作为快照基线版本,然后用事务提交 LSN 和事务内序号来判断增量顺序。
日志型 CDC 通常采用 at-least-once(至少一次),重连和投递重试仍然可能产生重复事件。所以下游应该按主键执行 UPSERT,并结合事务提交 LSN、事务内序号或事件 ID 来实现幂等。大表扫描持续时间较长时,既要监控 Slot 所需的 WAL 保留量,也要关注长事务快照对 VACUUM 和表膨胀的影响。
如果 Publication 用了行过滤或列列表,快照查询也应该用等价的过滤条件和字段投影,保证存量基线和增量范围一致。
MySQL 用的具体原语不同,但原则是一样的:先取得一致性快照及其对应的 binlog position/GTID,再完成存量读取和增量接续。
9. 工程实现的正确性核心
一个健壮的 CDC 消费者,不光要能解析协议,还得正确处理确认、并发、顺序和异常数据。核心问题可以归纳成下面这样:
| 工程问题 | 处理原则 |
|---|---|
| LSN 确认时机 | 完整事务在下游确认成功后,才能上报该事务的 Commit.end_lsn |
| 并发完成与安全水平 | 只推进连续完成的事务边界;在途窗口满了就暂停读取,形成背压 |
| 同一行的变更顺序 | 用 schema.table.primary_key 作为稳定路由键,多实例时统一哈希算法、编码和种子 |
| unchanged TOAST | 保留“字段未变化”的语义,由下游合并旧值,不能把它当成 NULL 或空值 |
| 大事务 | 限制在途数据并做好背压;PG 14+ 可以结合客户端和插件对流式解码大事务的支持 |
| 超大单行 | 用 Chunking 分片、Claim Check 外置内容,或者在接受时序差异的前提下回源查询 |
| DDL 与 Schema 演进 | 用 Event Trigger、pg_catalog、Schema History 或数据库迁移事件来同步结构变化 |

这里面最关键的,是区分“已经收到”、“已经处理”和“已经安全确认”这三个位置。推荐的推进顺序是:
复制代码接收 → 解码 → 投递 → 下游确认 → 记录安全水平(如需要)→ 上报安全 LSN
并发不会改变这个原则。如果较晚提交的事务已经完成了,但更早的事务还在重试,安全水平必须停在较早事务之前。整体投递通常采用 at-least-once(至少一次),然后通过下游的幂等来吸收重连和重试产生的重复事件。
PostgreSQL 逻辑复制默认不传播 DDL;MySQL binlog 通常能看到 DDL Query Event,但两者都需要显式处理 Schema 事件和行数据之间的顺序及兼容性问题。
10. 消息模型设计
快照事件和增量事件最好用同一套结构,只通过 op 来区分。一个实用的事件通常包含下面这些字段:
复制代码{
"key": "public.products:42",
"source": "production-pg",
"op": "UPDATE",
"commit_lsn": "16/B374D848",
"event_index": 1,
"xid": 123456,
"commit_ts": "2026-07-23T10:00:00Z",
"schema": "public",
"table": "products",
"primary_key": {"id": 42},
"before": {"price": 100},
"after": {"price": 99},
"unchanged_toast": ["attributes"]
}
设计时需要重点考虑:
- 路由稳定:key 能稳定标识同一行,保证分区内的顺序;
- 位置可比较:携带事务提交 LSN 和事务内序号,方便排错、排序和去重;
- 事务可追踪:保留 xid,必要时支持事务级聚合;
- 旧值语义明确:区分缺失、
NULL和 unchanged TOAST; - Schema 可演进:用 Protobuf/A vro 时遵守兼容性规则;
- 快照增量同构:下游不需要维护两套完全不同的处理逻辑。
消息格式不是越完整越好。完整的 before image、列元数据和事务信息都会增加体积,具体应该由下游需求来驱动。
11. 开源工具与客户端
下面列出采用 OSI 认可许可证、可用于 PostgreSQL CDC 的开源工具:
| 工具 | 形态/语言 | 主要用途 | 开源许可证 |
|---|---|---|---|
| Debezium + Kafka Connect | Ja va · Connector 运行时 | 捕获多种数据库的变化并写入 Kafka | Apache-2.0 |
| Debezium Server | Ja va · 独立运行时 | 将数据库变化直接输出到多种消息系统 | Apache-2.0 |
| Apache Flink CDC | Ja va · 分布式数据管道 | 将数据库快照和增量接入 Flink 数据管道 | Apache-2.0 |
| Apache SeaTunnel | Ja va · 数据集成平台 | 在多种数据源和目标之间同步批量与增量数据 | Apache-2.0 |
xataio/pgstream | Go · CLI/库 | PostgreSQL 到 Kafka、OpenSearch、Webhook 或 PostgreSQL | Apache-2.0 |
ConduitIO/conduit | Go · Connector 框架 | 通过 Connector 连接 PostgreSQL 与其他数据系统 | Apache-2.0 |
| PeerDB | Go/Rust · 复制平台 | PostgreSQL 到分析系统、队列和对象存储 | AGPL-3.0 |
| Sequin | Elixir · 自托管服务 | PostgreSQL 到队列、搜索引擎和 Webhook | MIT |
Airbyte 的主要代码采用 ELv2,Materialize 当前版本采用 BSL 1.1。两者源码可见,但许可证不是 OSI 认可的开源许可证,所以没有放进上面的开源工具清单里。PeerDB 已经在 2026 年改为 AGPL-3.0,属于开源软件,但使用时需要遵守比较强的 copyleft 条款。
学习和排障时,可以用 pg_recvlogical、test_decoding 或 wal2json 直接观察解码结果。
用于自研的常见客户端也可以直接对比:
| 语言 | 客户端/API | 作用 | 开源许可证 |
|---|---|---|---|
| Go | jackc/pglogrepl + jackc/pgx | 复制协议、pgoutput 消息和 PostgreSQL 连接 | MIT |
| Ja va | PostgreSQL JDBC PGReplicationStream | JDBC 驱动内置的逻辑复制接口 | BSD-2-Clause |
| Python | psycopg2 | LogicalReplicationConnection 等逻辑复制接口 | LGPL-3.0-or-later |
| Rust | supabase/etl | 构建 PostgreSQL 逻辑复制与实时数据管道的 Rust 框架 | Apache-2.0 |
| Node.js | kibae/pg-logical-replication | 支持 pgoutput、wal2json 等输出插件 | MIT |
| C/C++ | libpq | PostgreSQL 官方底层客户端库 | PostgreSQL License |
12. 监控、故障恢复与上线检查
生产运维主要关注三件事:消费者跟不跟得上、Slot 需要保留多少 WAL,以及故障后能不能从正确位置恢复。
复制代码SELECT
slot_name,
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
WHERE slot_type = 'logical';
restart_lsn 用来估算 Slot 当前需要保留的 WAL 范围,confirmed_flush_lsn 用来观察消费者的确认进度。除此之外,保留下面这几组指标就够了:
| 监控方向 | 关键指标 |
|---|---|
| Slot 状态 | active、保留的 WAL 大小、确认 LSN 推进速度 |
| 端到端延迟 | 事件提交到下游确认的 P95/P99 |
| 消费者状态 | 在途数量、重试次数、背压时间和重连次数 |
| 主库资源 | WAL 目录大小和磁盘剩余空间 |
故障恢复不能只做断线重连。恢复前应该确认仍然连接到预期的集群、timeline 和 Slot 状态有效、所需的 WAL 仍然存在,并核对本地安全水平与 confirmed_flush_lsn 是否符合预期;不一致时应停止消费并告警。正常的主备切换不一定改变 systemid,所以还需要提前验证所用 PostgreSQL 版本下的 Slot 同步或 failover slot 方案。
上线前可以归纳为六项检查:
- 实例参数、复制用户权限和 Publication 范围是否正确;
- 各表主键与
REPLICA IDENTITY是否满足下游需求; - 快照与增量衔接是否经过并发写入验证;
- LSN 是否只在下游确认后连续推进;
- unchanged TOAST、DDL、大事务和重复事件是否有明确处理策略;
- Slot、端到端延迟和磁盘告警是否已配置,主备切换与废弃 Slot 清理流程是否已演练。
小结
PostgreSQL WAL CDC 的核心,是把数据库里的变化,转化成一条可以持续消费、确认和恢复的事件流。WAL 提供变化来源,Logical Decoding 恢复行级语义,Publication 确定捕获范围,LSN 标记日志位置,Replication Slot 保存消费者进度,REPLICA IDENTITY 则决定 UPDATE 和 DELETE 能提供哪些旧值。
第一次启动 CDC 时,还需要处理存量与增量的边界:导出快照负责边界之前的数据状态,逻辑复制流负责边界之后提交的变化。只要两者使用同一个一致性起点,就不需要暂停业务写入。
工程实现中最重要的原则,是完整事务在下游确认成功后,才能推进安全 LSN。并发消费只能确认连续完成的事务边界,同一行的变化需要保持顺序,重复事件则通过主键 UPSERT、事务提交 LSN 与事务内序号或事件 ID 来实现幂等。系统还需要明确处理 unchanged TOAST、DDL、大事务、超大消息和主备切换。
PostgreSQL 和 MySQL 的日志型 CDC 原理相近,但日志保留方式不同:MySQL 需要关注 binlog 是否在消费者恢复前过期,PostgreSQL 需要关注 Slot 进度带来的 WAL 保留量变化。
最终衡量一套 CDC 系统是否可靠,可以归结为三点:全量与增量之间没有遗漏,数据确认位置能够安全恢复,下游在重复、延迟和乱序情况下仍能收敛到正确状态。
