半夜两点,迁移群的欢呼声刚落下,业务方甩来一张对账报表:目标库少了一万三千条明细。监控大屏上,Kafka 消费延迟刚刚归零,Flink 作业的水位线一片平静,可账就是差。这种场面在异构数据同步项目里太常见了:大家盯了一整晚的“延迟”指标,眼见它终于安全,信心满满准备割接,结果一查库存、一拉流水,钱、单、号对不上。我做过的迁移方案不下十套,早期也栽过同样的跟头。今天这篇复盘,想完整拆解基于 KFS(Kafka + Flink + State/Checkpoint)的同步链路,重点不在“怎么把延迟压到 0”,而在“延迟归零之后,怎么守住每一笔账”。适合正在做数据库迁移、异构存储同步,或者正被业务催着“今晚必须切”的同学参考。
1. 延迟追平不是终点:账目对不上的四类根因
大多数迁移项目的监控面板里,Kafka consumer lag 是最醒目的一块。lag 降到 0,所有人都松一口气,误以为“源端产生的每一笔变更都已经落到目标库”。但真相是,lag 归零只说明“消息被消费端取走了”,不代表“消息被正确执行”,更不代表“目标库的数据和源库在业务语义上一致”。
我这些年遇到的账目对不上案例,归根结底基本都落在这四类:
1.1 位点语义被破坏:消费位置和写库结果没有绑定
Flink 作业从 Kafka 消费时,offset 提交和结果写库是两个动作。如果结果写库了、offset 没提交,任务一重启就会重复消费;反过来,如果 offset 先提交、结果还没落库,任务一挂就丢数据。很多人以为 Flink Checkpoint 能解决一切,但 Checkpoint 只保证算子状态的一致性,不保证你下游那个自己实现的写库函数天然幂等。一旦消费位点和目标库写入结果之间出现状态割裂,账目就会开始变得不可追查。
1.2 写库动作不幂等:同一笔数据执行两次结果不同
同一个 UPDATE 语句执行两次,有的场景没影响,有的场景直接翻倍。比如“金额 = 金额 + 100”的操作,重复执行一次结果就多 100。迁移链路里发生重复消费几乎不可避免,从 Kafka 的设计理念到 Flink 的故障恢复机制,都在默认“消息可能被重复投递”,如果你的目标端写入不带幂等保护,那账目对不上只是时间问题。
1.3 校验口径不一致:两边“讲的语言”不一样
源库是 MySQL,目标库是另一种存储,两边对 decimal 精度、varchar 尾部空格、时间时区的处理方式就可能不同。写个简单的 count 汇总往往显示行数一致,但金额汇总差几分钱,或某些字符串字段比对永远不一致。这个坑最隐蔽,因为它不一定是数据丢了,而只是两边对同一笔数据的描述方式不同。
1.4 切换窗口存在空洞:割接瞬间没有人搬数据
不停机迁移最后都要做流量切换。如果先切应用、再停同步任务,切与停之间的那段变更很可能两边都漏了;如果先停同步、再切流量,业务又可能出现不可写时段。这个窗口设计不好,账目就会出现一段“无主时段”,业务侧看到的现象是:这个时间段内的订单、流水在新旧两侧都对不上。
2. KFS 账本体系拆解:消息、计算、状态各守一摊
解决上述问题的关键,不是找一个更快的同步工具,而是把整个链路设计成一套“账本体系”。KFS 不是某个开源软件名,而是三个组件的配合方式:Kafka 负责消息账本,Flink 负责计算账本,State/Checkpoint 负责位点账本。三者各司其职,缺一不可。
典型链路长这样:
源库开启 binlog 或归档日志 → CDC 组件捕获变更 → 写入 Kafka Topic(按主键 hash 到分区)→ Flink 作业消费 Topic,做类型映射、清洗、幂等处理 → 写入目标库,同时写一张同步流水表用于对账 → Flink 周期性做 Checkpoint。
2.1 Kafka 这层账本:分区有序比全局有序更重要
Kafka 在整套架构里扮演的是“原始凭证仓库”。每条变更产生后先落 broker,按 key 的 hash 值进分区。key 的选取逻辑是这张账本的命门:必须用业务主键或业务唯一键,只有这样才能保证同一行数据的变更顺序在同一个分区内严格有序。
很多人纠结“要不要全局有序”,这个问题的答案是否定的。因为业务上只有同一行数据才关心先后顺序,跨行之间的顺序对最终一致性的影响通常可以忽略。如果为了全局有序把分区数设为 1,等于放弃了并行消费,延迟和吞吐会双双变差。分区有序、消费并行、按 key 路由,这才是在 Kafka 账本里“记好账”的正确姿势。
Kafka 的保留时间(retention.ms)也要刻意调大。迁移期间我通常至少保留 7 天,因为它是后面增量对账、故障重放的数据源。如果消息太早被清理,就算技术能力再强,也无法回答“某个 offset 上的消息原始长什么样”这个问题。
2.2 Flink 这层账本:Checkpoint 是对账的锚点
Flink 在链路里干的事,是把消息变成对目标库可执行的写操作。但它的职责远不止“消费-转换-写入”,更关键的是通过 Checkpoint 给整条链路一个可恢复的确定性。
每次 Checkpoint 相当于给账本拍一张快照,快照里记录了算子的状态和 Kafka 消费位点。任务崩溃后,能从最近的 Checkpoint 恢复,重新消费那一段消息。但要注意,Checkpoint 的默认语义是 at-least-once 还是 exactly-once,跟你的应用配置有关,更跟下游写入的实现有关。想要达到账目上真正的“不丢不重”,必须同时控制好 Checkpoint 配置和写库逻辑。
一个相对稳的 Flink 配置示例:
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); env.enableCheckpointing(60_000, CheckpointingMode.EXACTLY_ONCE); env.getCheckpointConfig().setMinPauseBetweenCheckpoints(30_000); env.getCheckpointConfig().setTolerableCheckpointFailureNumber(3); env.getCheckpointConfig().setCheckpointStorage("hdfs:///flink/checkpoints");Checkpoint 间隔 60 秒,最小暂停 30 秒,意思是最坏情况下每 60 到 90 秒生成一份快照。间隔太短会给状态后端带来压力,间隔太长则故障恢复时重放窗口会变大,恢复时间变长。这里需要根据业务容忍度去权衡,而不是照抄网上的模板。
2.3 三件套缺一不可:裸 Kafka 或应用双写差在哪
有人问,我不用 Flink,直接用 Kafka 做管道,消费者自己写库行不行?行,但你失去的是状态管理和恢复语义。裸 Kafka 只解决传输,不解决“任务挂了我从哪里继续”这个问题;消费者自己维护 offset,等于把账本逻辑打散在业务代码里,出了事要用人肉翻日志来对账。
也有人问,我不用 Kafka,直接在应用层双写源库和目标库行不行?风险在于双写不是原子的,应用一旦在写第一个库之后、写第二个库之前崩溃,这个窗口的数据就永久丢失。更麻烦的是,没有消息留存,就没有回放能力,一旦目标库数据有问题,你只能重新导全量,而不是精确重放某一小段。
KFS 的核心价值,是让“每一笔变更”都变成有据可查、可以回放的确定性操作。消息在 Kafka 里留底,计算过程在 Flink 里有状态,位点在 Checkpoint 里有锚点。这三样都齐了,才谈得上对账和收口。
3. 双轨对账落地:存量分片校验与增量流水核对的详细做法
有了链路还不够,迁移期间必须有一套完整的对账机制,而且是双轨:存量数据对一遍,增量数据持续对。
3.1 存量校验:把整库比对拆成百万个可并行的小账单
存量数据校验不能一口气全表扫描比对,尤其是亿级大表,一次查询就能把源库打挂。正确做法是按主键范围分片,每片几十万行,分片任务并发执行。
每片的任务逻辑是:读取源端该范围的原始数据,做规范化处理后拼成字符串,计算一个校验值(比如 CRC32 或 MD5,效率优先推荐 CRC32),再读取目标端同一范围的规范化数据计算校验值,两个值比对。不一致就进入明细比对,逐行找出差异字段,记录到差异单表。
差异单的字段至少要包含:分片 ID、主键值、源端值、目标端值、差异字段、检测时间。这张表是后续修复和跟进的重要依据。没有差异单,对账就变成了一锤子买卖,发现问题也无从下手。
3.2 增量核对:用左闭右开窗口锁住每一段流水
增量数据核对依赖同步流水表。每次 Flink 写目标库时,在同一条事务里写一张 sync_log 表,记录消息的 offset、源库事务 ID、主键、操作类型、写入时间。这相当于给每一笔变更在目标侧盖了一个“已执行”章。
对账任务定期扫描:从上一次核对位置开始,拉取 Kafka 对应时间段内的消息集合,和目标库 sync_log 表做两侧核对。这里有一个值得注意的细节,时间窗口必须设计成左闭右开,也就是起点包含、终点不包含,否则相邻窗口的边界数据容易被重复核对或漏掉。
窗口大小建议与 Flink Checkpoint 间隔对齐,或者取其整数倍。这样对账任务和快照节奏能形成一个稳定的对齐关系,不容易出现“对账追不上生产”的局面。
3.3 幂等键的账房规矩:哪些表能建唯一键,哪些表必须走版本号
增量核对只能发现问题,真正让“不丢不重”成立的是幂等写入。首选方案是目标表建业务唯一键,写库用 UPSERT 语义,例如 MySQL 的INSERT ... ON DUPLICATE KEY UPDATE,或者 PostgreSQL 的INSERT ... ON CONFLICT DO UPDATE。只要唯一键设计得对,同一笔消息重复投递十次,最终落库结果也只有一份。
但有些表确实没有天然业务唯一键,比如纯流水表、日志表。这时候需要在目标表加一个event_version字段,写入 SQL 变成带版本判断的条件更新:
UPDATE target_table SET amount = amount + #{delta}, version = version + 1 WHERE id = #{id} AND version = #{oldVersion};如果更新影响行数为 0,说明这条消息已经执行过,或版本已被更新的消息覆盖,直接丢弃即可。这套逻辑把“重复投递”变成了“无害投递”,是最后一道防线。
4. 延迟指标的正确读法:四个信号判断迁移能否安全收口
说了这么多账目问题,反过来再看延迟,它当然重要,但不能只看一个数。我把迁移监控里的信号归纳成四类,全部达标才谈得上“可以准备切换”。
| 指标 | 看什么 | 相对安全的信号 | 要注意的陷阱 |
|---|---|---|---|
| 消费 Lag 绝对值 | Kafka 积压未消费的消息量 | 持续下降,接近 0 | 归零只代表消息被取走,不代表执行完毕 |
| Lag 斜率 | 单位时间积压变化趋势 | 斜率为负,且稳定 | 大事务引发的瞬时飙升会误导判断 |
| Checkpoint 完成时间 | 快照写入耗时 | 稳定,持续低于间隔时间 | 连续失败说明恢复能力正在丧失 |
| 端到端水位线 | 从源库变更产生到目标库对账可见的耗时 | 秒级到分钟级,视业务容忍度而定 | 必须配合 sync_log 流水才能确认最终一致 |
4.1 消费 Lag 绝对值与斜率
我见过很多人只看 lag 的瞬时值,看到 0 就欢呼。但真正的做法至少要采样两次,看斜率。如果 lag 从 10 万降到 8 万再降到 5 万,说明消费能力大于生产速度,收敛方向是对的;如果 lag 在 0 和几百之间反复横跳,说明消费能力和生产节奏处在临界点,一旦来一个大事务就可能雪崩。
每次看 lag 最好用同一个消费组和同一个 topic,对比前后两次值。可以写一个小的轮询脚本,每 30 秒抓一次 lag,计算下降斜率,超过阈值就报警:
kafka-consumer-groups.sh --bootstrap-server $BROKER --group sync-job --describe生产环境建议接 JMX 或 Kafka Admin API,但在测试环境用命令先跑起来、建立对 lag 的量感,是很有效的入门方式。
4.2 Checkpoint 完成时间与失败次数
Checkpoint 是 Flink 作业的生命线。每次 Checkpoint 成功,才是真正“这之前的账目有据可依”的时刻。如果 Checkpoint 一直失败,说明状态后端存储有问题、或者作业负载过高,这时候即使 lag 是 0,也不能把作业状态视为健康。
运维上建议给 Checkpoint 完成时间加监控,指标名一般是flink_jobmanager_job_lastCheckpointDuration,超过 30 秒就要查一下是否有大状态、是否频繁 GC、存储是否出现瓶颈。同时监控 Checkpoint 的连续失败次数,连续 3 次失败就要人工介入。
4.3 端到端水位线
这个指标衡量的是“源库一条变更产生后,经过 CDC、Kafka、Flink,最终落到目标库并进入 sync_log”的完整耗时。它比 consumer lag 更能代表业务体感。计算方式不复杂:在源库变更里带上产生时间,目标侧 sync_log 记下写入时间,两者相减即可。如果目标库不支持额外时间字段,可以用 Kafka 消息时间戳和 sync_log 写入时间做近似估算。
4.4 收敛曲线与对账窗口覆盖率
最后一个信号偏工程管理:延迟不只要降下来,还要确定性地收敛。连续观察 3 到 5 个对账周期,如果每隔一段时间 lag 都会回到低位,且对账窗口的覆盖率是 100%(也就是每个时间段都有对应的 sync_log 记录可查),这时候才具备切换条件。瞬时为 0 不可靠,稳定可回放才可靠。
5. 延迟归零后的三次爆雷复盘:数据永远不会凭空消失
再讲几个真实踩过的坑。它们都发生在 lag 归零之后,每一次都实实在在地让账目出了问题。
5.1 爆雷一:decimal 精度不一致,校验任务从早跑到晚还是差一分钱
背景:源库 MySQL 的金额字段是DECIMAL(10,2),迁移到目标库时建表脚本写成了DECIMAL(10,4)。数据本身没问题,但做 CRC 校验时,源端把 9.99 拼成字符串 “9.99”,目标端则是 “9.9900”,hash 永远不一致。
排查链路:一开始怀疑数据丢了,但行数对得上;抽样看单条数据,发现值一样只是精度不同;最后看两侧 DDL,定位到建表脚本的类型映射错误。修复分两步:一是写校验函数时先统一规范化,所有 decimal 统一toFixed(2)再拼字符串;二是修掉建表脚本,重新跑增量。这个坑的启发是:对账不一致时不要急着怀疑数据,先看两侧口径。
5.2 爆雷二:Checkpoint 恢复后 UPDATE 被重复执行,金额翻倍
背景:一张用户余额变更流水表,同步作业消费 Kafka 后对目标库执行 UPDATE。当时以为加了事务就万无一失,但事务只保证“要么全做要么全不做”,不保证“只做一次”。某次作业重启后,Flink 从最近的 Checkpoint 恢复,那个 Checkpoint 是在一批消息写库之后、offset 提交之前完成的,结果这批复又被消费一遍,本应“金额+10”的 UPDATE 执行了两次,余额多了 10 块。
排查链路:先看 sync_log 表,发现同一 offset 的消息出现两次;再看 Flink 恢复日志,确认重复消费的起点;最后检查写库 SQL,发现完全没有幂等条件。修复就是前面提到的版本号方案,UPDATE ... WHERE version = #{oldVersion},重复消息如果版本不匹配会被自然过滤。从那以后我坚持一个习惯:任何迁移链路的写库代码,都必须假设“这条消息至少会看见两次”。
5.3 爆雷三:流量先切、同步后停,切换瞬间的变更没人接管
背景:凌晨预期是先停同步作业、再切流量,但上线脚本里两个步骤顺序写反了。流量先切到了新库,同步作业还在跑,结果切换后新写入的一部分业务数据进了新库,同步作业又把更早的变更追过来,两边互相覆盖。另一个方向上,同步停止后到完全停服期间产生的变更没人搬,目标库缺了一段数据。
排查链路:切换后做增量核对,发现某个时间段内的流水两侧都有缺口;翻运维日志确认步骤顺序和预案不一致;再手工比对切换时间点前后的变更,最终定位到那段“无主窗口”。修复方法是把最后切换设计成固定顺序:先停同步、再切流量、再对账、确认无误后关闭源库写入,并且对账窗口必须覆盖切换前后各一段,留足缓冲区。
这三个爆雷有个共同点:它们都不是被 lag 监控发现的。等 lag 归零之后,唯一的守护者就是那套对账体系和幂等机制。
6. 收口阶段的六项检查清单
迁移收口前,我每次都会过一遍自己的检查清单,省掉一次事故的概率比想象中更大。
- 存量校验全跑完,差异单已处理:所有分片校验任务状态为 SUCCEEDED,差异单要么已修复,要么有业务方确认的可接受说明。
- 幂等键全覆盖:逐一核查目标表,没有业务唯一键的表都已补上版本号方案,不存在“裸更新”。
- 延迟收敛趋势确认:连续观察至少 3 个对账周期,lag 在下降或者稳定在低位,没有再次抬升的迹象。
- 做过一次真实的故障演练:手动杀掉 Flink 作业再恢复,等链路追平后做一轮增量核对,确认 sync_log 没有缺失也没有重复。
- 切换窗口留出“最后体检”时间:流量切换后不要急着回收源端资源,至少保留一个完整对账周期,并让对账任务继续跑,直到确认不再有新增差异。
- Kafka 消息保留时间覆盖完整对账周期:迁移结束后不要马上清理 topic,至少再留一个保留周期,方便处理可能出现的迟滞发现的问题。
迟滞发现的问题在迁移项目里并不少见,有些差异要到业务跑完一个月的结算周期才浮出来。Kafka 里的消息就是你的后悔药,删了就真的再也查不到了。
做迁移方案这些年,我的最大感受是:方案文档里不仅要写“如何让数据跑得快”,还必须写清楚“如何知道数据没跑错”。延迟是结果,账目是底线。把对账设计成链路的一部分,而不是上线前的临时动作,KFS 这套组合的价值才能真正发挥出来。