当前位置: 首页
数据库
PostgreSQL WAL CDC系统构建原理与工程实现

PostgreSQL WAL CDC系统构建原理与工程实现

热心网友 时间:2026-07-24
转载

基于PostgreSQLWAL构建CDC系统,通过逻辑解码将WAL还原为行级事件,利用Publication指定捕获范围、ReplicationSlot维护消费者进度、LSN标记日志位置、REPLICAIDENTITY控制旧值。流式协议确保断点续传,与MySQLBinlogCDC相比,服务端保留日志的责任更重。

序言

数据库里的一条记录变了,影响面往往比想象中大得多:

  • 商品或内容改了,Elasticsearch 索引得跟上;
  • 用户资料或配置更新了,Redis 缓存和派生视图要刷新;
  • 业务数据源源不断地流入 Kafka、数据仓库或数据湖;
  • 账户、权限这些关键数据的修改,不光要留审计记录,还得触发后续流程。

下游系统五花八门,但它们依赖的能力其实是一样的:持续捕获数据库里的 INSERTUPDATEDELETE,并可靠地把这些变化送出去。这就是 CDC(Change Data Capture,变更数据捕获)

不过,真正让这件事变得棘手的,不是“读到一次变化”,而是当消费者重启、网络中断、全量扫描、下游拥塞或者主备切换发生时,仍然要保证:

  1. 已经提交的数据不能丢;
  2. 同一行的变化顺序要可控;
  3. 消费者能从正确的位置恢复;
  4. 存量数据和后续增量能无缝衔接。

这篇文章就从这些问题出发,一步步拆解 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 CDCPostgreSQL WAL CDC工程影响
日志基础binlog 原本主要服务于复制和增量恢复WAL 原本主要服务于崩溃恢复和物理复制PostgreSQL 需要先做逻辑解码,MySQL ROW event 已经有较强的逻辑语义
CDC 前置配置log_bin=ONbinlog_format=ROWwal_level=logical都需要实例级配置,部分参数变更可能得重启
解码方式解析 binlog eventLogical Decoding + Output PluginPostgreSQL 可以通过插件选择输出协议
内置输出格式ROW eventpgoutput 二进制协议两者通常都交给成熟客户端库去解析
数据范围通常在消费端按库表过滤Publication 在服务端声明表和操作范围;PG 15+ 支持行过滤与列列表PostgreSQL 能在源端缩小发布范围
消费位置binlog file/position 或 GTIDLSN都要持久化并用于断点恢复,但语义并不完全等价
进度保存通常由消费者保存位点Replication Slot 在服务端保存确认水平PostgreSQL 对消费者更友好,但把资源保留责任放到了数据库
日志保留binlog 通常按时间或空间过期Slot 会保留消费者仍然需要的 WALPG 需要关注 Slot 对 WAL 保留量的影响;MySQL 需要关注日志过期后的续传问题
UPDATE/DELETE 旧值binlog_row_image 控制行镜像逐表设置 REPLICA IDENTITYPostgreSQL 上线前要逐表检查旧值需求
全量与增量衔接一致性快照 + binlog 位点/GTID一致点 LSN + 导出快照原理相同:先确定日志边界,再读取同一时刻的存量
DDLbinlog 通常包含 DDL Query Event逻辑复制默认不传播 DDLPostgreSQL 需要单独设计 Schema 演进机制
大字段受 row image 配置影响需要识别 unchanged TOASTPG 消费端不能把“未变化”误认为 NULL
故障切换关注 GTID 连续性和拓扑切换关注 timeline、Slot 是否延续及新主库起始 LSN两者都不能只靠重连地址判断是否能安全续传

最值得记住的区别是:PostgreSQL 通过 Slot 在服务端保存消费者的确认水平,这使得断点恢复的逻辑更直接。但代价是,当消费者长时间没推进时,Slot 所需的 WAL 保留量就会增加,因此需要配合容量规划和监控来管理。


3. 一张图看懂 PostgreSQL CDC

先不看那些复杂的协议细节,一套 PostgreSQL CDC 链路可以概括成这样:

一次事务从数据库执行到变更投递,大致会经过下面这些环节:

  1. PostgreSQL 在修改数据页前,先把相关日志写入 WAL,提交时再写入提交记录;
  2. 事务确认提交后,Logical Decoding 把底层变化还原成行级语义;
  3. Output Plugin(通常是 pgoutput)把变化编码成协议消息;
  4. CDC 消费者通过 Replication Slot 持续读取;
  5. 消费者完成转换、路由和投递;
  6. 下游确认成功后,消费者才向 PostgreSQL 上报安全 LSN;
  7. 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

插件格式提供方式适合场景
pgoutputPostgreSQL 二进制逻辑复制协议PG 10+ 内置生产系统首选,标准、性能好、无需安装扩展
wal2jsonJSON第三方插件调试友好,消费端容易接入
decoderbufsProtobuf第三方插件二进制输出,Debezium 早期方案
test_decoding文本随 PostgreSQL contrib 提供,取决于安装包学习和测试,不面向生产

