ARTICLE DETAIL

资讯详情

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

PostgreSQL逻辑复制与WAL解析:实时CDC同步到Kafka实战

PostgreSQL逻辑复制与WAL解析:实时CDC同步到Kafka实战 简介本资源面向数据库开发与数据集成工程师提供一套基于PostgreSQL逻辑复制功能的实时数据变更捕获与同步系统源码。系统通过解析WAL日志捕获数据变更将其转换为可执行的SQL语句并借助Kafka消息队列实现PostgreSQL到异构数据源的实时同步适用于大数据分析、实时报表与数据仓库等场景。压缩包共24个文件约474KB以17个Java源文件为核心实现辅以2个XML配置、properties参数文件、README与说明文档另附docx说明与png示意图便于理解整体架构与部署流程。目前已有57人学习下载。读者可获得完整的CDC同步实现代码、WAL解析与SQL转换思路、Kafka解耦设计及容错恢复机制参考适合需要搭建跨平台数据同步管道的中高级开发者对照学习与二次开发。1. 从 WAL 到 Kafka一条 PostgreSQL 变更数据管道的真实落地路径线上跑着一套 PostgreSQL业务侧突然要求把订单表的每次变更实时推到 Kafka下游还有 HBase、ES 和另一个 MySQL 异构库要同步。这时候你大概会先想到触发器写审计表或者定时全量扫增量字段——前者拖慢写入后者延迟高还容易漏。真正在生产里扛得住的方案是直接读 PostgreSQL 的逻辑复制槽把 WAL 日志解析成带前后镜像的变更事件再转成 SQL 或 JSON 投递到 Kafka。这套「基于 PostgreSQL 逻辑复制功能的实时数据变更捕获与同步系统」核心就三件事用pgoutput或wal2json插件把 WAL 解出来按表做过滤和格式转换最后通过生产者把消息可靠地送进 Kafka。它适合有 PostgreSQL 运维基础、需要跨异构数据源做秒级同步的团队也适合想搞懂 CDC 底层到底怎么跑的一线工程师。下面我按自己踩过的路把选型、配置、代码和坑一次讲清。2. 逻辑复制与 WAL 解析为什么不用触发器而选复制槽2.1 逻辑复制的底层机制与物理复制的区别PostgreSQL 的 WAL 是预写日志所有数据变更先写 WAL 再落盘这是崩溃恢复的基础。物理复制把 WAL 原样传到备库要求主备完全一致逻辑复制则通过pg_recvlogical或复制协议把 WAL 里的记录按逻辑解码规则翻译成行级变更。逻辑解码依赖wal_level logical这个参数决定了 WAL 里是否记录足够的信息来重建元组。设成logical后WAL 会额外记录旧元组标识取决于REPLICA IDENTITY解码插件才能输出INSERT、UPDATE、DELETE的完整内容。复制槽是逻辑复制的核心对象。它记录消费者已经确认的 LSN日志序列号PostgreSQL 会保留槽位之后的所有 WAL直到消费者确认。这意味着即使消费者宕机几小时重启后仍能从断点继续不会丢变更。代价是 WAL 会堆积磁盘可能被撑爆——这是最常见的翻车点后面避坑章节细说。和触发器方案比逻辑复制的优势很明显不侵入业务 SQL不增加写入事务的锁竞争能拿到变更前后的完整行镜像。触发器要在每张表上挂函数高并发下函数执行开销叠加还容易因为异常导致主事务回滚。逻辑复制是旁路读取对主库写入路径几乎无影响这是它成为 CDC 主流方案的根本原因。2.2 开启逻辑复制并创建复制槽的完整命令先确认postgresql.conf里的关键参数。改完必须重启wal_level不是动态参数# postgresql.conf 关键配置 wal_level logical # 必须否则无法创建逻辑槽 max_replication_slots 10 # 按消费者数量预留每个槽占一个 max_wal_senders 10 # 并发发送 WAL 的进程数 wal_keep_size 1GB # 防止 WAL 被过早回收按业务峰值调重启后用 SQL 创建逻辑复制槽。插件选pgoutputPostgreSQL 10 内置无需额外安装还是wal2json需装扩展输出 JSON 更省事取决于下游消费端。如果下游是 Kafka我一般用wal2json因为直接出 JSON省一层转换-- 创建逻辑复制槽插件用 wal2json SELECT * FROM pg_create_logical_replication_slot(cdc_slot, wal2json); -- 查看槽位状态和确认位点 SELECT slot_name, plugin, slot_type, active, restart_lsn, confirmed_flush_lsn FROM pg_replication_slots;pg_create_logical_replication_slot的第一个参数是槽名全局唯一第二个是解码插件名。创建后槽位处于非活跃状态直到有消费者连接。restart_lsn是 PostgreSQL 保证还保留 WAL 的起点confirmed_flush_lsn是消费者最后确认的位置。这两个值差距越大说明消费者落后越多WAL 堆积风险越高。如果表需要UPDATE和DELETE的旧值还得设置REPLICA IDENTITY。默认是DEFAULT只记录主键要拿完整旧行得改成FULLALTER TABLE orders REPLICA IDENTITY FULL;FULL会把整行旧值写进 WALWAL 体积会明显增大只对确实需要旧值的表开。大部分同步场景只需要主键定位DEFAULT就够。2.3 用 pg_recvlogical 验证变更能否被正确捕获在写消费者代码前先用命令行工具验证槽位能出数据。开一个会话执行pg_recvlogical -d your_db -S cdc_slot \ --slotcdc_slot --start -f - -o pretty-print1 -o include-xids1然后在另一个会话对业务表做增删改观察输出。-f -表示输出到标准输出-o传插件参数。wal2json的pretty-print让 JSON 带缩进include-xids带上事务 ID方便排查。正常输出类似{ change: [ { kind: insert, schema: public, table: orders, columnnames: [id, amount, status], columntypes: [integer, numeric, text], columnvalues: [1001, 299.00, paid] } ] }看到这个结构说明 WAL 解析链路通了。注意columnvalues的顺序和columnnames严格对应下游转换 SQL 时按这个顺序拼。如果输出为空先查wal_level是否生效、槽位是否 active、表是否在publication范围内用pgoutput时需要 publicationwal2json不需要。3. 把变更事件转成 SQL 并投递到 Kafka 的工程实现3.1 消费者程序的结构与关键依赖消费者要做四件事连上复制槽拉取变更、解析 JSON、按表路由生成目标 SQL、通过 Kafka 生产者发送。我用 Python 写依赖psycopg2的逻辑复制接口和kafka-python。核心结构是一个循环拉一批、解析、发送、确认 LSN。确认 LSN 是关键只有发送成功才能确认否则重启后会重复消费——所以下游必须能幂等。import json import psycopg2 from kafka import KafkaProducer # 连接 PostgreSQL注意 replicationdatabase 参数 conn psycopg2.connect( dbnameyour_db, userrepl_user, password***, host127.0.0.1, port5432, connection_factorypsycopg2.extras.LogicalReplicationConnection ) cur conn.cursor() cur.start_replication( slot_namecdc_slot, decodeTrue, options{pretty-print: 0, include-xids: 1} ) producer KafkaProducer( bootstrap_servers[kafka1:9092, kafka2:9092], value_serializerlambda v: json.dumps(v).encode(utf-8), acksall, # 所有 ISR 确认防丢 retries5, linger_ms20 # 攒批提升吞吐 )start_replication的decodeTrue让 psycopg2 自动把字节流转成字符串。options里的参数直接透传给wal2json插件。Kafka 生产者这边acksall保证消息写入所有同步副本才返回配合retries应对瞬时抖动。linger_ms是攒批等待时间设 20ms 能在延迟和吞吐间取平衡追求低延迟可以设 0。3.2 解析 wal2json 输出并生成目标 SQLwal2json的输出里change数组每个元素是一个行变更。kind字段区分insert、update、delete。update会同时带columnnames/columnvalues新值和oldkeys旧主键。生成 SQL 时按kind分支def build_sql(change): kind change[kind] table change[table] cols change[columnnames] vals change[columnvalues] if kind insert: placeholders ,.join([%s] * len(vals)) return (fINSERT INTO {table} ({,.join(cols)}) fVALUES ({placeholders}), vals) if kind update: # oldkeys 里是定位条件通常是主键 old change.get(oldkeys, {}) old_cols old.get(keynames, []) old_vals old.get(keyvalues, []) set_clause ,.join([f{c}%s for c in cols]) where_clause AND .join([f{c}%s for c in old_cols]) return (fUPDATE {table} SET {set_clause} fWHERE {where_clause}, vals old_vals) if kind delete: old change.get(oldkeys, {}) old_cols old.get(keynames, []) old_vals old.get(keyvalues, []) where_clause AND .join([f{c}%s for c in old_cols]) return (fDELETE FROM {table} WHERE {where_clause}, old_vals) return None, None这里用参数化占位符%s不是直接拼字符串。下游如果是 MySQL 或 HBase占位符风格不同但思路一致列名和值分开传避免 SQL 注入和类型转换错误。oldkeys在REPLICA IDENTITY DEFAULT时只含主键FULL时含整行旧值生成WHERE条件时按需取。3.3 发送到 Kafka 与 LSN 确认的时序控制拉取和发送必须串起来确认 LSN 只能在 Kafka 确认之后。psycopg2 的read_message返回一个消息对象带payload和data_start即 LSN。发送成功后调send_feedbackdef consume_stream(cur, producer): while True: msg cur.read_message() if msg is None: continue payload json.loads(msg.payload) for change in payload.get(change, []): sql, params build_sql(change) if sql is None: continue # 发送到 Kafkakey 用表名保证同表有序 future producer.send( pg-cdc-topic, keychange[table].encode(utf-8), value{sql: sql, params: params, lsn: msg.data_start, table: change[table]} ) future.get(timeout10) # 同步等确认失败抛异常 # 全部发送成功后才确认 LSN cur.send_feedback(flush_lsnmsg.data_start)future.get(timeout10)是同步等待确保这条消息真的进了 Kafka 才继续。如果超时抛异常整个循环中断LSN 不确认重启后从上次确认位点重放。这就是「至少一次」语义下游必须幂等。Kafka 的 key 用表名保证同一张表的变更进同一分区消费端能按顺序处理。如果追求吞吐可以改成批量发送后统一flush再确认但延迟会上升。4. 异构同步与 Kafka 投递的避坑清单4.1 复制槽不消费导致 WAL 撑爆磁盘现象主库磁盘告警pg_wal目录持续增长pg_replication_slots里槽位activefalse但restart_lsn很旧。原因消费者进程挂了或网络断了槽位没被确认PostgreSQL 按槽位保留 WAL 不回收。解决监控pg_replication_slots的restart_lsn和当前 LSN 差值超过阈值告警给消费者加自动重启实在不行用pg_drop_replication_slot删掉废弃槽位但会丢未消费变更删前确认下游能接受。4.2 REPLICA IDENTITY 没设导致 UPDATE/DELETE 缺旧值现象update事件里oldkeys为空或只有主键下游WHERE条件拼不出来更新错行。原因表默认REPLICA IDENTITY DEFAULT只记录主键如果表没主键连主键都没有。解决有主键的表保持DEFAULT即可需要旧值全量的表设FULL无主键表必须加主键或设FULL否则DELETE无法定位。4.3 Kafka 生产者 acks 配置不当导致消息丢失现象消费者确认了 LSN但 Kafka 里查不到对应消息。原因acks1或acks0时leader 写入即返回leader 宕机且副本未同步就丢。解决acksall配合min.insync.replicas2确保至少两个副本确认。同时retries设够enable_idempotenceTrue防重试导致重复。4.4 大事务导致单条消息过大被 Kafka 拒绝现象批量UPDATE十万行wal2json输出一个巨大 JSON超过 Kafkamax.request.size或message.max.bytes发送失败。原因逻辑解码按事务输出大事务就是大消息。解决调大 Kafka 的message.max.bytes和生产者max_request_size或者在消费者侧按行拆分一个事务拆成多条 Kafka 消息但要注意保持事务边界下游按事务 ID 聚合。4.5 时间戳和序列化格式在异构端不一致现象PostgreSQL 的timestamptz解出来是 ISO 字符串下游 MySQL 或 HBase 解析失败numeric精度丢失。原因wal2json默认把时间转成字符串数值按文本输出异构端类型映射没对齐。解决在build_sql里做类型转换时间统一转成 epoch 毫秒或目标库接受的格式numeric用字符串传输避免浮点精度问题。下游建表时字段类型要能容纳源库的最大精度。5. 用 LSN 对账和幂等写入把同步可靠性兜住同步系统跑起来容易跑稳难。我一般会加两层保险LSN 对账和幂等写入。LSN 对账是定期在主库查pg_current_wal_lsn()和消费者最后确认的confirmed_flush_lsn比差值就是延迟。这个值持续增大说明消费跟不上得扩容消费者或优化下游写入。幂等写入是在目标端建一张cdc_applied_lsn表记录已应用的 LSN每条变更应用前先查这个 LSN 是否处理过处理过就跳过。这样即使「至少一次」导致重放也不会产生重复数据。验证同步是否真的可靠我会做一次故障演练停掉消费者十分钟期间主库持续写入然后重启消费者观察是否能从断点续上、数据是否完整。再模拟 Kafka 不可用看消费者是否阻塞在future.get而不确认 LSN恢复后是否自动继续。这两步能暴露大部分时序和确认逻辑的问题。一个具体技巧是给wal2json加include-transaction和include-timestamp参数输出里带上事务提交时间和事务 ID。下游按事务 ID 做去重按提交时间做乱序判断比单纯靠 LSN 更稳。参数在start_replication的options里传cur.start_replication( slot_namecdc_slot, decodeTrue, options{ include-xids: 1, include-timestamp: 1, include-transaction: 1, filter-tables: public.audit_log # 排除不需要同步的表 } )filter-tables能排除审计表、日志表这类不需要同步的表减少无效流量。多个表用逗号分隔。这个参数在表多的时候特别有用不然所有变更都往 Kafka 灌下游还得过滤。我自己踩得最狠的一次是没监控复制槽消费者进程被 OOM kill 后没人管第二天主库磁盘满了导致写入全部阻塞业务直接不可用。从那以后复制槽的restart_lsn延迟和磁盘水位成了必配告警。这套方案值得做但前提是把监控和幂等做扎实不然实时同步会变成实时故障。希望帮到你。本文还有配套的精品资源点击获取
返回列表