ARTICLE DETAIL

资讯详情

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

基于Kafka与Flink的异构数据库同步方案:从延迟监控到账实相符

基于Kafka与Flink的异构数据库同步方案:从延迟监控到账实相符 做数据库迁移最怕的不是停机窗口不够而是迁移完第二天业务方拿着对账单找上门“昨天下午那笔订单怎么没了”很多人习惯用同步延迟来评估迁移质量看着监控面板上几百毫秒的 lag觉得一切稳了结果一核对明细总有几条数据对不上。异构数据同步这事的真相是延迟只是面子账实相符才是里子。我参与过的几个不停机迁移项目最后能顺利收尾靠的都是一套叫 KFSKafka-Flink Sync的同步链路方案把每一笔数据都当成一笔账来管。这篇文章就把这套方案的思路、细节和踩过的坑完整拆给你。不管你是刚接触数据迁移的 DBA还是已经在做异构同步的工程师这篇文章都能给你一个可以直接落地的参考框架。下面我按从问题到解法、从设计到实操的顺序来聊。1. 不停机迁移的三大死穴为什么延迟正常还是会丢账1.1 异构系统之间格式能转语义难对齐很多团队在评估异构数据同步时第一反应是“字段类型转换一下不就行了”。但真正做过的人都知道格式转换只是最表层的问题。我碰到过一个真实案例源端 MySQL 里的订单状态是枚举值0/1/2目标库 PG 里定义成pending/paid/refunded如果只做字典映射看起来没问题但源端中间有一次业务升级把支付中的状态从1改成了3同步链路还按旧映射跑结果目标库里出现了大量空状态和unknown。这只是语义差异的冰山一角。更隐蔽的是时区问题、金额精度问题、空字符串和 NULL 的问题。比如 MySQL 里0000-00-00 00:00:00这种日期值同步到 PG 的 timestamp 字段会直接报错再比如金额字段源端是decimal(10,2)目标端为了性能用了double平时看着没事遇到特别大的订单金额尾数就可能差出分来。异构同步最难的点不是“怎么把数据搬过去”而是“怎么保证搬过去的每一条数据在目标系统里依然具备同样的业务含义”。1.2 追延迟不等于追账增量链路里的三类丢数据监控面板上最直观的指标就是延迟但延迟低只代表消息从源端到目标端的流动速度快不代表消息没有在流动过程中丢。我在生产环境里至少见过三类“看着正常、实则丢账”的情况。第一类是主键冲突导致的部分更新失败。源端订单表主键是自增 ID目标端也定义了同样的主键但同步程序在写入时没有做幂等处理。某一次同步任务重跑了同一批消息后到的消息因为主键冲突被目标库拒绝程序又没有捕获异常继续往下走结果这条数据就默默丢了监控上延迟依然很低但目标库里的记录数和源端差了一条。第二类是乱序覆盖。比如同一笔订单先发生了“创建”又立刻发生“状态修改”两条消息都进入了 Kafka。如果目标端写入时没有版本比较后消费到创建消息把它当成最新状态写进去就会把已经更新的状态给覆盖掉。这类问题在延迟监控上完全看不出来只有明细核对才能发现。第三类是事务边界被拆散。源端一个事务里同时更新了订单表和订单明细表CDC 工具把两张表的变更分别发到两个 topicFlink 任务分两条链路处理。目标库先收到明细后收到订单表的主记录如果刚好在最后一条明细写入前发生了任务重启就可能导致明细表和主表之间出现不一致。每一类丢数据都是在提醒我们延迟指标只能证明“链路通没通”不能证明“账全没全”。1.3 双写窗口停机窗口不可能为零不停机迁移的名字听起来很潇洒但实际操作上并不是完全不停。几乎所有方案都需要一个短暂的“业务切换窗口”可能是五分钟也可能是半个小时。在这个窗口里新老系统同时开放写入等确认一致后把流量切到新系统。问题在于如果前置的增量同步没有把数据追上切换窗口就只能往后拖。正常情况下五分钟的切换窗口意味着你要提前把延迟压到秒级同时还要在窗口内完成最后一轮数据校验。我见过一个团队为了赶上线时间把同步并发调得特别高延迟是压下去了但目标库因为写入压力过大导致主键冲突频繁数据反而对不上。最后只能回滚多花了整整两个晚上去补对账。这个教训很典型追延迟本身没有错但不能为了追延迟牺牲数据正确性。真正可靠的做法是先设计好“对账机制”再去优化同步速度。KFS 这套方案的核心思路就是把每一笔数据都当成账目来处理让同步链路具备可校验、可回放、可追补的能力。2. 理解 KFS 的核心设计Kafka-Flink Sync 的账本思路2.1 为什么选择 Kafka Flink 这个组合KFS 是我在项目里对这套链路的内部叫法全称是 Kafka-Flink Sync。选用这个组合不是因为它们“流行”而是因为它们在不停机迁移这个场景里提供了三个关键能力缓冲、回溯、状态。Kafka 天然是个大缓冲。源库的写入峰值往往不可控如果同步工具直接连源库拉数据一旦源库抖动同步进程很容易被打垮。把变更日志先写入 Kafka相当于给源库和目标库之间加了一道缓冲无论是目标库短时间不可用还是 Flink 任务需要重启消息都还在 Kafka 里不会丢。更重要的是 Kafka 的消息可以按 offset 回溯这为后面做对账和补数提供了基础设施。Flink 的价值在于流处理和状态管理。源端 CDC 出来的消息是分散的更新事件而目标库需要的往往是最新状态。Flink 可以把同一主键的多条变更事件按时间或版本号做合并实现幂等写入。同时 Flink 的 checkpoint 机制可以保证“处理到哪一条”这件事有据可查配合 Kafka 的 offset就能实现精确一次语义。整个方案的定位不是简单的搬运工具而是一套可追踪、可审计的数据流转管道。2.2 关键字段驱动用业务主键把链路串起来KFS 设计里最重要的一条原则是“每条消息都必须携带业务主键”。这里的业务主键不一定等于数据库主键而是目标端用来唯一识别一条记录的那个字段组合。比如订单表可以用order_id账户余额表可以用account_id关联明细表可能需要用(order_id, line_no)的复合键。所有环节都围绕这个关键字段展开。Kafka 消息的 key 就用业务主键这样可以保证同一笔账的变更事件进入同一个 partitionFlink 在处理时根据 key 分组天然保持顺序。写入目标库时upsert 语句也基于这个 key 做插入或更新。对账任务更是直接按主键集合做差集比对。如果主键设计得不对比如用了源端自增 ID而目标端已经重新定义了 ID 生成策略那整个链路的幂等和对账都会失去抓手。实际操作中还要注意业务主键一旦确定就不要轻易更换。我见过一个项目开始用订单 ID 做 key后来业务要求按用户维度聚合就把 key 改成了用户 ID结果同一个订单的多个变更事件被分散到不同 partition乱序问题立刻爆发。关键字段在一个同步周期内必须保持稳定。2.3 不止追延迟水位线和位点给每笔账做标记KFS 的第二个设计要点是把同步进度的粒度从“链路延迟”细化到“记录位点”。Kafka 的 offset 记录了消费者组读取到哪个位置Flink 的 checkpoint 记录了状态快照这两样东西结合起来就能形成一个“同步账本”。简单说每条消息进入 Kafka 时都会带着源库的 binlog 位点或时间戳Flink 处理完一批消息后会把这批消息中最大的位点记录到目标库的同步监控表里。这个账本有什么用遇到数据不一致时我们可以直接看目标库同步监控表里记录的位置减去 Kafka 当前最新 offset得到精确的“未完成位点范围”。而不是像以前那样只能看到“延迟了 5 分钟”这种模糊信号。更重要的是对账任务可以基于位点做定向补数比如发现某个主键的数据在当前 binlog 位点之后没有更新就可以单独从 Kafka 重放该主键的消息不需要把整张表重新同步一遍。说白了延迟监控回答的是“链路快不快”位点账本回答的是“每笔账到没到”。KFS 的设计目标就是让后者成为同步链路的默认能力。3. 核心细节解析与实操要点3.1 CDC 捕获阶段从 binlog 到 Kafka 的三条注意事项不管是 MySQL 还是 PostgreSQLCDC 工具的选择很多Debezium、Canal、Flink CDC 都用过。但工具只是第一步真正决定同步质量的是配置细节。这里分享三条最容易踩坑的注意事项。第一源库 binlog 格式必须设成ROW并且binlog_row_image要设成FULL。默认的MINIMAL模式只记录修改列导致 before 镜像里缺字段遇到需要根据旧值做判断的场景就会出问题。我遇到过因为没改这个配置导致更新事件里没有主键信息目标端无法定位记录只能转全量重刷。第二每个同步任务必须分配独立的server-id。CDC 工具本质上是一个模拟从库的客户端如果多个任务共用同一个 server-id源库会误判为同一个从库连接可能导致连接被踢掉binlog 发生断裂。建议把 server-id 的分配纳入部署脚本管理不要用手动配置。第三schema 变更必须提前演练。CDC 工具会把 DDL 事件也发到 Kafka如果目标端还没准备好对应的字段变更Flink 任务解析消息时就会报错。我们的做法是在生产环境变更表结构之前先在测试环境跑一遍完整链路确认目标库能正确接收新字段。这件事看起来和同步无关但在不停机迁移期间业务方很可能因为上线新功能而修改表结构这块不提前管理迁移窗口里很容易出大乱子。3.2 消息结构设计payload 里必须带版本号和来源标记Kafka 消息如果只存业务字段后面做排序、去重和对账都会很吃力。KFS 的实践里每条消息的 value 结构至少要包含下面几类信息。先看一段简化后的 JSON 示例{ schema: dbserver1.inventory.orders, payload: { op: u, before: { order_id: 1001, status: pending }, after: { order_id: 1001, status: paid, amount: 199.00 }, source: { version: 1.9.4.Final, connector: mysql, ts_ms: 1720000000000, snapshot: false, db: shop, table: orders, server_id: 223344, file: mysql-bin.000188, pos: 4231 }, op_ts: 1720000000000, event_ts: 1720000001000 } }这里面的op标记操作类型source.file和source.pos是 binlog 位点这两个字段必须保留因为它们是后面对账和定位问题的基础。除了 CDC 工具自带字段我强烈建议在消息进入 Flink 前增加一个自定义的event_ts字段记录业务实际发生时间而不是 CDC 捕获时间。这是因为如果源库有批量补数或人工修改捕获时间会和业务时间存在较大偏差后续做事件时间窗口计算时会出问题。另外如果有多套环境或业务来源payload 里还需要加一个env或source_system标记。我有一次同时接了生产和预发两套库因为没有来源标记Flink 任务把预发环境的测试数据也写进了目标库查了好久才发现问题。消息结构看似只是格式问题实际是同步链路可靠性的基石。3.3 写入端幂等设计异构目标库 Upsert 的正确姿势Kafka 消息是至少一次语义Flink 任务重启后可能重复消费同一批消息所以目标端写入必须做好幂等。不同数据库提供的能力不一样但核心思路都一样基于关键字段做 upsert并且用版本号判断是否覆盖。以 PostgreSQL 为例目标表建一个版本字段sync_version每次写入时用两段 SQLINSERT INTO target_orders (order_id, status, amount, sync_version) VALUES (?, ?, ?, ?) ON CONFLICT (order_id) DO UPDATE SET status EXCLUDED.status, amount EXCLUDED.amount, sync_version EXCLUDED.sync_version WHERE EXCLUDED.sync_version target_orders.sync_version;这里的sync_version可以直接用 binlog 位点拼接出来的数字比如file_index * 100000000 pos保证同一个事务里所有变更事件都具备单调递增的顺序。如果目标库是 MySQL可以用INSERT ... ON DUPLICATE KEY UPDATE加条件更新但要注意 MySQL 不支持WHERE子句作用于VALUES后面的更新条件处理起来比较绕。更稳妥的方式是先把消息写入 Kafka 前的 CDC 版本号变成目标表的版本字段用数据库的CASE WHEN判断INSERT INTO target_orders (order_id, status, amount, sync_version) VALUES (?, ?, ?, ?) ON DUPLICATE KEY UPDATE status IF(VALUES(sync_version) sync_version, VALUES(status), status), amount IF(VALUES(sync_version) sync_version, VALUES(amount), amount), sync_version IF(VALUES(sync_version) sync_version, VALUES(sync_version), sync_version);ClickHouse 这种分析型数据库没有真正意义上的 upsert但可以用ReplacingMergeTree引擎以版本列为排序键后台自动保留最新版本。无论用哪种数据库都要在测试环境里模拟“重复发送同一条消息”的场景验证幂等是否生效。我在项目里专门写了一个测试脚本往目标库插一条记录后再把相同主键的旧版本消息重新投递一次如果目标记录没有被旧版本覆盖说明幂等逻辑正确。4. 实操过程与核心环节实现4.1 最小可用部署拓扑三台机器也能跑起来很多人一听 Kafka Flink 就以为要搭一个大数据集群实际不停机迁移项目里三台 8C16G 的虚拟机就能跑起来。我们项目里采用的是最小部署拓扑大致如下源端 MySQL开启 binlog配置独立服务器账号。Kafka ConnectDebezium以插件方式运行负责读取 binlog 并写入 Kafka。Kafka 集群3 个 broker负责承接变更消息。Flink 集群1 个 TaskManager可先跑单机模式消费 Kafka 消息并写入目标库。目标端 PostgreSQL提供服务给业务系统。对账服务独立脚本负责周期性地从源库和目标库计算数据指纹并对比。这套拓扑不做高可用也能满足大部分迁移场景。迁移窗口内如果发生节点故障Kafka 的复制机制和 Flink 的 checkpoint 能兜住大部分问题如果确实需要多活再考虑给 Kafka 和 Flink 加节点。我的建议是不要一开始就把架构搞复杂先把同步链路的账目逻辑跑通再根据压力测试结果决定是否扩展。4.2 Flink SQL 同步任务示例MySQL 到 PG下面给一个可以直接参考的 Flink SQL 任务示例场景是 MySQL 的orders表同步到 PostgreSQL 的orders表。真实项目里通常会有几十张表这里演示核心写法。先创建 Kafka 源表注意format使用debezium-jsonFlink 能自动解析 Debezium 的 CDC 消息结构CREATE TABLE source_orders ( order_id BIGINT, status STRING, amount DECIMAL(10, 2), event_time BIGINT, version BIGINT, PRIMARY KEY (order_id) NOT ENFORCED ) WITH ( connector kafka, topic dbserver1.inventory.orders, properties.bootstrap.servers kafka-1:9092,kafka-2:9092,kafka-3:9092, properties.group.id kfs-order-sync, scan.startup.mode earliest-offset, format debezium-json );接着创建目标表使用 JDBC 连接器CREATE TABLE target_orders ( order_id BIGINT PRIMARY KEY NOT ENFORCED, status STRING, amount DECIMAL(10, 2), sync_time TIMESTAMP(3), version BIGINT ) WITH ( connector jdbc, url jdbc:postgresql://pg-target:5432/shop, table-name orders, username sync_user, password sync_password );然后写入时先把 Oracle 的二进制位点转成version再使用INSERT INTO ... SELECT将数据发送到目标表。但要注意标准 JDBC 连接器只支持 append 写入不支持 upsert。如果要用 Flink 完成幂等写入需要借助自定义 Sink 或者使用 Flink CDC 生态里的JdbcUpsertSink。最常用的做法是在 DataStream 里写 Flink SQL 的TableEnvironment配合upsert-kafka等连接器如果不想引入额外组件也可以直接在 DataStream 里写一个RichSinkFunction把upsert逻辑封装进去。下面是一个简化的RichSinkFunction示例用来演示版本控制的核心思路public class OrderUpsertSink extends RichSinkFunctionOrderRecord { private Connection conn; private PreparedStatement stmt; Override public void open(Configuration parameters) throws Exception { Class.forName(org.postgresql.Driver); conn DriverManager.getConnection( jdbc:postgresql://pg-target:5432/shop, sync_user, sync_password); String sql INSERT INTO orders(order_id, status, amount, sync_version) VALUES (?, ?, ?, ?) ON CONFLICT (order_id) DO UPDATE SET status EXCLUDED.status, amount EXCLUDED.amount, sync_version EXCLUDED.sync_version WHERE EXCLUDED.sync_version orders.sync_version; stmt conn.prepareStatement(sql); } Override public void invoke(OrderRecord record, Context context) throws Exception { stmt.setLong(1, record.getOrderId()); stmt.setString(2, record.getStatus()); stmt.setBigDecimal(3, record.getAmount()); stmt.setLong(4, record.getVersion()); stmt.executeUpdate(); } }这个类里最关键的是最后那行 SQL只有当新消息的sync_version大于目标库当前版本时才覆盖否则放弃。实际项目中下游收到业务系统在任何时刻发来的新数据只要 version 是递增的就可以安全地“后到优先”这也是 KFS 能处理乱序问题的根本原因。4.3 对账任务怎么做把“每一笔账”变成可校验的集合对账是整个方案里最容易被忽略、却最重要的一环。如果把同步链路比作一笔流水账那对账任务就是定期清点账目。我们需要保证源端和目标端的每条业务记录都最终一致而不仅仅是延迟接近零。KFS 的对账思路分为三个层次。第一层是计数对账。每隔固定周期分别查询源库和目标库某个表的COUNT(*)如果数量不一致立刻告警。这个方法最快但只能发现整表层面的差异定位不到具体的主键。第二层是主键差集对账。把源表主键集合和目标表主键集合分别查询出来在内存中做MINUS操作找出只在源端出现的主键列表。这个列表就是我们常说的“丢账明细”。可以用临时表或中间表加速避免在业务高峰期做全表扫描。拿 MySQL 举例-- 源库侧抓取主键快照 SELECT order_id FROM source_db.orders WHERE updated_at 2025-01-01 00:00:00;再把目标库侧的主键集合同样导出放到同一张对账临时表里做差集。脚本可以每分钟跑一次也可以按迁移窗口需要在切换前高强度跑。第三层是数据指纹对账。为了发现“数量对但没有更新到最新状态”的问题需要对关键字段计算哈希指纹。最简单的做法是把一行记录按固定分隔符拼接成字符串计算 MD5例如SELECT order_id, MD5(CONCAT_WS(|, COALESCE(status, ), COALESCE(CAST(amount AS CHAR), ), COALESCE(updated_at, ))) AS row_hash FROM source_db.orders;目标库跑同样的指纹逻辑然后两边的order_idrow_hash做比对结果不一致的就是需要从 Kafka 重放的主键。这个方法成本较高建议只对核心交易表开启并且只在迁移切换前、以及迁移后的一两天内运行。日常运行计数对账和主键差集对账已经足够。5. 常见问题与排查技巧实录5.1 问题速查表延迟、积压、数据漂移我在多个异构数据同步项目里沉淀了一张问题速查表遇到问题先按表排查能省大量时间。症状可能原因排查方法解决方案延迟持续走高但 Kafka 无积压Flink 任务消费能力不足或目标库写入慢查看 TaskManager 的 CPU、内存和数据库连接池增加并发或把批量提交大小调大减少频繁 commitKafka 积压明显延迟正常消息生产速度远大于消费速度查看 consumer lag 和 topic 分区数增加 topic 分区同时提高 Flink 并行度目标库记录比源库少主键冲突被忽略或事务边界处理问题用主键差集对账定位具体主键升级幂等逻辑开启错误日志补数脚本回放目标库记录数比源库多重复消费且未做幂等查目标库重复主键修正写入逻辑为 upsert并清理重复记录同一主键状态被旧值覆盖乱序低版本消息后到查看消息的 version 字段在目标端 SQL 中加入版本比较低版本不覆盖新字段同步不过去schema 变更未同步到目标库查看 CDC 任务日志检查 DDL 事件提前在测试链路演练 DDL确认字段映射时间字段相差 8 小时时区不一致对比源库时区和 Flink 进程时区统一使用 UTC 存储展示层再转业务时区这张表覆盖了迁移项目里 80% 的“同步异常但延迟正常”场景。排查时记住一个原则先确认 Kafka 是否完整再确认 Flink 是否处理正确最后确认目标库是否写入成功每一步都要有可观测的日志和位点记录。5.2 追账实操手工比对和自动补数的经验即便做好了前面的每一步迁移过程中还是难免需要人工追账。这里分享几个实操经验。第一补数任务必须和目标库写入任务隔离。假设你发现某个主键的账漏了如果直接在同一个 Flink 任务里重新投递消息很可能因为和目标库正在执行的更新冲突导致新账变旧账。我们通常在 Kafka 里新建一个retry-topic补数消息发送到这个 topic由单独的补数 Flink 任务消费并且在写入目标库前统一加上当前时间戳。这样才能保证补数操作不会打乱正常链路的位点。第二追账之前一定要先做全量快照对比。漏一条账可能是漏了一批账的冰山一角。我遇到过一次表面上只有 3 条订单对不上对账脚本只重放了这 3 条。结果第二天业务方又反馈少了 20 条原因是这些订单都来自同一次批量导入源库当时没有正确生成 binlog 事件。后来调整策略对账发现差异后先查这一批记录的source.file和source.pos确认是否集中在同一个 binlog 区域再决定是单条补数还是批量重刷。第三给所有补数操作留下审计日志。谁在什么时间补了哪些主键、补数前后的值是什么都要落到一张表里。迁移结束后如果业务方质疑数据这张日志表就是你最有力的解释。我在项目里专门建了一张sync_repair_log字段包括repair_time、table_name、pk_value、old_version、new_version、operator上线后几乎每天都会用到。5.3 让迁移后期更省心的两件小事最后分享两件不复杂但很有用的小事。第一在目标库建一张监控表记录每个同步任务的“当前已处理位点”和“当前 Kafka 最新位点”。这个表不仅给你一个精确的延迟指标还能在排查问题时直接回答“这个任务到底处理到哪了”。很多团队只盯 Kafka 自带 consumer lag但那个指标反映的是消费者组的读取位置不反映 Flink 是否已经完整写入了目标库。自己维护位点表才能做到每一笔账有据可查。第二迁移完成后的前三天不要立刻删掉 Kafka 里的原始 topic。Kafka 的过期时间通常设置成 7 天如果迁移刚结束就发现漏了很多历史数据kafka 里的原始变更日志还在就能直接重放。如果有条件最好在切换窗口结束后先保留 topic 数据一周等业务方确认无差异后再做清理。这个操作不花什么成本但能让你在“刚上线最慌的那几天”多一重保障。我个人在这些项目里的体会是不停机迁移最难的不是技术方案选型而是把“数据同步”从搬运任务升级成账目管理任务。追延迟靠监控工具就够了守住每一笔账得靠机制和流程。KFS 这套思路最重要的价值就是把可校验、可追补的能力提前设计进同步链路而不是等出问题后再去救火。如果你也在准备类似迁移建议先搭一条最小链路把对账脚本跑通再谈并发和性能优化。
返回列表