新系统通常优先选用 pgoutput。它的缺点是人眼读不了,但这属于客户端库该解决的问题,没必要为了调试方便,长期承担第三方插件的运维成本。


5. 四个核心概念

开始消费逻辑复制流之前,最好先理解四个彼此关联的概念:

概念可以理解为回答的问题
Publication发布清单哪些表和操作需要输出?
Replication Slot消费者书签这个消费者已经读到哪里了?
LSNWAL 坐标某次变化在日志的哪个位置?
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 可用的旧值典型用途
DEFAULTDELETE 携带旧主键;UPDATE 仅在主键变化时携带旧键下游按主键定位
USING INDEXDELETE 携带索引旧值;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 时,还需要解析 BeginCommitRelationInsertUpdateDeleteTruncate 这些逻辑消息。一个典型的事务可以理解为:

 复制代码Begin → Relation(必要时)→ Insert/Update/Delete... → Commit

默认模式下,逻辑复制流是按事务提交顺序输出的,BeginCommit 之间的内容属于同一个事务。PG 14+ 开启大事务流式解码后,一个事务可能被拆成多个 Stream Start/Stop 片段,最后收到 Stream CommitStream Abort;消费端应该按 xid 归组,确认提交后再应用。生产环境里,还是应该用成熟的客户端库来处理 CopyData 和插件协议,不要假设一次网络读取就恰好对应一条业务变更。

需要特别澄清的是:正常的 PostgreSQL 主备切换通常仍然属于同一个数据库集群,systemid 不一定变化。恢复时除了检查 systemid,还必须确认 timeline、逻辑 Slot 是否已经同步到新主库,以及请求的 LSN 是否仍然有效。


7. 服务端配置与最小示例

在继续讨论快照和工程实现之前,先用一组明确的示例对象,把最小链路跑通。后续命令统一使用以下名称:

对象示例值
数据库appdb
Schemapublic
业务表products
复制用户cdc_user
Publicationcdc_pub
Replication Slotdemo_slot
客户端网段10.0.0.0/8

配置 PostgreSQL 实例

 复制代码# postgresql.conf
wal_level = logical
max_replication_slots = 10
max_wal_senders = 10

这些值至少应该覆盖实际的 Slot 和复制连接数量,并预留一些运维空间。修改 wal_level 这类启动参数后,需要重启实例才能生效。

准备数据库、示例表和复制用户

先创建示例数据库,然后用 psqlconnect 命令切到该数据库。以下示例里的所有对象都位于 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 时,需要完成两件事:

  1. 读取当前已有的数据,给下游建立一个完整的基线;
  2. 从一个确定的 LSN 开始,消费此后的增量变化。

真正的难点不在于分别完成全量扫描和增量消费,而在于让两者共享同一个边界。可以把这个边界想象成一次切分:

比如,把 appdb 里的 public.products 全量导入搜索引擎,可能需要一小时,但扫描期间商品价格可能还在变。如果等扫描结束才临时确定增量起点,中间提交的变化就可能漏掉;如果没控制好应用顺序,较旧的快照数据还可能覆盖掉较新的增量结果。

PostgreSQL 通过两个相互对应的值来建立这条边界:

返回值作用
consistent_point逻辑复制流的安全起始 LSN,边界之后提交的变化可以从这里开始读
snapshot_name和该边界对应的 MVCC 快照,让一个或多个扫描连接能看到同一时刻的数据

完整的衔接过程

  1. 建立边界:通过复制协议创建逻辑 Slot,并要求导出快照,获得 consistent_pointsnapshot_name
  2. 导入快照:用一个或多个普通数据库连接开启只读的 REPEATABLE READ 事务,执行查询前导入同一个 snapshot_name
  3. 扫描存量:各个连接并发读取 public.products,把快照数据写入下游;在所有快照数据确认完成之前,不应用更晚的增量事件。
  4. 接续增量:从 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 FULLCLUSTER 等表重写操作而发生变化。如果用 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 ConnectJa va · Connector 运行时捕获多种数据库的变化并写入 KafkaApache-2.0
Debezium ServerJa va · 独立运行时将数据库变化直接输出到多种消息系统Apache-2.0
Apache Flink CDCJa va · 分布式数据管道将数据库快照和增量接入 Flink 数据管道Apache-2.0
Apache SeaTunnelJa va · 数据集成平台在多种数据源和目标之间同步批量与增量数据Apache-2.0
xataio/pgstreamGo · CLI/库PostgreSQL 到 Kafka、OpenSearch、Webhook 或 PostgreSQLApache-2.0
ConduitIO/conduitGo · Connector 框架通过 Connector 连接 PostgreSQL 与其他数据系统Apache-2.0
PeerDBGo/Rust · 复制平台PostgreSQL 到分析系统、队列和对象存储AGPL-3.0
SequinElixir · 自托管服务PostgreSQL 到队列、搜索引擎和 WebhookMIT

