ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

PostgreSQL逻辑复制实战:WAL解析到Kafka的CDC同步方案

PostgreSQL逻辑复制实战:WAL解析到Kafka的CDC同步方案 简介这是一份围绕PostgreSQL逻辑复制机制实现的实时数据变更捕获与同步系统源码适合数据工程师、后端开发人员以及需要构建异构数据实时同步管道的技术团队。系统通过解析WAL日志并将变更转换为SQL语句借助Kafka消息队列实现解耦和异步分发覆盖从数据捕获、转换到跨平台同步的完整链路同时支持扩展至NoSQL、数据仓库等异构数据源场景。压缩包共含24个文件其中Java源码多达17个核心逻辑清晰另有XML配置、properties配置文件、Markdown说明、TXT与DOCX文档及一张示意图便于快速读懂项目结构与运行方式。资源整体仅474KB轻量紧凑但实现了核心功能框架。目前已有57人学习适合用于研究PostgreSQL逻辑复制、WAL日志解析及Kafka集成场景也可作为二次开发或课程设计参考直接导入IDE结合构建文件即可调试运行。1. 为什么把同步做进 WAL 层从双写翻车到逻辑复制做过数据同步的团队大概率都经历过两条路应用层双写业务代码里多一条 MQ 发送逻辑上线三个月后因为某个事务漏发导致两边数据对不上定时批处理凌晨跑一次 ETL数仓里的数据永远滞后一天。等到要接实时数仓、缓存更新、订单风控这类场景才发现需要的是一个能把数据库变更在秒级内送出去的通道而不是靠业务代码自觉。PostgreSQL 的逻辑复制就是这条通道的底座它直接从 WAL 日志里解析出每一行数据的插入、更新、删除把变更还原为结构化事件再推给 Kafka最后由任意异构数据源消费。这套方案把同步逻辑从业务层下沉到了数据库层与触发器和双写无关适合正在做异构同步选型、或者已经被双写搞得焦头烂额的团队。2. 逻辑复制是怎么把 WAL 变成变更事件的2.1 物理复制与逻辑复制的本质差异很多人在理解逻辑复制时会先想到一主一备的物理复制。物理复制就像把主库的磁盘块变化原样搬到备库备库拿到 WAL 后直接重放双方的二进制状态完全一致。这种复制对备库的要求很高备库必须也是 PostgreSQL而且要保证同样的版本和编译选项否则 WAL 重放就可能失败。更重要的是物理复制默认是整库级、集群级的你想把某几张表的数据同步到另一个异构系统物理复制根本插不上手。逻辑复制则是在 WAL 之上加了一层翻译。它不关心页面怎么变只关心逻辑上的变更哪张表、哪个主键、什么操作、字段值是多少。PostgreSQL 9.4 引入逻辑解码机制10.0 把 Publication/Subscription 做成正式功能到了现在的版本逻辑复制已经成为把 PostgreSQL 数据送出去最稳的出口。配置逻辑复制需要的参数很少但每个都关键。装好 PostgreSQL 之后默认的 wal_level 是 replica只支持物理复制必须改成 logical 并重启实例ALTER SYSTEM SET wal_level logical; ALTER SYSTEM SET max_replication_slots 10; ALTER SYSTEM SET max_wal_senders 10; -- 需要重启 PostgreSQL 实例生效这三行参数的作用wal_level 决定 WAL 里是否写入逻辑解码所需的额外信息max_replication_slots 限制同时存在的复制槽数量每个连接器或消费端占一个槽max_wal_senders 则是允许多少个 WAL 发送进程同时工作。重启后可以用 SHOW wal_level 确认看到 logical 说明已经就绪。这一步是整套方案的入口遗漏任何一个参数后面创建复制槽时都可能报错错误信息还不太直观我见过有人在这里卡了几个小时。2.2 WAL 里没有 SQL逻辑解码与元组变更这里必须纠正一个常见误解逻辑复制并不是从 WAL 里读出原始 SQL。PostgreSQL 的 WAL 记录的是页面的物理变更比如某个数据块在偏移量 X 处写入了一段二进制数据它根本不关心这条变更来自哪条 SQL、由谁执行的。逻辑解码要做的事情是把这种物理变更反向还原成逻辑变更这张表的主键是 42UPDATE 操作把 name 字段从 A 改成了 B。还原出来的是一个结构化的变更事件而不是可执行的文本 SQL。输出插件就是干这个翻译活的人。PostgreSQL 内置了 test_decoding 和 pgoutput 两个插件社区里最常用的是 wal2json。test_decoding 输出纯文本格式适合调试pgoutput 是二进制协议效率高是 Publication/Subscription 和 Debezium 等工具的默认选择wal2json 输出 JSON 文本可读性强适合自研消费端时做解析。无论用哪个插件本质都是一样的插件订阅 WAL 里的逻辑解码消息把它转换成特定格式的输出流。这里还有一个容易被忽略的点逻辑解码只能在有复制槽的情况下工作。复制槽记录了消费端当前消费到了哪个 LSN日志序列号PostgreSQL 会认为这个位置之前的所有 WAL 都还有消费者需要因此不会被自动清理掉。这是它能实现断点续传的基础但也是磁盘爆满事故的根源一个长期没有消费者拉取的复制槽会在后台悄悄把磁盘占满这个话题在第 5 章单独展开。2.3 发布订阅模型与复制槽的关键角色理解了输出插件接下来看逻辑复制的工作模型。发布Publication定义哪些表允许被复制订阅Subscription定义谁来消费、消费什么复制槽则负责记录消费进度。三者缺一不可。在测试环境里可以用一条命令查看当前实例上所有的复制槽SELECT slot_name, plugin, slot_type, active, restart_lsn FROM pg_replication_slots;slot_name 是槽的名字plugin 是解码插件名称active 表示当前是否有消费者连在上面。如果 active 长期为 false就要警惕了说明没有任何进程在消费它但它仍然在阻止 WAL 清理。如果是自研方案也可以用 pg_recvlogical 手动建槽和消费pg_recvlogical -h localhost -U repl_user -d mydb \ --slot demo_slot --create-slot -P pgoutput # 监听并输出逻辑变更-o 参数传给输出插件 pg_recvlogical -h localhost -U repl_user -d mydb \ --slot demo_slot --start -o proto_version1 -o publication_namespub_all操作前需要准备一个具备 REPLICATION 权限的角色。pg_recvlogical --start 是阻塞式的会一直等待新的变更并打印到终端测试时看到初始化输出就可以 CtrlC 结束槽和 LSN 会保留下次再启动仍然从上次的位置继续读。这段命令展示了逻辑复制的底层层面但生产环境很少有人直接用 pg_recvlogical因为它只解决了能读到变更这一步后面的序列化、分发、断线重连都要自己写。于是有了更完整的落地工具。3. 把 WAL 解析成可下发的 SQL插件选型与最小实现3.1 pgoutput、wal2json、Debezium 的定位差异把 WAL 变成变更事件只是第一步真正让这套系统落到生产还得解决三个问题怎么把二进制变更转成方便下游消费的形式、怎么处理断线重连时的位点管理、怎么把变更稳定地写入 Kafka。pgoutput 和 wal2json 都只是输出插件它们解答了第一个问题后面的问题需要框架来兜底。常见做法是直接上 Debezium。它是一套完整的 CDC 框架内置了 PostgreSQL 连接器底层默认使用 pgoutput 插件把 WAL 变更解析成固定的 JSON 消息结构通过 Kafka Connect 写入 Kafka 集群。相比自己用 wal2json 写消费端Debezium 把最有难度的部分都封装好了快照机制、位点管理、schema 变更处理、断线恢复。它不是一个输出插件而是一个完整的集成框架。wal2json 适合的场景是你只是想快速在本地看一下变更长什么样或者团队想完全掌控解析逻辑不愿意引入 Kafka Connect 这层依赖。方案输出格式位点管理生产可用性适用场景pgoutput二进制需自建中基于 Publication/Subscription 的原生复制wal2jsonJSON 文本需自建中自研消费端、调试观察DebeziumJSON/AVRO内置高对接 Kafka 的正式同步链路3.2 用 Debezium 把变更事件接入 Kafka最小可跑配置前置条件是一套能跑的 Kafka 集群外加一个 Kafka Connect 节点。如果你只搭过单机 Kafka 没装过 Connect确认 kafka-connect 的插件目录里放入了 Debezium 的 PostgreSQL 连接器 JAR 包即可。连接器以 REST API 的方式注册到 Connect 集群配置如下{ name: pg-cdc-connector, config: { connector.class: io.debezium.connector.postgresql.PostgresConnector, database.hostname: localhost, database.port: 5432, database.user: repl_user, database.password: secret, database.dbname: mydb, topic.prefix: cdc, table.include.list: public.users,public.orders, plugin.name: pgoutput, slot.name: debezium_cdc_slot, publication.name: debezium_cdc_pub, value.converter: org.apache.kafka.connect.json.JsonConverter, value.converter.schemas.enable: false } }注册命令如下curl -X POST http://localhost:8083/connectors \ -H Content-Type: application/json \ --data pg-connector.json curl http://localhost:8083/connectors/pg-cdc-connector/status几个参数的取舍要注意。topic.prefix 会出现在最终 topic 命名里后续 topic 名就是topic.prefix.schema.table一旦上线后不要改改了等于所有下游重新接一遍。plugin.name 用 pgoutput 是首选它是 PostgreSQL 内置插件不需要在数据库服务器上安装额外东西wal2json 需要在服务器上编译安装能不用就不用。slot.name 和 publication.name 对应数据库里的复制槽和发布对象Debezium 在首次启动时会自动创建如果名字和已有对象冲突会直接报错。这条链路启动后在数据库里执行一条 INSERTKafka 里立刻能看到对应消息。用控制台消费验证最快kafka-console-consumer.sh \ --bootstrap-server localhost:9092 \ --topic cdc.public.users \ --from-beginning消息体是一个 JSON包含 schema 和 payload 两块。如果配置里关掉了 schemas.enable则只输出 payload结构更轻量下游解析成本低。生产环境如果对消息体积敏感可以先关掉 schema 输出但要自己维护字段变更的兼容性这是一笔要算清楚的账。3.3 生成的 SQL 为什么不能直接回放看到 Kafka 里的变更消息后有些团队会兴奋地写一个消费者把 UPDATE 事件转成 SQL 塞进目标库结果在目标库上疯狂报错。原因很简单逻辑复制还原出来的是变更后的数据状态不是可以安全重放的 SQL。一条 UPDATE 消息里包含的是主键、变更后的字段值、以及变更前的旧值要生成 SQL 得自己拼。拼 SQL 本身不难UPDATE public.users SET email after.email WHERE id before.id;但是有几个坑。第一如果 UPDATE 语句修改了主键本身before.id 和 after.id 不同生成的 UPDATE 语句 WHERE 条件用哪个 id 都对不上常见做法是把原主键找出来用原主键 WHERE、新主键赋值但框架不会替你判断主键是否变了。第二目标库必须有主键约束否则 UPDATE 只能退化成先 DELETE 再 INSERT而这是最伤性能的写法。第三如果源表没有主键WAL 里根本不会记录完整旧值需要依赖 REPLICA IDENTITY FULL 才能把旧行的所有字段写进日志代价是每个变更都携带整行数据写放大非常明显。这些问题留在第 5 章展开反正记住一点逻辑复制只保证把变更告诉你不保证你转出来的 SQL 在目标端一定能正确执行。4. Kafka 作为分发中枢topic 设计、消息格式与 Lag 排查4.1 topic 划分与分区键设计Kafka 在这套系统里的角色是变更事件的交换总线。源端只负责把 WAL 变更写进 Kafka目标端各自按需消费因此 topic 怎么划、消息怎么落键直接决定了下游扩展性和消费性能。Debezium 默认的 topic 命名是topic.prefix schema table也就是前面配置里的cdc.public.users。默认规则下一张表对应一个 topic这种设计贴合大部分场景下游想单独处理某个表直接订阅对应 topic 即可互不干扰。如果表特别多每个 topic 只有一个分区那消费并行度就顶死在一了。建议给高频表单独设定分区数比如 users、orders 这类核心表设 8 到 16 个分区日志型、低频表保持单分区即可。分区键的选择决定了同一行的多次变更是否落在同一分区从而保证顺序。Debezium 默认用表主键作为消息 key如果你手动写生产者务必沿用主键或业务唯一键做 key这样同一条记录的所有变更都会路由到同一分区消费端能按顺序重放。topic 计划里另一个常见的需求是表合并比如把多个表的变更汇总到一个 topic减少下游订阅数。Kafka Connect 的 RegexRouter 可以做到transforms: route, transforms.route.type: org.apache.kafka.connect.transforms.RegexRouter, transforms.route.regex: cdc\\.(.*)\\.(.*), transforms.route.replacement: all_changes这个配置把所有prefix.schema.table格式的 topic 路由到同一个 all_changes代价是下游要自己根据消息里的表名字段区分来源失去了按表独立扩容的能力。生产环境我一般不建议路由到单 topic除非下游是 Flink 之类需要统一流入口的引擎那里的场景另当别论。4.2 before/after 消息结构与消费端幂等Kafka 里的每一条变更消息都包含操作类型和前后镜像。以一个典型的 UPDATE 为例关闭 schema 输出后的 payload 结构如下{ before: { id: 1001, email: oldexample.com, status: 1 }, after: { id: 1001, email: newexample.com, status: 2 }, source: { lsn: 63125392, ts_ms: 1715112000000, table: users }, op: u }op 字段标识操作类型c 表示插入u 表示更新d 表示删除r 表示快照阶段的初始读取。source.lsn 是这条变更在源库 WAL 里的位置source.ts_ms 是变更在源库发生的时间戳。消费端做幂等最直接的办法就是记录每个主键最后处理到的 LSN当再次收到相同或更小的 LSN 时直接跳过。Kafka 在默认配置下本身是至少一次投递重复消费是常态不加幂等就是给自己埋雷。消费端的提交策略要和幂等配合。推荐的做法是把自动提交关掉业务处理成功后再手动提交 offset这样崩溃后重新消费的只是未提交的部分配合按主键的 LSN 去重即使有重复也不影响最终结果from kafka import KafkaConsumer import json consumer KafkaConsumer( cdc.public.users, bootstrap_serverslocalhost:9092, group_idcdc-syncer, enable_auto_commitFalse, value_deserializerlambda m: json.loads(m) ) for msg in consumer: payload msg.value # 这里执行目标库的写入逻辑 apply_to_target(payload) # 处理成功后手动提交保证 at least once 语义 consumer.commit()这里的手动提交会带来一个副作用如果 apply_to_target 成功但 commit 前进程崩溃重启后会重新消费这一批消息所以目标库写入语句必须天然幂等比如 INSERT ... ON CONFLICT DO UPDATE。只依赖提交机制而不处理重复业务数据早晚会在某个凌晨被重复账单或重复订单教育一次。4.3 消费 Lag 的排查路径链路跑久了最常见的故障就是消费 Lag 越来越高消息延迟从秒级变成分钟级甚至小时级。排查这类问题有一套固定顺序第一步永远是看 group 的积压情况。kafka-consumer-groups.sh \ --bootstrap-server localhost:9092 \ --describe --group cdc-syncer输出里的 LAG 一列表示每个分区还剩多少条消息没被消费。如果 LAG 持续增大优先检查消费端自身的处理耗时。很多人忽略了分区数量对消费并发的限制一个分区只能被同一个 group 里的一个消费者线程消费如果你只有 1 个分区却有 10 个消费者线程剩下 9 个全是闲置的。单消费者处理一条消息如果耗时 100ms每秒最多处理 10 条而源库高峰期一秒写入几百条Lag 不涨才怪。还有一种隐蔽的情况Kafka 可视化工具面板上看到 LAG 正常但业务端仍然觉得延迟高。这时候要检查的是 Producer 端到 Connect 端是否有积压有时是 Kafka Connect 的 worker 数不够或者 Debezium 的 max.queue.size 被改小了导致它从 WAL 读取的速度赶不上写入速度。检查 Debezium 连接器日志时多关注 slow source request 和 WAL 读取相关的 warning这两个关键字出现说明它已经跟不上源库的变更速度了需要优先扩容而不是纠结消费者。5. 同步链路避坑复制槽膨胀、大事务、DDL 与无主键表5.1 复制槽不消费把磁盘打爆现象数据库磁盘使用率曲线在几周内持续上涨排查后发现 WAL 目录占了几百 GB删除归档日志也没有用。原因复制槽保存着一个 restart_lsnPostgreSQL 会认为是复制槽的消费者还没处理到那个位置因此这个位置之前的 WAL 都不能清理。如果复制槽创建了但从不消费或者消费者断连后一直没有恢复数据库就会默默积累 WAL 文件直到磁盘满。这个问题在测试环境特别常见开发和测试人员建了槽忘了删几个月后磁盘悄悄被填满。解决先用一条 SQL 找出是哪个槽在占空间再决定是修复消费端还是删除槽。下面是常用的排查语句SELECT slot_name, active, pg_size_pretty(pg_wal_lsn_diff(pg_current_wal_lsn(), restart_lsn)) AS retained_wal FROM pg_replication_slots ORDER BY pg_wal_lsn_diff(pg_current_wal_lsn(), restart_lsn) DESC;确认是无用槽后直接删除SELECT pg_drop_replication_slot(debezium_cdc_slot);生产环境的防御性配置是在 postgresql.conf 里设置 max_slot_wal_keep_size比如限制为 8GB超过后 PostgreSQL 会把该复制槽标记为失效避免磁盘被打满。代价是后续消费者需要重建槽并重新完成一次全量快照但这比整个数据库磁盘不可用要好得多。这个参数当初我以为是可有可无的直到线上实例因为一个未消费的槽差点故障转移从那以后每个实例我都强制要求配置。5.2 大事务刷屏把下游压垮现象某一时刻源库只执行了一条 UPDATE更新了上百万行紧接着 Kafka 消息量暴增消费端 Lag 从个位数涨到几十万broker 的带宽也被打满。原因逻辑解码是按行输出的。一条 UPDATE 影响一百万行WAL 里就有一百万个元组变更事件Debezium 会把这百万条消息全部发送到 Kafka。下游面对瞬时流量洪峰要么消费不过来要么写目标库时主从延迟变大最终表现为整体链路雪崩。解决从源头上减少大事务出现的频率。业务侧如果能改造尽量把大 UPDATE 拆成小批量提交每批几千行这个收益是最直接的。如果业务侧不能改Debezium 里有一个参数 max.batch.size 可以限制每次读取的变更行数但要注意它只是一次轮询的限制不会把一条 UPDATE 拆开。真正能缓解下游压力的做法是让下游具备按表合并变更的能力比如只关心最终状态的目标系统可以跳过部分中间变更类似 Flink 的去重聚合。对于非关键场景直接接受大事务带来的短时延迟给它几分钟追平比强行优化更省事。5.3 表结构变更导致解析中断现象某天执行了 ALTER TABLE ADD COLUMN 之后Debezium 连接器进入 FAILED 状态日志里报 schema 变更相关的错误重试多次也起不来。原因Debezium 在内存里维护了一张基于 WAL 事件推断出来的表结构。当源库表结构变更而连接器配置没有同步更新时它会认为随后到来的变更事件与预期结构不匹配从而停止解析。这不是 PostgreSQL 逻辑复制的问题而是连接器对 schema 变化敏感。PostgreSQL 自身处理 DDL 的策略是DDL 本身不通过逻辑复制传播订阅端需要自己手动同步结构变更。解决ALTER TABLE 尽量放在业务低峰期执行执行前把连接器停掉结构变更同步到目标端之后再重启连接器。重启后它会在上次记录的一致点继续解析一般不会丢数据。如果你用的 Debezium 版本较新连接器会发出 schema change event可以监听它来自动更新下游结构但这是需要额外开发量的别指望开箱即用。最省心的办法是屏蔽 DDL 权限只允许通过受控的变更流程去改表这条链路的稳定性会立刻上一个台阶。5.4 无主键表与异构类型翻车现象同步到目标端的数据插入正常更新和删除却频繁报错或者影响行数为 0数据对账怎么都对不上。原因更新和删除操作在目标库定位行时需要依赖源库 WAL 里记录的旧值作为 WHERE 条件。PostgreSQL 默认只记录主键或非空的唯一索引作为定位条件。如果表没有主键WAL 变更消息里的 before 字段可能只有空值下游拿到空 WHERE 条件自然无法更新。另一个常见问题是类型映射PostgreSQL 的 timestamptz 在 JSON 里是带时区的 ISO 字符串目标端是 SQL Server 时需要转成 datetimeoffset转错时区差就是八小时的数据错误numeric 类型转成 double 会丢精度多位数金额同步过去再对账结果一定对不上。解决没有主键的表用 REPLICA IDENTITY FULL 强制 WAL 记录整行旧值。用这个命令前要考虑清楚它会让每个变更事件体积增大几倍日志写入量也明显上升只对确有必要保留的历史表开这个开关ALTER TABLE public.history_log REPLICA IDENTITY FULL;类型映射方面Debezium 的 decimal.handling.mode 设为 string把 DECIMAL/NUMERIC 用字符串传输让目标端自行解析避免精度丢失timestamp 类字段尽量以字符串形式下发由下游显式转换。另外注意 SQL Server 的 datetime 类型只精确到 3.3 毫秒而 PostgreSQL 的 timestamptz 精确到微秒截断误差在每天几百万条变更的场景下会积累成秒级的延迟差对账时要用容差而不是精确相等。6. 用 LSN 与对账把这条链路验证死一个落地收尾的技巧链路搭建完成后最怕的不是功能不可用而是你无法证明它是对的。验证方案要我推荐首选是 LSN 定位法它把源库的 WAL 位置和 Kafka 消息直接关联起来能精确到某一条变更。具体做法是在源库执行一个事务插入一条带标记的记录同时记录事务开始时的 LSNBEGIN; INSERT INTO sync_check(id, marker, created_at) VALUES (nextval(sync_check_id_seq), marker-20240501, now()); SELECT pg_current_wal_lsn(); COMMIT;假设返回的 LSN 是 63125392那么到 Kafka 里消费 cdc.public.sync_check 这个 topic找到 marker 为 marker-20240501 的消息它的 source.lsn 如果大于或等于 63125392说明这条变更已经完整走通了 WAL 解析到 Kafka 的整条链路。用 ts_ms 和业务时间相减还能精确计算出端到端延迟。这个验证方式比手动检查目标库有没有多一行数据可靠得多因为它是逐条精确对位的。百分比对账建议做成日常巡检。源库和目标库每天跑一次抽样校验用 count 加字段哈希值对比SQL 格式大致如下SELECT count(*), sum(hashtext(email::text)) AS email_hash, sum(hashtext(status::text)) AS status_hash FROM public.users;源库和目标库分别执行三个值完全一致才算通过。别在意哈希碰撞这类理论风险这种校验的目的是快速发现大面积漏同步或错位不是密码学验证。一旦发现不一致回查上一次校验点位之后 Kafka 里的变更消息基本就是断点发生的位置。最后说一个我自己的习惯首次接入新表时先让它空跑一天第二天跑对账确认一致后再把告警阈值收紧。空跑一天看起来浪费时间但它能过滤掉数据结构不合理、类型映射错误、目标库约束缺失这些隐藏问题。这套系统的观感是搭起来不难跑稳定不容易真正稳住靠的是把 LSN 校验和对账做成日常机制而不是上线当天测通了就完事。希望这些经验对你的落地过程有帮助。本文还有配套的精品资源点击获取
返回列表