1. 迁移时的数据账本,为什么比延迟数字更值得盯
上个月做支付核心库从 MySQL 迁移到分布式数据库,业务方第一天就追着问:现在延迟多少毫秒?我说你先别盯延迟,去看对账单有没有平。异构数据同步这个圈子,大家津津乐道的往往是 RTT、堆积量、消费速率这些指标,可真到不停机迁移的时候,延迟只是面子,每一笔账能不能对齐才是里子。这套保障方案我们内部叫它 KFS,它不是某个新开源框架,而是 Kafka-Flink-Sink 三个环节组成的同步守护链路:Kafka 当总线解耦异构数据源,Flink 做变更处理和数据整形,Sink 层负责幂等落账和审计。今天就把这套系统的设计思路、切换流程和踩过的坑完整写一遍。
1.1 延迟指标暴露不了“丢账”和“错账”
延迟低,只能说明消息从源端流到目标端的通路顺畅,但它完全无法回答三个更重要的问题:源端产生了一百笔订单,目标端是不是也收到了一百笔?收到的那一百笔里,有没有哪笔的金额被字段截断了?订单状态从“已支付”到“已发货”的变更,在目标端是不是也按顺序发生了?
我见过一个真实的惨例。某团队做订单表迁移,延迟一直稳定在 1 秒以内,监控面板很漂亮。结果切流后第二天,线上出现大量订单卡在“已支付、未发货”状态。排查到最后发现,问题出在同步链路丢了一条 update 语句:源端订单状态先更新为“已支付”,随后又被补偿流程更新为“已发货”,但由于目标端唯一键冲突,第二条 update 被静默丢弃。业务看起来一切正常,延迟也一直是绿的,可账在那一秒已经坏了。
这就是典型的“管线视角”和“账本视角”的差别。只看管线,你关心的是速率、积压、网络抖动;看账本,你关心的是每一条变更是否都按时、按序、按原值落到了目标端,并且能随时回答“两边是不是一样的”。
1.2 不停机迁移放大了不一致的代价
日常双跑阶段,源端还是权威系统,目标端数据错了可以重刷。但不停机迁移的可怕之处在于:切流之后,目标端在某一个瞬间成为唯一的业务真相。那时候再发现历史数据有差异,就不是重跑一个同步任务那么简单,而是要面对资损、客诉、甚至回滚的连锁反应。
KFS 在设计时默认了一个原则:迁移期间要把每一笔数据都当成一次“跨系统资金划转”来对待。划转需要凭证、需要账目、需要可追溯,数据迁移也一样。每条变更记录都要能回答:我从哪里来(源端位点)、我经历了什么(处理状态)、我最终落到哪里(目标端主键和落库结果)。这三个信息连起来,就是一条完整的审计轨迹。
2. KFS 的组件拼图:Kafka 解耦、Flink 守序、Sink 兜底
KFS 不是一套新写的同步软件,而是把三个成熟组件按迁移场景重新组织起来:Kafka 负责接入和缓冲,Flink 负责处理与状态管理,Sink 层负责落库与幂等。每一层都有明确的边界,也都有针对“每一笔账”的专门设计。
2.1 Kafka 层:为什么不建议源端直连目标端
很多团队的异构同步最初是“源端直写目标端”:写一个脚本从 Oracle 抽数,直接 JDBC 灌到目标库。这种方式在数据量小的时候很爽,但迁移一旦涉及几十张表、多个数据源,问题立刻就来了:源端的一次抖动会直接传导到目标端,目标端的一次锁等待也可能反过来拖垮源端业务。两个系统耦合在一起,账出问题了都不知道该查哪边。
Kafka 在中间当总线,本质上是给两个系统之间加了一个“缓冲隔离带”。Kafka 能长期保存数据,源端产生变更后只要写进 Kafka 就算成功,目标端消费的快慢不会反向影响源端。这个特性在不停机迁移中特别重要:迁移期间你会频繁地对目标端做表结构调整、索引重建、数据校验,目标端随时可能要停住,缓存层正好给了你从容操作的时间窗口。
Topic 的划分也有讲究。我建议按业务域拆分,比如订单域、会员域、支付域各一个 topic,分区键一律选业务主键。这样同一个订单 id 的所有变更永远落到同一个分区,Flink 消费时天然有序,目标端才不会出现同一主键的乱序覆盖。
2.2 Flink 层:把“乱序变更流”变成“可重放账本”
Kafka 能保证单分区内有序,但跨分区、跨表、跨 topic 的变更之间仍然存在顺序问题。比如订单主表和订单明细表是先删明细再删主表,还是先删主表再删明细,这个顺序一旦错位,外键约束就会直接拦你。Flink 在这里干的活,是把 Kafka 里的原始变更流整理成一份“可重放的账本流水”。
具体来说有三个关键职责。
第一是去重与排序。通过 Flink 的 keyed state 记录每个业务主键最近一次处理的时间戳和位点,遇到乱序事件要么丢弃、要么延迟处理,避免旧数据把新数据覆盖掉。第二是维表补全。异构系统之间字段命名和编码往往不一致,源端写的状态值是 1、2、3,目标端要求的是字符串枚举,这种转换放到 Flink 里统一处理,Sink 层就能专心落库。第三是脏数据隔离。格式不完整、字段缺失、主键为空的记录,不要直接抛异常让整个链路卡死,而是写入特殊的死信 topic,同时记录原始位点,方便事后补偿。
迁移期间 Flink 有一个参数我建议特殊对待:Checkpoint 间隔。日常同步可以设成 60 秒一次,但迁移切流前后,我会把间隔缩短到 10 秒到 15 秒。Checkpoint 越频繁,故障恢复时回放的数据量越少,两边账目的差距越小。对应的代价是状态后端压力变大,所以存储目录最好用 SSD 而不是机械盘。
2.3 Sink 层:幂等是守住账目的最后一道防线
很多人以为同步任务只要“消费成功”就等于“写库成功”,这是误解。Flink 的 Checkpoint 机制只能保证“这条消息被 Flink 处理了”,不能保证“这条消息对应的事务已经提交到目标库”。两者之间一旦出现断电、网络闪断、目标库回滚,下游就可能出现重复数据或丢失数据。要守住账,Sink 层必须自己具备幂等能力。
我们这里的做法是强制的:目标表的每条业务记录都带一个sync_uk字段,取值是源端实例ID + binlog文件名 + 位点 + 业务主键。写入时使用INSERT ... ON DUPLICATE KEY UPDATE或等价语义。这样即使 Flink 因为故障从上一个 Checkpoint 重放了一批数据,重复执行也不会产生重复记录,而是原地更新。Sink 层同时维护一张“同步流水表”,每条数据落库成功后写一条流水,记录处理时间、位点、影响行数。这张流水表就是前面说的“审计轨迹”的实体,后续所有对账查询都从它出数。
| 层级 | 核心职责 | 账号目安全对应的关键点 | 迁移期重点参数 |
|---|---|---|---|
| Kafka | 数据接入与缓冲 | 消息留存时间、分区有序性 | retention.ms调大至 7 天以上 |
| Flink | 清洗、排序、状态管理 | Checkpoint 恢复点、去重状态 | checkpoint.interval=10s,状态后端用 RocksDB |
| Sink | 幂等落库与审计 | 唯一键设计、流水表记录 | 开启INSERT ... ON DUPLICATE KEY UPDATE |
3. 全量与增量怎么衔接,才能让两套系统同时算对账
不停机迁移最核心的技术难点,不是全量数据怎么抽,也不是增量数据怎么同步,而是全量和增量交界的那个瞬间,怎么保证同一笔数据不会算了两遍或者漏算一遍。KFS 在处理这个问题时用了一条非常朴素的策略:先锁定账本起点,再抽取存量,最后回放增量。
3.1 全量基线:快照读与分批拉取
全量抽取阶段,我们对源库的压力控制得很严格。首选方案是从只读从库拉数,如果业务允许,也可以在备库上做,目的就是不要把主库的 IO 打满。每张表的抽取按主键范围分批执行,一批 5000 条左右,拉完一批记录一个批次的“最大主键值”作为断点。这样即使任务中途挂掉,也能从断点继续,不需要重头再来。
这一阶段最容易犯的错是:全量抽取时用默认的查询隔离级别,导致同一张表在不同时间点读到了不同快照。比如钱表读了 100 万条,订单表却是在那之后 5 分钟才开始读的,此时订单表已经新增了 3000 条新数据。两边基线时间不一致,接下来的对账就会一直对不上。KFS 的解决办法是:启动全量任务前,先记录一个数据库统一的“水位时间”(MySQL 可以用SELECT NOW(6),同时记录 binlog 位点),所有全量抽取 SQL 都加上这个水位时间的查询条件。虽然无法完全替代事务快照,但在业务低峰期操作,配合从库,误差已经可以做到可接受范围。
3.2 增量接续:从位点而不是从时间点开始
全量基线跑完后,还缺一个关键步骤:把这段时间里源端新产生的增量变更补到目标端。如果你在 T0 时刻记录了 binlog 位点,然后全量跑到 T1 结束,那么 T0 到 T1 之间的变更就需要从 T0 位点开始回放。
这里的细节在于,Kafka 中已经存在 T0 之前的数据,直接从头消费会造成重复入库。所以我们的做法是:Kafka 消费者在初始化时,按照记录下来的源端位点去定位——用 Flink 的 Kafka Source 的setStartFromTimestamp或自定义 partition discover 机制,找到对应时间戳的 Offset 再开始消费。同时,正在运行的全量任务不能立即停止,要等增量消费追平到“全量结束时间点”之后,才做一次数据合并校验。
数据合并的规则我们称之为“位点后写覆盖”:对同一个业务主键,如果增量变更的位点晚于全量快照的位点,以增量结果为准。实现上并不需要逐条比对位点,只要顺序合法,直接让增量执行的 update 覆盖目标表即可。唯一要注意的是,全量任务和增量任务可能同时写同一条数据,Sink 层必须保证最终位点落在更新更晚的那一侧。
3.3 双跑期的日切对账:三账户配上哈希校验
迁移双跑期不是只跑数据,而是每天都在对账。KFS 每天凌晨固定跑一轮对账批处理,规则很简单:笔数守恒、金额守恒、状态机守恒。
- 笔数守恒:每张业务表在源端和目标端的记录数相等,按天按状态字段分组各比一次。
- 金额守恒:所有金额字段的 SUM 值在两边一致,按币种、按渠道维度分别核对。
- 状态机守恒:对“订单状态”这类有明确状态流转的字段,统计每种状态值在两边分布的条数。
光比对总数还不够,还需要防止“两张表总数相同但具体某几条数据内容不一致”的情况。我们会对每张表的主键和关键业务字段拼接后做哈希,源端和目标端各自算出哈希再比对。一致性比对只抽查 5% 到 10% 的数据,成本不高,但能有效覆盖字段值被截断、时区错位、字符编码转换出错这类问题。
对账发现差异时,先别急着灌数修复。第一步是查同步流水表,看差异数据最近一次落库的位点、时间和影响行数;第二步是反查 Kafka 原始消息,确认源端的变更记录到底是怎样的;第三步才是人工判断是忽略还是补数。这样做的目的,是避免直接修正数据把真正的问题掩盖掉。
4. 切换那一刻:灰度切流、校验窗口与回退预案
迁移项目做九十天,真正让人睡不着的可能只有切流那两小时。KFS 在切换环节没有搞“一键切换”,而是把整个过程拆成了可验证、可回退的小步骤。
4.1 切流前置检查清单
切流前我要跑一遍固定的检查清单,全部通过了才允许动开关:
- Kafka 消费积压接近归零:目标端的消费位点已经追上源端最新位点,差距在几百条以内。
- Flink Checkpoint 连续成功至少 50 次,最近一次恢复演练成功。
- 前一天的自动对账结果是零差异,或者所有差异都已经人工确认并修复。
- 源端与目标端统计出的关键业务表行数完全一致。
- 同步流水表记录数与源端 binlog 事件数在可解释的误差范围内。
- 已通知业务方维护窗口,并且备份了切流前后各一张核心表的快照。
这份清单看起来平平无奇,但它最大的价值不是“防止出错”,而是“在出错时知道自己走到了哪一步”。一旦后续发现问题,往回看清单,就能判断问题到底出在切流前还是切流后,从而决定是修复继续,还是直接回退。
4.2 灰度切流的操作顺序
切流从来不是“源端停写、目标端开写”这种二元操作。KFS 建议的思路是按流量维度灰度:
第一步,只读流量先切。把报表查询、管理后台这类只读请求指向目标端,观察目标端在真实读压力下的表现。这个阶段允许误差,发现数据不对可以随时把读流量切回源端。
第二步,小比例写流量切。选择某个租户或某一批用户 id,把他们的新写入直接落到目标端,同时源端仍然持续同步。此时要特别盯一个指标:源端同步链路是否还在正常工作。实践中很容易出现一种诡异场景:小流量切到目标端后,源端与目标端两边同时写入同一条业务主键,两边的同步链路开始互相覆盖。为了避免这种“双主冲突”,我们将已切流量的业务主键范围在流水中打标签,同步过滤掉这部分变更,目标端只保留新写入的数据。
第三步,全量写流量切换。在所有小流量验证通过后,源端停写、目标端接管全部读写。这个过程建议放在业务低峰期执行,并在切换前提前通知所有上游应用刷新配置。
4.3 KFS 的水位对齐判定标准
“延迟归零”并不能直接说明“两边已经一致”,因为延迟只代表最近一条消息被消费了,不代表所有消息都按顺序完成了处理。KFS 在切换前还会做一个“水位对齐”检查:选一张核心业务表,记录源端最新变更的位点,然后在目标端找到该位点对应的数据,确认它已经可以查询到。对比两侧的 watermark 时间,正常应该相差不超过 30 秒。
切流后的观察窗同样重要。我们固定观察 15 分钟,期间每分钟跑一次最小化对账(只比最新五分钟内的增量和关键字段哈希)。任何一次对账出现差异,立即暂停剩余切流步骤,并触发回退预案。这里我特别想说明,回退不是“把流量切回源端”这么简单,还需要把目标端在此期间产生的新数据反向同步回源端。KFS 的做法是:目标端也开启一套反向的 Kafka-Flink 同步任务,只是平时处于暂停状态,一旦需要回退,立即启动。这套反向通道平时不花钱,但关键时刻能救整个项目。
5. 跑稳一年之后,我在这个方案里踩过的坑
架构说得再漂亮,最终都要落到一次次的故障排查里。下面这几个坑,都是我在真实迁移项目中踩过、并且最后在 KFS 设计上做了针对性改进的。
5.1 Kafka 大事务导致的顺序反转
第一个坑发生在消费超时。当时源端有一笔批量退款,单个事务里更新了上千条记录。Kafka 客户端处理这批消息时,单个事务的总耗时长于max.poll.interval.ms默认值,触发了消费者 rebalance。Rebalance 之后分区归属变化,有一部分消息被重新分配给了另一个消费者实例。因为另一个实例的本地状态没来得及同步,导致同一条记录的旧事件反而比新事件更晚被写入目标端。目标端看到的结果就是:一笔订单先变成了“已退款”,而后又变回“待退款”。
这个坑排查了两天才确认根因。修复措施有三条:把max.poll.interval.ms调整到 5 分钟,同时把max.poll.records调小到 500,避免单次 poll 处理时间过长;另一个措施是在 Flink 侧对相同业务主键做基于事件时间的窗口去重,确保旧位点的事件不会覆盖新位点。更重要的是,从业务侧限制了单事务操作行数,大事务拆成小批次。事实证明,业务侧配合比技术侧硬扛要有效得多。
5.2 Flink 状态后端选型对恢复时长的影响
项目初期我们用默认的 HashMap 状态后端,同步的表少时没问题。随着迁移表数量增加到上百张,Flink 做 Checkpoint 的时间越来越长,最夸张的一次整整花了 6 分钟。那段时间只要有一个 TaskManager 宕机,整个作业恢复就要从最后一个 Checkpoint 重放大量数据,恢复期间目标端数据缺口急剧拉大。
换成 RocksDB 增量 Checkpoint 之后,情况好了很多。但 RocksDB 也有自己的脾气:本地磁盘 IO 如果跟不上,反而会拖慢正常处理。我们把状态目录挂在 SSD 上,同时给 TaskManager 配置了独立的临时目录,不让 Checkpoint 与日志写同一块盘。这里给个经验值:单作业状态在 10GB 以下,用 HashMap 就行;超过 10GB 并且恢复时间超过 10 分钟,果断切 RocksDB。
5.3 重复消费导致的唯一键设计失误
典型场景是 Flink 作业手动重启时,由于 Checkpoint 没有成功保存,从上一个位点重新消费了一批数据。我们最初给目标表设置唯一键时只用了业务主键,第一次重复消费时数据原地更新没有问题,但第二次重复消费同一批数据时,涉及“先删后插”逻辑的表就出现了问题:删除事件重复执行,把不该删的记录删掉了。
排查之后我们把唯一键改成了sync_uk,也就是前面介绍过的那个“源端实例ID + binlog文件名 + 位点 + 业务主键”的联合值。这样重复消费同一批数据时,每一次的sync_uk都一样,数据库的 upsert 语义会让后面的执行结果覆盖前面的,不会产生额外的删除。这个改动听起来很小,但它彻底消除了重复消费这批数据时所有潜在副作用。
5.4 时区字段差八小时引发的“假不一致”对账
另一件好笑又耽误事的问题:对账脚本跑出来全是差异,最后发现是时区。源端 MySQL 的datetime字段不带时区,目标库是分布式数据库,JDBC 连接默认时区是 UTC。同步框架在写入时把 MySQL 的本地时间当作 UTC 处理,导致目标端时间字段整体比源端晚了 8 小时。业务时间字段错了,对账脚本按小时分组统计时自然对不上。
修复方式分两层:第一层是技术口径统一,所有 JDBC URL 显式指定serverTimezone=Asia/Shanghai,并且在 Flink 的序列化器里对时间类型做统一转换,不再依赖环境变量;第二层是业务口径明确,时间字段分成“业务时间”和“技术时间”,业务时间在迁移中保持原值不转换,技术时间统一用 UTC 存、展示时再转换。这样以后再看到对账差异,先检查是不是时区问题,省下大量排查时间。
写在最后的一个小建议
每次有新项目来咨询迁移方案,我都会让对方先回答三个问题:你的数据账本是什么?哪些字段能定义“一笔账是同一笔”?如果对账有差异,你的止损线在哪里?这三个问题不想清楚,再好的同步工具也只是给错误加速。KFS 这套方案真正值钱的地方,不在于它用了多新的技术,而在于它把“每一笔账都守得住”当成设计的第一原则。如果你也在准备不停机迁移,建议先从你最重要的一张订单表开始,手工模拟一遍全量加增量、切换加回退的完整流程。跑通一次之后,你会对“延迟”这两个字有完全不一样的理解。