Airbyte 的主要代码采用 ELv2,Materialize 当前版本采用 BSL 1.1。两者源码可见,但许可证不是 OSI 认可的开源许可证,所以没有放进上面的开源工具清单里。PeerDB 已经在 2026 年改为 AGPL-3.0,属于开源软件,但使用时需要遵守比较强的 copyleft 条款。

学习和排障时,可以用 pg_recvlogicaltest_decodingwal2json 直接观察解码结果。

用于自研的常见客户端也可以直接对比:

语言客户端/API作用开源许可证
Gojackc/pglogrepl + jackc/pgx复制协议、pgoutput 消息和 PostgreSQL 连接MIT
Ja vaPostgreSQL JDBC PGReplicationStreamJDBC 驱动内置的逻辑复制接口BSD-2-Clause
Pythonpsycopg2LogicalReplicationConnection 等逻辑复制接口LGPL-3.0-or-later
Rustsupabase/etl构建 PostgreSQL 逻辑复制与实时数据管道的 Rust 框架Apache-2.0
Node.jskibae/pg-logical-replication支持 pgoutputwal2json 等输出插件MIT
C/C++libpqPostgreSQL 官方底层客户端库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 系统是否可靠,可以归结为三点:全量与增量之间没有遗漏,数据确认位置能够安全恢复,下游在重复、延迟和乱序情况下仍能收敛到正确状态。

来源:https://juejin.cn/post/7665553541961121802

游乐网为非赢利性网站,所展示的游戏/软件/文章内容均来自于互联网或第三方用户上传分享,版权归原作者所有,本站不承担相应法律责任。如您发现有涉嫌抄袭侵权的内容,请联系youleyoucom@outlook.com。

同类文章
更多
自增主键值从何而来?深入理解原理,告别只会auto_increment

自增主键值从何而来?深入理解原理,告别只会auto_increment

KingbaseES推荐使用serial、bigserial、显式sequence或identity列实现自增主键。serial创建integer并关联序列,bigserial对应bigint;显式sequence可自定义起始值等参数;identity有generatedbydefault(允许指定值)与always(禁止)两种模式。

时间:2026-07-25 22:22
Linux下瀚高数据库授权文件过期及替换解决方案

Linux下瀚高数据库授权文件过期及替换解决方案

在银河麒麟系统下,瀚高数据库hgdb-4 5试用授权20天到期后需替换正式授权文件。正确操作:停止服务,备份旧文件,将授权文件复制到 opt highgo hgdb-4 5 etc lic 并命名为hgdb lic,设置权限600和属主highgo:highgo,再启动服务。禁止直接修改data目录下的license info文件。

时间:2026-07-25 22:22
Oracle BLOB实时同步的5大技术挑战与难点解析

Oracle BLOB实时同步的5大技术挑战与难点解析

OracleBLOB实时同步面临分片组装、多列隔离、长事务跨窗口、事务回滚及大对象资源控制等技术挑战,必须在日志中精确还原完整字段值,才能保证源端与目标端数据完全一致,这对同步系统的稳健性提出了高要求。

时间:2026-07-25 22:22
MySQL禁用redo日志导致全备失败

MySQL禁用redo日志导致全备失败

MySQL全量备份失败是由于数据定义语言操作触发排序索引构建,禁用重做日志导致XtraBackup无法获取一致性备份。测试验证表明,优化表语句即使无数据也会触发该问题。根本原因在于排序索引构建过程跳过了重做日志记录,破坏了备份的一致性。

时间:2026-07-25 20:35
Kafka架构图优化与改进的全面详细步骤与实践指南

Kafka架构图优化与改进的全面详细步骤与实践指南

Kafka作为实时数据流处理的核心中间件,其底层架构虽已相当成熟,但在实际生产环境中,要充分发挥其性能潜力,仍需落实到具体的调优与架构改造上。核心目标可归纳为三点:如何承载更高的吞吐量、如何保障数据不丢失、以及故障发生时如何快速恢复。本文将从这几个关键方向出发,深入探讨如何真正榨干Kafka集群的性

时间:2026-07-25 20:35
热门专题
更多
刀塔传奇破解版无限钻石下载大全 刀塔传奇破解版无限钻石下载大全
洛克王国正式正版手游下载安装大全 洛克王国正式正版手游下载安装大全
思美人手游下载专区 思美人手游下载专区
好玩的阿拉德之怒游戏下载合集 好玩的阿拉德之怒游戏下载合集
不思议迷宫手游下载合集 不思议迷宫手游下载合集
百宝袋汉化组游戏最新合集 百宝袋汉化组游戏最新合集
jsk游戏合集30款游戏大全 jsk游戏合集30款游戏大全
宾果消消消原版下载大全 宾果消消消原版下载大全
  • 热门数据榜