简介:本资源面向数据库开发与数据集成工程师,提供一套基于PostgreSQL逻辑复制功能的实时数据变更捕获与同步系统源码。系统通过解析WAL日志捕获数据变更,将其转换为可执行的SQL语句,并借助Kafka消息队列实现PostgreSQL到异构数据源的实时同步,适用于大数据分析、实时报表与数据仓库等场景。压缩包共24个文件,以17个Java源码为核心,辅以2个XML配置、properties参数文件及md、txt说明文档,整体约474KB,结构清晰便于二次开发。目前已有57人学习下载。读者可获得完整的CDC实现思路,包括WAL日志解析、SQL转换、Kafka发布订阅、事务回滚与故障恢复等关键模块,并附有说明文档与预览图辅助理解,适合需要构建跨平台数据同步方案的中高级开发者参考借鉴。
1. 从 WAL 到 Kafka:一条被低估的实时同步链路
线上库刚写入一条订单,三秒后风控系统就要拿到它做规则判定,五秒后数仓要落进明细表,十秒后搜索索引要能查到——这种场景下,定时轮询全表基本等于自杀,触发器写审计表又会把主库拖垮。基于 PostgreSQL 逻辑复制功能的实时数据变更捕获与同步系统,解决的正是这件事:不去业务表上加任何东西,而是让 PostgreSQL 自己把 WAL 日志里的行级变更吐出来,解析成带表名、操作类型、新旧值的结构化消息,再转成下游能吃的 SQL 或 JSON,推到 Kafka 这类异构数据源。它适合做 CDC 的中间件开发者、做数仓实时入仓的数据工程师,也适合被"主从延迟 + 异构同步"折磨过的 DBA。核心链路就三段:逻辑复制槽产出变更流 → 解析并组装成消息 → 投递到 Kafka。下面按这条链路拆开讲。
2. 逻辑复制与 WAL 解析:先搞懂数据从哪来
2.1 逻辑复制槽到底吐出了什么
PostgreSQL 的物理复制传的是 WAL 原始字节,接收端只能还原成同样的数据页,没法跨版本、跨异构。逻辑复制走的是另一条路:通过pgoutput这类输出插件,把 WAL 里的记录解码成逻辑意义上的 INSERT / UPDATE / DELETE,带上表 OID、列值、事务边界。这个解码结果不是 SQL 文本,而是一套二进制协议消息,常见的有Begin、Relation、Insert、Update、Delete、Commit。Relation消息很关键,它携带表的列定义,解析端必须先缓存它,否则后面拿到列值时不知道哪一列对应哪个字段。
复制槽(replication slot)是这套机制的锚点。它记录消费者已经确认到哪个 LSN,保证 PostgreSQL 不会提前回收还没被消费的 WAL。这既是可靠性来源,也是最大的坑:槽一旦创建,即使没有消费者,WAL 也会一直堆积,磁盘迟早爆。所以任何生产方案都必须有槽的监控和清理策略。
2.2 开启逻辑复制的最小配置
先确认wal_level是logical,这是前提,改完要重启。
# 查看当前 wal_level psql -c "SHOW wal_level;" # 若不是 logical,编辑 postgresql.conf # wal_level = logical # max_replication_slots = 10 # 按消费者数量留余量 # max_wal_senders = 10 # 至少大于槽数量 # 改完重启实例 pg_ctl restart -D $PGDATA参数说明:max_replication_slots决定能同时存在多少个槽,每个下游消费者通常占一个;max_wal_senders是 WAL 发送进程上限,必须大于等于槽数量,否则创建槽会失败。这两个值调小容易,调大要重启,规划时宁可多留。
2.3 建表、建槽、验证变更流
逻辑复制对表有硬性要求:必须是有主键或 REPLICA IDENTITY 的表,否则 UPDATE / DELETE 无法定位行。
-- 业务表,必须有主键 CREATE TABLE orders ( id bigserial PRIMARY KEY, user_id bigint NOT NULL, amount numeric(12,2) NOT NULL, status text NOT NULL DEFAULT 'created', updated_at timestamptz NOT NULL DEFAULT now() ); -- 创建逻辑复制槽,使用内置 pgoutput 插件 SELECT * FROM pg_create_logical_replication_slot('cdc_orders_slot', 'pgoutput'); -- 查看槽状态,确认 active 和 confirmed_flush_lsn SELECT slot_name, plugin, active, restart_lsn, confirmed_flush_lsn FROM pg_replication_slots;逻辑说明:pg_create_logical_replication_slot第二个参数是输出插件名,内置的pgoutput不需要额外安装,配合CREATE PUBLICATION使用最省事。建完槽后,用pg_logical_slot_get_changes可以手动拉一批变更验证链路是否通:
-- 先做一次写入 INSERT INTO orders (user_id, amount) VALUES (1001, 99.50); -- 手动消费变更(会推进 confirmed_flush_lsn,测试环境用) SELECT * FROM pg_logical_slot_get_changes( 'cdc_orders_slot', NULL, NULL, 'proto_version', '1', 'publication_names', 'pub_orders' );注意:pg_logical_slot_get_changes会消费并推进槽位,生产环境别拿它做调试,否则真实消费者会丢数据。调试用pg_logical_slot_peek_changes,它只读不推进。
2.4 用 Publication 圈定同步范围
不建 publication 也能用槽,但pgoutput需要它来指定同步哪些表。
-- 只同步 orders 表,避免把整库变更都推下去 CREATE PUBLICATION pub_orders FOR TABLE orders; -- 后续要加表 ALTER PUBLICATION pub_orders ADD TABLE payments;选型理由:按表建 publication 而不是FOR ALL TABLES,是因为下游异构系统往往只关心部分业务表,全库同步会把无关变更也塞进 Kafka,topic 膨胀、消费端过滤成本高。粒度控制在源头做,比在消费端做便宜得多。
3. 把变更流组装成 Kafka 消息:解析与投递
3.1 解析 pgoutput 消息的字段映射
拿到二进制流后,解析端要按协议逐条读。以 Python 为例,用psycopg2的逻辑复制游标能直接拿到解码后的消息,省去手写协议解析。
import psycopg2 import psycopg2.extras import json # 关键:connection_factory 用逻辑复制专用工厂 conn = psycopg2.connect( dbname="appdb", user="repl", password="***", host="10.0.0.10", port=5432, connection_factory=psycopg2.extras.LogicalReplicationConnection ) cur = conn.cursor() cur.start_replication( slot_name='cdc_orders_slot', options={'proto_version': '1', 'publication_names': 'pub_orders'}, decode=True # 让 psycopg2 帮忙解码成 dict ) def handle_msg(msg): payload = msg.payload # 已是 dict,含 action / schema / table / columns # 只处理数据变更,跳过 begin/commit 心跳 if payload.get('action') in ('insert', 'update', 'delete'): record = { 'op': payload['action'], 'table': payload['table'], 'data': {c['name']: c['value'] for c in payload['columns']}, 'lsn': str(msg.data_start) } produce_to_kafka(record) msg.cursor.send_feedback(flush_lsn=msg.data_start) # 确认消费位点 cur.consume_stream(handle_msg)逻辑说明:decode=True让 psycopg2 把 pgoutput 二进制转成 Python dict,字段结构里action是操作类型,columns是列数组。send_feedback是命门——只有调用它,PostgreSQL 才会推进confirmed_flush_lsn,否则槽位不动、WAL 堆积。参数上flush_lsn传当前消息的data_start,表示"这条我已处理完"。
3.2 转成 SQL 还是 JSON:下游决定格式
标题里提到"转换为 SQL 语句",这在异构同步里确实常见——目标端是另一个关系库时,直接重放 SQL 最省事。但推到 Kafka 时,JSON 更通用。两种都给你。
def to_sql(record): t = record['table'] d = record['data'] if record['op'] == 'insert': cols = ', '.join(d.keys()) vals = ', '.join(f"'{v}'" if isinstance(v, str) else str(v) for v in d.values()) return f"INSERT INTO {t} ({cols}) VALUES ({vals});" if record['op'] == 'update': sets = ', '.join(f"{k}='{v}'" for k, v in d.items() if k != 'id') return f"UPDATE {t} SET {sets} WHERE id={d['id']};" if record['op'] == 'delete': return f"DELETE FROM {t} WHERE id={d['id']};" def to_json(record): return json.dumps(record, ensure_ascii=False, default=str)参数说明:to_sql里对字符串值加引号、对数字不加,是最朴素的类型判断,生产环境要按列的真实类型走映射表,别用isinstance硬猜,numeric、timestamptz、jsonb 都会翻车。to_json的default=str兜底处理 datetime 和 Decimal,避免序列化报错。
3.3 投递到 Kafka 的幂等与顺序
Kafka 生产者要开幂等,否则重试会写出重复消息。
from kafka import KafkaProducer producer = KafkaProducer( bootstrap_servers=['kafka1:9092', 'kafka2:9092'], acks='all', # 所有 ISR 确认,防丢 enable_idempotence=True, # 幂等,防重 max_in_flight_requests_per_connection=5, # 幂等开启时上限就是 5 key_serializer=lambda k: k.encode('utf-8'), value_serializer=lambda v: v.encode('utf-8') ) def produce_to_kafka(record): # 用主键做 key,保证同一行变更落到同一分区,顺序不乱 key = str(record['data'].get('id', '')) producer.send('cdc.orders', key=key, value=to_json(record))逻辑说明:acks='all'配合enable_idempotence=True是防丢防重的标准组合,代价是延迟略高。用主键做分区 key 是关键——同一行的 INSERT、UPDATE、DELETE 必须进同一分区,否则消费端看到的顺序是乱的,更新可能先于插入到达。max_in_flight_requests_per_connection在幂等开启时不能超过 5,超了会直接报配置错误。
3.4 消费端如何保证不重复落库
Kafka 至少一次投递,消费端必须自己幂等。目标端是关系库时,用主键 UPSERT。
-- 目标端 PostgreSQL / MySQL 通用思路:按主键 upsert INSERT INTO orders (id, user_id, amount, status, updated_at) VALUES (?, ?, ?, ?, ?) ON CONFLICT (id) DO UPDATE SET user_id = EXCLUDED.user_id, amount = EXCLUDED.amount, status = EXCLUDED.status, updated_at = EXCLUDED.updated_at;DELETE 操作则直接按主键删,重复删不报错即可。这样即使同一条消息被消费两次,结果也一致。
4. 避坑与排查:那些让同步链路半夜报警的细节
4.1 复制槽不推进,磁盘被 WAL 撑爆
现象:pg_replication_slots里active=false,restart_lsn长时间不动,pg_wal目录疯涨。原因:消费者挂了或没调send_feedback,槽位不推进,PostgreSQL 保留所有未确认 WAL。解决:先确认消费者存活,再检查代码里是否每条消息都发了 feedback;确认不再需要的槽用SELECT pg_drop_replication_slot('槽名')删掉。监控上给pg_replication_slots的restart_lsn和当前 LSN 差值设告警,比看磁盘更早发现问题。
4.2 UPDATE 拿不到旧值,下游无法做差异比对
现象:解析出的 UPDATE 只有新值,没有变更前的旧值。原因:表的 REPLICA IDENTITY 默认是DEFAULT,只带主键,不带旧列值。解决:需要旧值时把表设为ALTER TABLE orders REPLICA IDENTITY FULL;,代价是 WAL 体积变大,只对确实需要比对的表开。
4.3 大事务把内存打满
现象:解析进程内存飙升甚至 OOM。原因:一个事务里批量更新几十万行,Begin到Commit之间的消息全堆在内存里等提交。解决:解析端按事务边界流式处理,别把整个事务缓存成列表;业务侧尽量把大事务拆小,单事务变更行数控制在万级以内。
4.4 表结构变更后解析错位
现象:加了一列之后,解析出的字段对不上,值串位。原因:Relation消息缓存了旧列定义,DDL 后没刷新。解决:解析端收到新的Relation消息时必须覆盖缓存,别只认第一次;DDL 变更尽量在低峰做,变更后观察一批消息确认字段正确。
4.5 Kafka 消息延迟高,消费端积压
现象:kafka-consumer-groups显示 lag 持续增长。原因:分区数太少,单分区吞吐到顶;或消费端单条处理太慢。解决:topic 分区数按峰值吞吐规划,一般不少于消费者线程数;消费端批量拉取、批量 upsert,别一条一条提交。分区 key 用主键时,热点行的变更会集中到一个分区,这是顺序性的代价,接受它或改用表名 + 主键哈希。
5. 进阶:用 LSN 做断点续传与一致性校验
链路跑通只是开始,真正决定这套系统能不能上生产的,是断点续传和一致性校验。复制槽本身帮你记住了消费位点,但消费者重启后从哪继续、怎么确认没丢没重,得自己设计。
一个实用技巧是把 LSN 落进下游。每条消息投递时带上lsn字段,消费端处理完把最大 LSN 写进一张cdc_checkpoint表。重启时先读这张表,再决定从 Kafka 的哪个 offset 或从槽的哪个位置继续。这样即使 Kafka 和 PostgreSQL 两侧的位点对不上,也有个统一的进度基准。
CREATE TABLE cdc_checkpoint ( consumer text PRIMARY KEY, last_lsn pg_lsn NOT NULL, updated_at timestamptz NOT NULL DEFAULT now() ); -- 消费端每批处理后更新 INSERT INTO cdc_checkpoint (consumer, last_lsn) VALUES ('orders_sync', '0/1A2B3C4D') ON CONFLICT (consumer) DO UPDATE SET last_lsn = EXCLUDED.last_lsn, updated_at = now();一致性校验则定期做:拿源表和目标表按主键比对行数和关键字段校验和。行数对不上,多半是 DELETE 丢了;字段校验和对不上,多半是 UPDATE 的旧值/新值处理有问题。校验别全表扫,按时间分区抽样,比如每天校验最近一小时变更过的行。
| 校验项 | 方法 | 常见偏差原因 |
|---|---|---|
| 行数 | 按主键 count 比对 | DELETE 未同步、重复消费 |
| 字段值 | 关键列 md5 聚合比对 | UPDATE 新值解析错、类型转换丢精度 |
| 变更时序 | 比对 updated_at 最大值 | 分区 key 不当导致乱序 |
我自己的习惯是:任何 CDC 链路上线前,先跑一周影子同步,源库正常写,目标库只读比对,确认零偏差再切流量。这套系统最贵的不是写代码,是上线后半夜被 WAL 撑爆磁盘叫醒。把槽监控、LSN 落库、定期校验这三件事做扎实,比堆任何花哨功能都值。希望帮到你。
本文还有配套的精品资源,点击获取