
简介本资源面向数据库开发与数据集成工程师提供一套基于PostgreSQL逻辑复制功能的实时数据变更捕获与同步系统源码。系统通过解析WAL日志捕获数据变更将其转换为可执行的SQL语句并借助Kafka消息队列实现PostgreSQL到异构数据源的实时同步适用于大数据分析、实时报表与数据仓库等场景。压缩包共24个文件约474KB以17个Java源文件为核心实现辅以2个XML配置、properties参数文件、README说明文档及配套说明材料结构清晰便于二次开发。目前已有57人学习下载。读者可获得完整的CDC同步实现思路包括WAL日志解析、SQL转换、Kafka发布订阅、容错与高可用设计等关键模块适合需要搭建跨平台数据同步管道的开发者参考借鉴。1. 从 WAL 到 Kafka一条让 PostgreSQL 变更实时流出去的链路线上跑着一套 PostgreSQL业务侧突然提了个需求订单表一有变更下游的风控、数仓、缓存要能秒级感知。你第一反应可能是写触发器往消息队列里塞或者定时扫表比对增量。前者侵入业务、DDL 一改就崩后者延迟高、还容易漏。真正稳的做法是让 PostgreSQL 自己把变更吐出来——逻辑复制就是干这个的。它把预写日志WAL里的行级变更解码成可读的 INSERT/UPDATE/DELETE再转成 SQL 或 JSON 推到 Kafka 这类异构数据源。这套方案适合做增量同步、CDC、缓存失效、审计留痕的团队前提是你愿意把复制槽、发布订阅、消息顺序这几件事吃透。下面按我实际落地的顺序从原理到参数到踩坑讲一遍。2. 逻辑复制凭什么能替代触发器WAL 解码的底层账2.1 物理复制和逻辑复制的分水岭PostgreSQL 的 WAL 本质是一串二进制记录记录的是「哪个数据页的哪个字节被改成了什么」。物理复制直接把这串字节流搬到备库要求主备完全一致跨版本、跨平台都不行。逻辑复制走的是另一条路它借助wal_level logical让 WAL 里额外带上足够的信息再通过输出插件output plugin把二进制记录翻译成逻辑层面的行变更。翻译出来的东西是「表 t 的 id5 这行name 从 a 变成了 b」跟数据页无关所以能跨版本、能只同步部分表、能落到完全不同的目标端。这个翻译过程有两个关键角色。一个是复制槽replication slot它记录消费者读到哪个 LSN 了保证主库不会把还没消费的 WAL 清掉——这是「不丢」的根基也是「磁盘涨满」的元凶。另一个是输出插件内置的pgoutput是给原生订阅用的输出的是逻辑复制协议格式如果你要自己解析成 SQL 或 JSON通常用test_decoding或者第三方插件比如wal2json先把变更取出来。2.2 发布订阅模型和它的边界原生逻辑复制用CREATE PUBLICATION/CREATE SUBSCRIPTION一对概念。发布端定义哪些表、哪些操作insert/update/delete/truncate要发订阅端连过去拉。它开箱即用但有几个硬边界你得先知道不复制 DDL。表结构变了订阅端不会自动跟着变得自己管。不复制序列sequence的当前值自增主键在目标端可能撞车。初始快照和增量之间有一致性窗口大表同步期间要小心。订阅端必须是 PostgreSQL。要落到 Kafka原生订阅帮不上得自己写消费者去读复制槽。所以「PostgreSQL 到 Kafka」这条链路主流做法是主库开逻辑复制 → 建一个逻辑复制槽 → 用程序或 Debezium 这类工具连上槽用输出插件拉变更 → 转成消息发到 Kafka。下面按这个思路落地。2.3 先把主库参数配对改postgresql.conf核心就三个# 必须开 logical否则连复制槽都建不了 wal_level logical # 给逻辑复制留的 WAL 发送进程上限一个槽至少占一个 max_replication_slots 10 # 发送进程总数要 槽数量 max_wal_senders 10改完重启。wal_level从replica调到logical需要重启实例这一步没有后悔药生产上要排窗口。max_replication_slots别抠门每个消费者一个槽槽满了新连接直接报错。改完用SHOW wal_level;确认返回logical才算数。提示wal_level logical会让 WAL 体积比replica略大写入密集的库要留意磁盘增速别等槽把盘撑爆才发现。3. 建槽、拉变更、转 SQL把 WAL 变成 Kafka 消息3.1 建一个逻辑复制槽并验证输出先建槽。槽名自己定插件用test_decoding方便调试生产上换成wal2json或直接用 Debezium 的pgoutput。-- 建槽plugin 决定输出格式这里先用 test_decoding 看原始变更 SELECT * FROM pg_create_logical_replication_slot(cdc_slot, test_decoding); -- 看一眼当前所有槽确认 active 和 restart_lsn SELECT slot_name, plugin, slot_type, active, restart_lsn FROM pg_replication_slots;建完槽对业务表做一次变更再用pg_logical_slot_get_changes把变更取出来。注意get_changes会消费并推进槽位取过就没了调试时想反复看用pg_logical_slot_peek_changes。-- 造一条变更 INSERT INTO orders (id, status) VALUES (1001, paid); UPDATE orders SET status shipped WHERE id 1001; -- 取出变更NULL 表示从槽当前位点开始 SELECT * FROM pg_logical_slot_get_changes(cdc_slot, NULL, NULL);返回的data字段就是解码后的文本类似table public.orders: INSERT: id[integer]:1001 status[text]:paid。这一步能跑通说明主库侧链路是活的。参数说明第一个参数是槽名第二个是起始 LSNNULL 从槽位点开始第三个是条数上限NULL 不限第四个往后是传给插件的选项。3.2 用 Python 消费槽并转成 SQL调试通了换成程序常驻消费。核心逻辑是循环调pg_logical_slot_get_changes把返回的 data 解析成结构化变更再拼成 SQL 或 JSON 发 Kafka。下面是最小可跑版本用psycopg2。import psycopg2 import json from kafka import KafkaProducer # 消费端连接必须用复制模式autocommit 打开 conn psycopg2.connect( hostpg-primary, port5432, dbnameapp, usercdc_user, password***, connection_factorypsycopg2.extras.LogicalReplicationConnection ) conn.autocommit True cur conn.cursor() producer KafkaProducer( bootstrap_servers[kafka1:9092, kafka2:9092], # 按主键哈希分区保证同一行的变更进同一分区顺序才不乱 key_serializerlambda k: k.encode(utf-8), value_serializerlambda v: json.dumps(v).encode(utf-8), acksall, # 等所有 ISR 确认防丢 enable_idempotenceTrue # 幂等防重试导致重复 ) while True: # 每次拉一批100 条或等 1 秒 cur.execute( SELECT data FROM pg_logical_slot_get_changes(%s, NULL, %s), (cdc_slot, 100) ) rows cur.fetchall() for (data,) in rows: # data 形如 table public.orders: INSERT: id[integer]:1001 ... # 真实项目建议用 wal2json 直接拿 JSON省去文本解析 msg parse_change(data) # 自定义解析函数 producer.send( pg.cdc.orders, keystr(msg[pk]), # 用主键做 key valuemsg ) producer.flush()逻辑说明LogicalReplicationConnection是 psycopg2 专门给复制协议用的连接工厂普通连接调不了pg_logical_slot_get_changes。acksall加enable_idempotenceTrue是防丢防重的组合拳代价是吞吐下降量大的库要权衡。key用主键是为了让同一行的变更落到同一分区——Kafka 只保证分区内有序key 选错同一行的 UPDATE 和 DELETE 可能乱序下游就翻车了。参数说明pg_logical_slot_get_changes的条数上限别设太大一次拉太多处理慢会拖住槽位点WAL 积压。我一般 100 到 500 之间配合flush频率调。3.3 用 wal2json 省掉文本解析test_decoding输出的是给人看的文本正则解析容易在字段含冒号、引号时崩。生产上换wal2json直接吐 JSON。-- 换成 wal2json 插件建槽 SELECT * FROM pg_create_logical_replication_slot(cdc_slot_json, wal2json); -- 带选项拉取include-pk 带上主键format-version 2 输出结构化 JSON SELECT data FROM pg_logical_slot_get_changes( cdc_slot_json, NULL, NULL, include-pk, 1, format-version, 2, include-transaction, 0 );返回的data直接是{change:[{kind:insert,schema:public,table:orders,columnnames:[id,status],columnvalues:[1001,paid]}]}。解析成本几乎为零字段类型也保住了。include-pk让你拿到主键做 Kafka key 正好。include-transaction设 0 可以不带事务边界信息简化下游处理但如果你要做事务一致性得设 1 并自己按事务聚合。注意wal2json是第三方插件得先在主库装好扩展包不同发行版包名不一样装完重启或CREATE EXTENSION视安装方式而定。别在生产上现装现试。4. 顺序、重复、积压CDC 链路的三个老大难4.1 Kafka 分区键选错导致乱序现象下游数仓里一条订单先出现shipped再出现paid状态倒流。原因生产者没设 keyKafka 轮询分区同一行的两条变更进了不同分区消费端并行处理就乱序了。解决必须用能唯一标识一行且不变的字段做 key通常是主键。如果表没有单列主键用联合主键拼字符串。别用会变的字段比如 status做 key那等于没做。4.2 消费失败重试引发重复写入现象Kafka 里同一条变更出现两次下游计数翻倍。原因消费者处理失败后没提交 offset重启后从上次位点重拉或者生产者acks没配好broker 没收到却以为成功客户端重试。解决下游写入做成幂等——用主键做 upsert或者维护一张已处理 LSN 的去重表。Kafka 侧开enable.idempotencetrue消费侧手动提交 offset 且「处理成功再提交」。4.3 复制槽不消费把磁盘撑满现象主库磁盘告警pg_replication_slots里某个槽active falserestart_lsn停在很久以前。原因消费者挂了或停了槽位点不推进主库为了不丢数据一直保留从该位点起的 WAL越积越多。解决监控槽的pg_wal_lsn_diff(pg_current_wal_lsn(), restart_lsn)超过阈值告警。消费者下线要主动pg_drop_replication_slot别留着。这是逻辑复制最经典的血泪坑没有之一。4.4 大事务把内存打爆现象消费者进程 OOM或者拉取时卡死。原因主库一个事务改了上百万行解码后一次性堆在内存里。解决pg_logical_slot_get_changes分批拉别一次NULL拉全部业务侧尽量避免超大事务批量操作拆成小批提交。这个坑在数据订正、批量导入时最容易踩。5. 避坑与排查上线前必须过的五道关5.1 槽建了但拉不到变更现象pg_logical_slot_get_changes返回空明明有写入。原因表没设REPLICA IDENTITYUPDATE/DELETE 时旧值拿不到或者表压根没被发布。解决给表设ALTER TABLE orders REPLICA IDENTITY FULL;或DEFAULT配合主键确认wal_level logical已生效。FULL会记录整行旧值WAL 变大只在需要旧值时用。5.2 目标端主键冲突现象同步 INSERT 报 duplicate key。原因初始快照和增量重叠或者序列没同步目标端自增主键和源端撞了。解决初始同步用一致性快照pg_export_snapshot配合SET TRANSACTION SNAPSHOT增量从快照对应的 LSN 开始序列值单独同步一次之后靠应用层保证。5.3 时间戳字段时区错乱现象源库timestamptz到目标端差 8 小时。原因解码输出的是 UTC目标端按本地时区解释。解决全链路统一用 UTC 存储展示层再转或者解码时显式带上时区信息别让下游猜。5.4 消费者重启后重复消费现象每次重启都重放一批老消息。原因offset 提交策略是自动提交重启时最后一批没提交成功。解决改手动提交处理完一批再commit配合幂等写入重复也不怕。5.5 监控缺失出事才发现现象槽积压几天没人知道直到磁盘满。原因没监控。解决至少盯三个指标——槽的restart_lsn落后量、消费者 lag、Kafka topic 积压。用pg_replication_slots加 Prometheus 采集设阈值告警。这套监控比同步本身还重要。6. 把链路做稳的两个进阶技巧第一个技巧是用 LSN 做端到端去重。每条变更都带一个 LSN把它一起写进 Kafka 消息和下游表。下游写入前先查这个 LSN 是否处理过处理过就跳过。这样即使消费者重启、Kafka 重投也不会重复。代价是下游多一次查询但换来的是精确一次语义值。实现上在消息体里加lsn: 0/16B3748下游用INSERT ... ON CONFLICT (lsn) DO NOTHING。第二个技巧是给复制槽做健康检查脚本定时跑发现异常自动告警甚至重建。下面这个查询能一眼看出哪个槽在拖后腿SELECT slot_name, active, pg_size_pretty(pg_wal_lsn_diff(pg_current_wal_lsn(), restart_lsn)) AS retained_wal, pg_wal_lsn_diff(pg_current_wal_lsn(), restart_lsn) AS lag_bytes FROM pg_replication_slots ORDER BY lag_bytes DESC;retained_wal超过几个 GB 就该查了。我一般设 5GB 告警10GB 人工介入。active false且 lag 在涨基本就是消费者挂了。还有个容易忽略的点DDL 变更的同步。逻辑复制不管 DDL表加字段后解码输出里会多出新字段下游解析如果写死了字段列表就会崩。我的习惯是下游解析用「按 columnnames 动态映射」别硬编码下标。加字段前先在目标端加好再在源端加顺序反了中间会有一段解析失败窗口。最后说个我自己的教训早期做这套链路图省事没设REPLICA IDENTITYUPDATE 和 DELETE 全丢旧值下游拿到的变更缺字段排查了两天才定位到。从那以后建表规范里第一条就是「参与 CDC 的表必须显式设 REPLICA IDENTITY」。这套方案值不值得做如果你的场景是「PostgreSQL 变更要实时到异构端」它比触发器干净、比轮询及时代价是要管好槽和顺序。把上面这些参数和坑过一遍基本能稳。希望帮到你。本文还有配套的精品资源点击获取