
异构数据同步做了这么多年我一直有个观点延迟数字是给老板看的账实相符才是给自己保命的。尤其是不停机迁移这种场景业务一分一秒都不能断数据却要从旧库搬到新库两边结构还不一样这时候如果只盯着同步延迟等到切换那天发现订单少了、金额多了事故等级直接拉满。这篇文章要说的 KFS是一套用 Kafka、Flink 配合目标存储StarRocks、MySQL 都可以搭起来的异构数据同步链路。核心不是把延迟压到多低而是解决不停机迁移过程中最要命的问题数据不丢、不重、每一笔都能对上。我会把整个链路的选型思路、关键配置、一致性保障手段以及我在真实项目里踩过的坑一次性讲清楚。适合正在做 MySQL 迁移到 StarRocks/ClickHouse/新实例、实时数仓同步或者想把应用双写改成日志级同步的团队参考。1. 迁移场景里为什么说延迟不是唯一指标1.1 不停机迁移的本质搬一个还在营业的仓库先打个比方。不停机迁移就像你开着一家仓库白天客户进进出出、货物不停流转这时候你还要把仓库从 A 栋搬到 B 栋而且搬家期间一个客户都不能赶走。很多人第一反应是我多叫几辆卡车拉货跑快点就好了。但真正的难点根本不是卡车跑多快而是搬过去的货不能少、不能多、不能串。同步数据也是一样。老订单在 MySQL 里新系统在 StarRocks 或者另一个 MySQL 实例里业务还在持续写入。延迟从 10 分钟压到 10 秒看起来很好看但这 10 秒里系统故障了丢了 500 条 binlog 没消费你根本察觉不到直到第二天对账发现少了 50 万营业额那时候再去补数据就非常被动了。所以我在评估同步方案时第一优先级永远是数据能不能追回来。延迟只是体验指标数据完整性和准确性才是生死指标。KFS 这套架构从一开始就把这几件事分开处理Kafka 保证消息不丢Flink 保证计算状态可恢复目标端的幂等模型保证重复数据能被吸收。1.2 延迟追平之后真正的工作才刚刚开始我见过太多团队一看到监控面板上同步延迟归零就欢呼追平了然后准备切换。实际上延迟归零只代表当前这一刻两边进度一样不代表历史数据完全一致。存量数据和增量数据之间的接缝有没有处理好同步过程中产生的重复数据有没有被正确去重源库 binlog 里的历史变更有没有在 Kafka 里积压到被清理这些问题在延迟数字上是看不到的。一个合格的不停机迁移同步链路要同时满足四个条件低延迟业务变更能在秒级内到达目标端这是体验底线。Exactly-Once 语义消息从采集到写入目标端不能丢也不能重或者至少下游具备幂等吸收重复的能力。可对账任意时刻都能拉出源端和目标端各有多少数据、差异在哪。可回滚切换出问题后能快速回到源库继续服务而不是被困在半路。KFS 这套组合里的每一项都在为这四个条件服务。Kafka 的持久化和 offset 机制解决了不丢Flink 的 checkpoint 和幂等 Sink 解决了不重而目标端的唯一键模型和对账任务解决了账实相符。1.3 为什么选 KFS 而不是应用双写或定时批处理先列一下市面上常见的几种异构同步方案方案实现成本对源库影响延迟水平不停机迁移适配度应用双写高改代码多无额外影响但有业务耦合实时差回滚要发版Trigger/存储过程中源库性能下降明显准实时较差侵入性强定时批处理低高峰期有压力分钟到小时级差追不上增量日志采集Canal/Debezium中高略微影响类似多一个从库秒级好天然异步解耦日志采集KafkaFlinkKFS高略有影响但有缓冲削峰秒级最好断点续传、可重放应用双写是我最不建议的方案。你在每个业务方法里多写一行调用当时看起来简单等迁移完成想把这个双写代码删掉需要所有业务方重新发版。而 binlog 采集方式对业务代码完全无感切换目标端只是改一下消费端配置回滚也只是把消费暂停、把流量切回旧库业务代码一行不用动。至于 KFS 里的Flink是不是必须的如果你只是小表同步Kafka 直接配一个 Connector 写到目标端也可以。但一旦涉及复杂转换、多表关联、维表补全、窗口聚合或者目标端偶尔抖动需要重试Flink 的状态管理和 Checkpoint 机制能让你省下大量手工写代码的工作。我在后面的章节会具体拆开讲。2. 拆解 KFS 同步链路每一层都在解决什么问题2.1 采集层Binlog 位点整条链路的账本起点从 MySQL 同步数据出来目前主流做法是用 Canal 或者 Debezium 伪装成 MySQL 的从库拉取 binlog。这一步最核心的不是怎么连而是位点怎么管。位点就是你读到 binlog 的哪个位置类似看书看到第几页第几行。Canal 每次拉取都会维护一个位点这个位点必须持久化到 ZooKeeper、MySQL 或者本地文件里。跑批任务的朋友可能觉得断点续传很简单但在 binlog 场景里有个麻烦binlog 文件是有生命周期如果 Canal 挂了 3 天位点对应的 binlog 文件已经被源库清理掉了续传就无从谈起。实操上我会做两件事源库 binlog 保留时间调长。默认可能是 1 天迁移期间至少调到 7 天大促前或长假期前还要再评估。用 SQLSHOW VARIABLES LIKE binlog_expire_logs_seconds查看低于 604800 的一律改掉。Canal 的位点存储独立出来。不要落在 Canal 服务本身所在的机器磁盘上否则机器一换位点就丢了。Debezium 还提供了一个很实用的能力把位点作为消息的一部分写进 Kafka Record Header。这样下游消费的时候不光知道这是什么数据还知道这数据在源库 binlog 的哪个位置产生。对后续做对账和问题回溯非常关键。2.2 缓冲层Kafka 的分区、顺序性与消息格式Kafka 在链路里不是简单的中转站它承担了削峰、缓冲和重放三个职责。源库高峰期每秒几千条变更目标端可能扛不住这么大的写入压力Kafka 先把消息积压下来Flink 按自己的节奏消费这就是削峰。下游处理逻辑出了问题只要 Kafka 里的消息还没过期就可以重置 offset 重新消费这就是重放。但 Kafka 用不好会引入一个很隐蔽的问题顺序乱。同一行数据比如同一个订单号如果被改了三回它的三次变更应该按顺序写到目标端否则最终值可能就是错的。要保证顺序必须严格控制分区规则Topic 按表拆分至少也要按库拆分不要把几十张表的大杂烩混在一个 Topic 里。每条消息的 Key 用业务主键或者唯一键Flink 消费后按 Key 做keyBy这样同一行数据的变更会进入同一个分区、同一个 Flink 子任务天然有序。消息体我建议至少带上这些字段before、after、opinsert/update/delete、ts、binlogFile、binlogPos。别觉得字段多浪费空间等你要排查数据问题时没有这些信息你会非常痛苦。注意Kafka 的 Topic 保留时间retention.ms在迁移期间也要调大建议至少 72 小时。因为一旦 Flink 任务出了故障重放窗口不够Kafka 里数据被清理了你只能从源库 binlog 重新补步骤繁琐且容易出错。2.3 计算与写入层Flink 的 Checkpoint 与幂等 SinkFlink 在这套链路里做三件事消费 Kafka 消息、做转换清洗、写入目标端。其中最容易出问题的不是转换逻辑而是状态一致性的保障。Flink 的 Checkpoint 机制简单理解就是定时把你的任务状态拍个快照存到外部存储HDFS、S3 或本地。假设任务处理到 Kafka offset 1000 时快照成功了突然宕机重启Flink 会从 1000 这个位置重新消费数据。注意是重新消费也就意味着 1000 之前已经写进目标端的数据可能会再写一遍。这就是 Kafka 的至少一次语义带来的重复。要解决重复光靠 Flink 不够必须让目标端具备幂等写入能力。对于 StarRocks用 Unique Key 模型写入时按主键去重对于 MySQL/ClickHouse对应的分别是INSERT ... ON DUPLICATE KEY UPDATE和 ReplacingMergeTree。宗旨是同一条数据写多少次最终结果都一样。我见过有人试图用 Flink 的精确一次语义强行解决重复把两阶段提交协议开到目标端。这个方案在部分存储上能跑通但代价是吞吐下降 30%~50%而且目标端必须支持事务。多数情况下没必要用好下游幂等就够了实际效果一样成本低得多。3. 守住每一笔账一致性保障的实操细节3.1 三个游标如何配合Binlog 位点、Kafka Offset、Flink Checkpoint一条数据从 MySQL 流到目标端会经过三个游标游标作用存储位置故障恢复方式Binlog 位点记录 Canal 读到源库的哪个位置ZooKeeper/MySQL/本地Canal 重启后从该位点继续拉取Kafka Offset记录 Flink 消费到 Topic 的哪个位置Kafka 内置 Flink CheckpointFlink 从 Checkpoint 恢复并重置 OffsetFlink Checkpoint记录任务处理状态和已提交位置HDFS/S3/本地任务重启后自动恢复这三个游标没有谁优谁劣它们在不同层级上解决不同问题。理解它们的关系你就明白为什么说只要链路设计对了数据一定能找回来MySQL 到 KafkaCanal 管位点确保 binlog 变更全部进了 Topic。Kafka 到目标端Flink 管 Checkpoint确保即使消费一半挂了也能从最近快照续跑。如果链路整体落后太多你还可以手动把 Kafka Offset 往回调让目标端重放一批数据配合幂等机制消除重复。做迁移项目时我会单独建一张心跳表每隔几秒在源库更新一条记录看看这条心跳记录经过链路到达目标端用了多久。这样比单纯看 Kafka Lag 更直观Lag 只能说明 Kafka 里积压了多少不能说明端到端延迟。3.2 幂等写入目标端必须有唯一键否则一切都是空谈我再强调一遍没有唯一键任何一致性方案都是空谈。Flink 重启后重复消费 Kafka 是常态不是意外。如果目标表没有唯一键约束同样的插入执行两次数据就翻倍了。我自己踩过一个很深刻的坑。当时迁移一张日志表建表时没有建唯一键觉得日志表不会更新插入就行。结果 Flink 任务因为源库抖动重启了几次日志表里出现了大量重复数据。因为是日志表对业务影响不大但这件事让我养成了一个习惯所有通过 KFS 同步的表目标端必须设计唯一键哪怕业务上不需要更新也要拿源表主键做唯一键。不亏这是买保险。StarRocks 的实现方式是建 Unique Key 模型MySQL 用唯一索引加ON DUPLICATE KEY UPDATEClickHouse 用 ReplacingMergeTree。Flink SQL 里的一句话就能体现设计思路INSERT INTO target_table SELECT id, order_no, amount, status, update_time FROM source_view目标表如果是 StarRocks Unique Key 模型相同主键的多条记录只会保留最新版本多次写入天然去重。3.3 对账机制迁移验收的硬指标对账不能等同于跑一下 count(*) 看数量一样。数量一样金额可能差金额一样明细可能多一条少一条。我建议至少分三层对第一层总量级对账。每天统计源表和目标表的行数、关键字段求和比如订单金额总和。这层能快速暴露大问题比如表没同步、任务没起来。第二层分片对账。按主键范围或者按天把数据切成多个分片每个分片分别统计行数和校验和比如对主键做 XOR或对某几个字符型字段做 SUM。分片对账能定位到具体哪一段数据有问题。第三层抽样明细对账。每个分片随机抽 100~1000 条比对每个字段的值是否完全一致。抽样明细能发现类型转换、时区处理、浮点精度这类数量对但内容不对的问题。一个可落地的简易 SQL 对账思路是这样-- 源端分片统计 SELECT CONCAT(DAY(create_time)) AS day_partition, COUNT(*) AS row_cnt, SUM(amount) AS amount_sum FROM source_order WHERE create_time 2024-06-01 GROUP BY day_partition; -- 目标端分片统计假设已有对应字段 SELECT day_partition, COUNT(*) AS row_cnt, SUM(amount) AS amount_sum FROM target_order WHERE create_time 2024-06-01 GROUP BY day_partition;把两个结果放在同一个报表里用 JOIN 比较差值。差值超过阈值就触发告警由值班人员来看。迁移期间的第一个星期对账任务建议每半小时跑一次后面逐步降低频次。对账本身也会给源库带来压力量大的表要控制并发别为了对账把线上搞挂了。3.4 延迟监控和背压处理既要看数字也要看趋势延迟监控很容易走两个极端要么只盯一个 Kafka Lag 数字要么搞一堆复杂指标最后没人看。我的经验是关注三个核心项端到端延迟通过心跳表计算反映真实链路健康度。Kafka Lag反映缓冲区积压情况Lag 持续上涨说明下游消费能力不足。Flink 背压反映 Sink 写入目标端是否成为瓶颈。背压过高时整个管道吞吐下降Lag 自然上涨。迁移过程中最怕的不是系统性延迟而是间歇性抖动。比如每到整点源的备份任务启动binlog 量暴增Kafka Lag 突然从 1000 跳到 10 万过半小时又回落。这种抖动本身不可怕怕的是你把它当成正常现象没有给它的处理留出缓冲。所以 Kafka 容量至少要按高峰期流量的 3 倍设计Topic 分区数也要留够。实操中还要注意给源库留余量。Canal 默认拉取速率很高高峰期可能把源库 IO 打满。我一般会给 Canal 设置限流控制在源库峰值 IOPS 的 30% 以内。Flink 端并行度也不是越大越好并行度太大目标端写入压力飙升容易把新库写挂。先用小并行度跑稳定再加并行度这是一个慢启动的策略。4. 不停机迁移完整落地流程从盘点评估到灰度切换4.1 迁移前的盘点与目标端设计迁移不是把同步管道搭起来就开始搬而要先把搬哪些东西搞清楚。我接手任何一个迁移项目第一周基本不做任何技术选型而是在做盘点表清单哪些表要迁哪些表其实已经废弃哪些表是日志表可以归档不迁。很多团队一上来就说把线上 300 张表全部同步过去结果一半是没用的表。数据量与增长趋势每张表的行数、日均增长量、大字段情况。这个数据直接决定 Kafka 分区数、Flink 并行度和目标端存储规格。主键与唯一键情况没有主键的表要特别标注迁移时要额外处理建议跟业务方确认能不能补一个业务唯一键否则后续对账和幂等非常麻烦。源库 binlog 参数binlog 格式必须是 ROWbinlog 保留时间要足够。目标端的表结构不建议手工一张张去写 DDL工作量太大而且容易漏字段。我通常会写一个脚本读源库的information_schema.COLUMNS自动生成目标端建表语句然后根据目标库的模型特点做调整。比如 MySQL 的decimal(10,2)到 StarRocks 可以保留 decimal 或改成 double但要注意精度损失的问题涉及金额的字段我从不降精度。4.2 存量导入与增量追平接缝怎么拼才不漏不停机迁移最核心的难点就是存量数据和增量数据怎么拼起来中间不能有缝也不能重复到不可控。很多人在这里栽跟头。我给一个经过多次验证的处理思路第一步先启动增量同步管道。让 Canal 开始拉 binlog 写入 Kafka但 Flink 先不要消费或者消费后直接丢弃。此刻记录一个 Kafka 位点记为 B。第二步做存量快照导出。导出工具如 mysqldump 或 DataX执行时确保是一致性快照。MySQL 支持在可重复读事务里做快照导出这样导出的数据对应源库的某一个时间点把这个时间点对应的 binlog 位点记为 S。第三步把存量数据导入目标端。这一步可以用 DataX、Spark 或者直接并行 SQL 插入量大的要分批不要一条大 SQL 塞到底。第四步处理接缝。存量导入完成后让 Flink 从位点 B 开始消费并写入目标端。这里关键来了因为 B 是存量导出之前就记录的位点所以 B 到 S 之间的增量数据可能已经包含在存量快照里消费时会和存量数据重复。怎么办靠目标端的幂等去重。这个方案的精髓是宁可重复不可丢失。重复数据由目标端唯一键吸收如果目标端幂等没做那就老老实实从 S 位点开始消费但前提是 S 之前的所有增量都已经包含在存量里这要求导出期间业务变更的 binlog 不能被清理掉。存量导入完成后观察延迟追平情况。所谓追平不是延迟变成 0 就算完而是延迟保持低位比如 5 秒以内至少 30 分钟且对账任务跑出来零差异才算暂时稳了。4.3 双跑、灰度切换与回滚预案增量追平以后先别急着把线上流量全部切到新库。我习惯分三步走第一步只读流量灰度。把报表查询、后台管理类的只读流量切一部分到新库让真实业务跑几天。这一步既能验证目标端的查询性能和 SQL 兼容性又能通过业务方的日常使用暴露数据问题。第二步半写灰度。一些非核心的写操作比如日志写入、埋点上报切到新库观察同步链路反向是否通畅。这一步不建议核心交易类写入直接切风险太大。第三步正式切换。选择业务低峰期把写流量整体切换到新库。切换前停掉对源库的同步写入或者保持只读等待目标端追平到最近点位然后一键切换。回滚预案必须在切换前写好并演练过。最稳的回滚方式是新库承担写入后同步链路反向再把数据同步回旧库一旦发现问题把流量切回旧库即可。这里需要注意反向同步时旧库也要有幂等机制否则新库的重复写入会被带回旧库。5. 常见问题与排查技巧实录5.1 问题速查表现象可能原因排查思路解决方式Kafka Lag 持续上涨Flink 消费能力不足查看 Flink 背压指标、目标端写入耗时增加 Flink 并行度或拆分 Topic 分区目标端数据重复目标表无唯一键按主键查目标表是否存在多条给表加唯一键/改用 Unique Key 模型某张表数据量对不上表没有纳入同步或者被过滤检查 Canal/Debezium 配置的过滤规则修改过滤规则重新消费该表数据时间差 8 小时时区处理不一致比较源库 datetime 与目标端时间统一 Flink 任务时区用TIMESTAMP_LTZ处理Flink Checkpoint 失败目标端连接不稳定或超时查看 Checkpoint 失败点确认 Sink 重试配置增大检查点超时时间Sink 加重试更新字段丢失目标表列名或类型不匹配对比源表和目标表 DDL修正目标表结构从 Kafka 重放消息源库 IO 飙升Canal 拉取速率过高看源库 IOPS 监控给 Canal 加限流降低拉取速率5.2 几个真实踩坑案例案例一Kafka 重平衡导致同一主键乱序。当时一个订单表 Topic 分区数是 12Flink 并行度也是 12一切看似正常。但某天 Kafka 有节点重启触发了分区重平衡消费组重新分配分区后同一订单的 update 和 delete 消息被不同并行度消费到了不同子任务而且 Flink 没有按 Key 做 keyBy结果目标端出现旧值覆盖新值的情况。修复方案Flink 消费后必须keyBy(order_id)再写入同时在目标端保留更新时间字段用UPDATE时加条件判断只有新数据的版本号大于旧数据才覆盖。案例二Checkpoint 一直失败任务反复重启。排查发现目标端 StarRocks 连接池配置太浅Checkpoint 期间大批量写入把连接池打满导致 Sink 超时。调大连接池、增加写入批次大小、延长 Checkpoint 超时时间到 5 分钟之后问题消失。这个案例的教训是目标端的写入性能是整个链路的天花板上线前一定要做流量摸底。案例三对账发现订单金额对不上但没有重复和丢失。最终定位是浮点精度问题。源库金额字段是decimal(10,2)同步时通过 JSON 传到 KafkaFlink 解析成 double 再写入目标端double 的精度损失导致某些金额差了几分钱。从那以后所有金额相关字段一律用 decimal 或 string 传递绝不用 double。案例四迁移完成后发现一张状态表漏了几天数据。原因是这张表在源库被某个定时任务批量UPDATE了上百万行binlog 瞬间暴涨Canal 限流导致拉取慢Kafka 积压而这张表又没有单独的 Topic 和优先级被其他大表的流量挤掉了。后来给核心表单独拆分 Topic、单独配置消费并行度问题才解决。核心表一定要有独立的处理通道不能和非核心业务挤在一起。6. 我在实际项目里最受益的几个小习惯最后分享几个纯粹从项目里摔出来的习惯不一定写在哪本教科书上但确实帮我省了很多事。第一每张同步表都维护一个水位线记录。在目标端建一张sync_watermark表记录每张表最近成功同步到的 binlog 位点、Kafka offset 和检查时间。出问题的时候看一眼水位线就知道倒在哪一步不用翻一堆日志。第二所有 DDL 变更要预留重放窗口。源库如果要做表结构变更最好提前评估对同步链路的影响。有的 DDL 会导致 binlog 事件变大进而影响 Canal 拉取性能。迁移期间尽量减少源库大表的 DDL 操作必须做的提前通知同步链路值班人员。第三对账任务一定要比业务切换更早上线。很多项目把对账当成事后验证来做切换完了才开始写对账逻辑这完全反了。对账工具应该在迁移的第一天就上线全程盯着同步质量而不是等到切换前才慌慌张张地补。要说这套方案有什么缺点那就是前期建设和调优投入确实不小尤其是 Kafka 和 Flink 的运维门槛摆在那里。但如果你做的是交易、订单、账户这类核心数据的迁移这点投入比事故带来的损失划算太多。每一步都有人盯着每一笔账都对得上迁完库业务无感这种踏实感只有真的做过一次不停机迁移的人才能体会。