不停机数据库迁移实战:基于Kafka与Flink的异构数据同步方案 📅 发布时间:2026/9/15 21:28:28 👁 浏览次数: 搞数据迁移的人最怕听到业务方说一句话“系统不能停你看着迁吧。”我第一次接这种活的时候以为有工具做异构数据同步把数据从A库搬到B库就行。后来被坑了几次才明白真正难的从来不是“搬”而是搬完之后两边每一笔账能不能对上。尤其在订单、流水、账户余额这类核心场景里错一笔就是生产事故漏一笔就是数据资产损失。这篇文章我们就聊聊不停机迁移场景下的异构数据同步重点讲一套我线上用过的方案——KFSKafka Flink Sync。它不是一个神乎其神的中间件而是一套把延迟、幂等、校验都变成可量化指标的方法论。适合正在做数据库迁移、实时数仓建设、或者被增量同步折磨的工程师看完至少能少走半年弯路。1. 不停机迁移的底层逻辑为什么“追延迟”是最容易走偏的目标1.1 业务方要的“不停机”到底是什么很多刚接触迁移的人会把“不停机”理解成“迁移过程中数据库完全不能被碰”。实际上业务方真正要的是应用不能停用户请求照常处理只是背后数据在悄悄搬家。最常见的就是一个交易系统要从MySQL换到分布式数据库或者从旧集群迁到新集群。按老办法凌晨两点停服四个小时全量导出再导入听起来可行但放到现在7x24的业务上别说四个小时哪怕十分钟的订单写入中断都会引起客诉。不停机迁移的本质是让新旧两套系统在某个时间段内“共存”等两边数据完全一致了再把流量平滑切过去。这个过程中源库继续接受业务写入同步工具把增量变化实时搬到目标库。听起来简单但真做起来你会发现事情远不是复制粘贴那么简单。尤其是异构场景——源库MySQL目标库可能是PostgreSQL、OceanBase、TiDB甚至可能是Hive、Iceberg这样的数仓存储。库和库之间不只是版本不同而是底层逻辑都不一样。1.2 异构数据同步的三个硬骨头第一个硬骨头是结构差异。MySQL里的datetime和PostgreSQL的timestamp精度不一样Oracle的NUMBER和MySQL的DECIMAL精度上限不一样还有自增主键、枚举类型、JSON类型稍微不注意就映射错。我在线上见过最典型的问题源库一个字段是tinyint(1)同步过去被自动识别成boolean应用层查出来直接变true/false业务逻辑当场崩掉。第二个硬骨头是语义差异。数据库不是简单的KV存储它有事务、有外键、有唯一约束、有触发器。同步到异构目标库时这些语义不一定能原样保留。比如源库做了一次先update后delete如果同步链路没有保证顺序目标库就可能出现先delete后update最终数据结果完全不同。又比如源库一个事务里写了10条数据同步端如果一条一条并发发出恰好其中一条失败就会造成部分提交目标库就多了一笔脏账。第三个硬骨头是位点一致性。增量同步永远要回答一个问题我到底从哪个位置开始消费源库的binlog position、Kafka里的offset、Flink的checkpoint这三个位点必须对齐。任何一方没对齐要么丢数据要么重复消费。这也是为什么我反复跟团队说追延迟之前先追位点延迟是表现位点才是根。2. KFS的整体架构以账目为准而不是以延迟为准2.1 KFS是什么一套围绕Kafka生态的同步框架我们内部说的KFS全称是Kafka Flink Sync核心思路很简单用Kafka当消息主干用Flink CDC当增量捕获引擎再用独立校验服务兜底。整个链路是源库 - Flink CDC - Kafka - Flink SQL/Jar任务 - 目标库 - 校验补偿服务。这六个环节各司其职形成一个闭环。有人可能会问为什么不直接用Canal或者Debezium往目标库灌因为同步链路中间一定要有一个可回溯、可重放的缓冲层。Kafka在这里的作用就是缓冲。它本身持久化消息默认可以保留几天甚至几周如果下游任务挂了修完之后还能从之前的位点重新消费不会丢数据。没有这个缓冲层源库的binlog过期了或者消费端崩了几天想补数据都无从下手。这也是我坚持在同步架构里放Kafka而不是让CDC直连目标库的原因。2.2 三层结构通道层、转换层、对账层把KFS拆开看它其实是三个逻辑层。通道层解决“数据怎么传”。所有源库变更事件先被序列化成统一的JSON格式再写入Kafka对应topic。topic按业务域拆分比如订单域一个topic用户域一个topic每张表按主键或者业务键做分区保证同一行数据的变更事件永远落到同一个分区这样消费端能严格按顺序处理。转换层解决“数据变成什么样”。Flink任务从Kafka消费事件做类型映射、字段过滤、脏数据清洗然后写入目标库。这一层是最灵活也最容易出问题的地方。我的经验是所有转换逻辑必须显式定义不要让框架自己猜。比如MySQL的timestamp转目标库的timestamp时区怎么处理字段精度不够时是截断还是报错这些都要在任务配置里写死。对账层解决“数据到底对不对”。KFS会在每个同步任务旁边挂一个校验引擎定时把源端和目标端的同一份数据拿出来做对比。对比的维度不是逐行比对——那样太慢也太贵而是按照分片统计行数、计算checksum、抽样对比关键字段。发现差异就自动生成一条补偿记录触发对应的修复任务。2.3 设计哲学延迟是结果不是目标KFS最反直觉的一点是它不追求把延迟压到最低。为什么因为延迟太低有时候反而危险。举个例子如果源库正在跑一个超大的批量更新同步端如果实时地把每一条变更都立刻写到目标库目标库会被打个措手不及CPU和磁盘IO瞬间拉满反而拖慢整个同步链路。这个时候主动让Kafka消费端“排队”反而更合理等一批变更集中起来批量写入目标库整体吞吐反而更高。还有一个更要命的问题就是事务中间状态。MySQL一个事务里改了100行binlog会记录100条变更事件如果同步端消费一条写一条正好在第50条的时候任务崩溃目标库就留下了半个事务的数据。这种数据是最难查的它既不违反唯一约束也不会报错但就是跟你源库对不上。为了避免这种情况我们经常故意让消费端延迟一段时间再处理事件等一个事务的所有变更都到齐了再一次性应用。所以后来我在项目里经常看到“kafka 如何延迟30分钟消费”这类搜索其实就是这个场景的需求。延迟能压到秒级当然更好但它是正确性解决之后才考虑的事情。KFS的默认策略是在账目一致的前提下把延迟控制在30秒以内。如果超过这个阈值触发告警优先排查链路瓶颈而不是无脑加并行度。3. 基于KFS的完整迁移实战从存量到增量到切换3.1 阶段一存量全量数据搬迁不停机迁移第一步一定是先把源库已有的历史数据搬到目标库。这个阶段最怕的是把源库压垮。网上很多教程会说开32个并发一次性select全表听起来很快实际上生产环境这么搞源库CPU直接飙到100%业务查询全堵住。正确的做法是分片抽样、分批搬迁。我们当时的订单表有10亿行按主键id范围切分成1000个切片每个切片只查1万条用Flink批任务并发16去搬每批500条提交一次事务。这套参数跑下来源库负载只增加了5%左右业务基本无感。切片大小怎么定我用的是一个经验公式单批数据量控制在5000到10000行或者5MB以内保证一次查询不会触发源库的大事务复制延迟。并发数控制在CPU核数的一半以下给业务留足余量。全量搬迁的时候源库还在写入新数据所以必须在全量开始前记录一个源库的binlog位点或者GTID保证后续增量能从正确的位置接上。这个位点记录得越早全量过程中变更的数据就越需要靠增量链路去补。KFS的做法是把位点作为一个独立事件写入ZooKeeper或者etcd全量任务和增量任务都读这个位点保证它们之间没有缝隙。3.2 阶段二增量日志订阅与消费位点管理全量搬完之后增量任务启动。Flink CDC是这里的主力它对MySQL binlog和PostgreSQL WAL的解析都比较成熟。启动参数里最关键的startupOptions我们用的不是默认的latest而是specific-offset也就是从全量阶段记录的那个binlog位点开始消费。直接加上特定偏移。这样能确保全量搬完的数据和增量产生的数据完美衔接不重也不丢。增量事件进入Kafka之后topic的参数也得提前规划。我们用三个分区、三个副本retention.ms设置成604800000也就是7天。为什么分区数不多不少用三个因为下游Flink的并行度一般是2到3分区数再多了会造成有的分区空闲、有的分区堆积白白浪费资源。为什么保留7天因为线上任务出问题恢复时间往往不是按小时算的是按天算的七天足够我们从容排查和补数据。这里还要特别强调Kafka位点和Flink checkpoint的对齐。Kafka的offset是消费端的位点Flink的checkpoint是任务状态的快照。每次checkpoint成功之后我们应该把“当前消费到的Kafka offset”和“处理完成的数据范围”一起保存下来。这样即使任务崩溃重启从checkpoint恢复后它知道自己处理到哪个位置了不会把已经写入目标库的数据再写一遍。KFS的位点管理表里会同时记录这三个信息Kafka offset、binlog position、目标库已提交事务ID。三方对不上立刻告警。3.3 阶段三回放、追平与安全切换增量链路跑起来之后新旧两边的差距会逐渐缩小最终达到一个稳定状态。但这个状态不是“停住不动”而是源库每产生一笔变更目标库几乎同时追上一笔。此时就进入切换决策阶段。很多人以为切换就是把域名改一下流量指到新库就行。实际上我们内部有一套完整的切换红线。第一连续30分钟同步延迟小于10秒不能只盯着某一分钟看峰值延迟说明链路还不稳定。第二校验任务连续三次零差异。第三最近一小时内增量任务没有一次失败重试。三条全过才允许切第一批流量。切换一定要灰度。从5%的流量开始观察15分钟看新库的写入延迟、报错率、慢查询有没有异常再切到20%、50%、100%。这个过程中旧的同步链路绝对不能停因为一旦新库出问题需要立刻把流量切回源库。我们管这个叫回切预案。有一次线上切到50%的时候新库的分布式事务锁出现死锁我们直接执行回切整个切换在10分钟内完成业务无感知。如果没有这条后路那就是大事故。3.4 阶段四并行校验与持续观察流量全部切到新库之后不代表迁移结束只能算“基本完成”。真正的考验在切换后的头七天。这期间源库还在正常服务但已经不再承担写入流量所以我们可以放心大胆地对源库和目标库做最终的全量校验。KFS在这个阶段跑的不是增量校验而是全量校验。具体做法是把源库的每张表按主键hash分片每片大概10万行计算行数、数值字段的sum、以及一个组合checksum写入校验结果表。然后对目标库做同样的计算调度任务负责逐片对比。有差异的分片会进入再次校验队列连续三次都不一致就需要人工介入了。这个“并行校验”的好处是它可以在不阻塞任何读写的情况下把迁移的最终结果用数据说话。我见过太多团队切换成功后三天就把同步链路拆了结果一个月后业务查出来一笔老数据对不上那时候源库binlog早就清了数据彻底找不回来。所以我的建议是同步链路保留最少一周每天自动跑一次校验零差异持续48小时再拆。这个习惯帮我避过好几次“假成功”的坑。4. 守住每一笔账一致性保障与延迟补偿机制4.1 幂等写入与Exactly-Once语义为什么同步链路会重复写数据最常见的原因就是Flink任务重启。只要任务重启checkpoint回放就会把上一次处理过的事件再消费一遍。如果目标库写入不是幂等的就会出现重复数据。解决这个问题第一道防线是目标库的写入操作必须基于主键或唯一键做upsert而不是简单的insert。KFS在写入目标库时统一使用on duplicate key update语义同一个主键的数据无论来多少次最终状态都一致。第二道防线是Flink端开启exactly-once语义。以JDBC sink为例需要在建表语句里设置sink.semantic为exactly-once配合Flink的两阶段提交协议。这里要注意两阶段提交不是免费的它会增加目标库的事务开销写入吞吐会下降20%到30%。所以具体开不开要看场景。我们的经验是核心交易类表必须开报表类、日志类表可以不开采用at-least-once加上目标库幂等键兜底。有一次我们排查数据重复发现是目标库的分布式自增主键生成策略有问题两个节点生成了相同的主键导致upsert互相覆盖。所以说幂等不光是同步工具的事目标库的表结构设计也要配合。KFS在迁移检查阶段就会自动扫描目标库的所有表凡是主键或唯一索引缺失的表直接判定为不合格不允许进入增量阶段。4.2 延迟升高时优先做什么同步链路最常出现的告警就是“消费延迟持续上涨”。很多人第一反应是加并行度让Flink任务多开几个并发去消费。但延迟升高的原因往往不在消费端而在写入端。目标库的事务提交卡住了索引维护太慢或者磁盘IO被打满这时候你加再多的并发只会让情况更糟因为所有并发任务都在抢同一个瓶颈资源。我总结了一套排查顺序。第一步看Kafka的lag分布是均匀分布在所有分区还是集中在某一个分区。集中在某一个分区说明那个分区的热点键数据量异常需要排查是不是有单表大事务。第二步看Flink任务的backpressure状态如果source端出现高水位就是下游处理不过来直接去查目标库的监控面板。第三步看目标库的活跃会话如果有大量的锁等待那多半是目标库的索引设计有问题同步任务在不断重试。延迟处理的另一个原则是不要为了追延迟而跳过数据。有人会想把积压的消息先丢弃一部分只同步最新的。这在KFS里是绝对禁止的。丢一条数据可能现在看不出来一个月后对账就会发现那笔账永远找不回来了。正确的做法是暂时降低延迟目标比如从30秒放宽到5分钟先保证数据不丢等目标库性能恢复或者扩容完成之后再逐步追平。4.3 校验引擎设计滑动窗口与分块对比亿级表做全量count非常不现实所以在KFS的增量运行期校验引擎跑的是滑动窗口模式。所谓滑动窗口就是每次只校验最近15分钟内变更过的那些分片而不是每次都全表扫描。怎么知道哪些分片变更过源库的binlog事件里带主键信息增量任务在写入目标库的同时会把变更行的主键hash到一个分片号存到一张变更记录表里。校验任务每隔5分钟扫一次这张表找出最近15分钟内被touch过的分片然后只对这些分片做源端和目标端的checksum对比。这样做有两个好处。第一对源库压力极小。增量高峰每小时可能有几百万行变更但落到分片维度也就几十个分片每个分片10万行以内校验一次的成本非常低。第二发现问题的时效性好。如果同步链路出了故障导致某个分片的数据不一致最多15分钟内就会被滑动窗口扫到而不是等到每天的全量校验才发现。分片校验的SQL也很简单。源端查一条 SELECT COUNT(*) AS cnt, IFNULL(SUM(CRC32(CONCAT_WS(|, id, order_no, amount))),0) AS chk FROM orders WHERE id BETWEEN 100001 AND 200000;目标库执行同样的SQL然后把cnt和chk分别写入校验结果表。调度任务拿两边结果做对比任何一项对不上就触达补偿机制。注意CRC32本身有碰撞概率所以对于校验不通过的分片KFS不会直接判定为数据错误而是进入逐行对比流程定位到具体哪一行不一致再由修复任务去补齐。5. 常见问题与排查实录我踩过的坑和对应解法5.1 增量同步消息堆积严重怎么办有一次线上Kafka的lag涨到了几百万条Flink任务的checkpoint一直失败。查了半天最后发现是目标库有一张表的唯一键冲突。原因是源库那批数据本身有重复的id但源库因为某种历史原因没有建唯一索引数据照样能插进去。到了目标库我们建了唯一索引重复数据就被拒了同步任务一直重试导致位点不前进。这个问题的解法分两步。第一步先跟业务确认重复数据的处理规则如果是脏数据就在同步任务里加一个过滤逻辑遇到重复主键只保留最新一条。第二步目标库不急着建唯一索引等数据校验通过之后再补索引。这里给一个通用建议新增索引这类DDL操作一定不要在同步跑量期间做否则锁表锁死同步任务全堵住。5.2 数据对不上先查这五个地方数据不一致的排查是有套路可循的。我自己遇到的对不上的情况90%以上逃不出这五个原因类型映射错误、事务乱序、过滤条件误伤、目标库触发器二次修改、以及消费位点不对齐。我整理了一个速查表团队新人排查问题的时候照着查就行。排查方向常见现象确认方法类型映射时间字段差8小时、精度丢失抽一行源端和目标端原始值做对比事务乱序目标库最终状态与源库不一致检查Kafka分区键设计同一主键是否落到同一分区过滤条件同步过去的数据比源库少核对同步任务SQL里的where条件目标库触发器写入后被二次修改查看目标库触发器与默认值定义位点不对齐切换后新库缺一部分增量对比Flink checkpoint与Kafka offset记录这五个坑里位点不对齐最坑因为它不会报错只会“少一点数据”而且通常要等到对账才能发现。我们后来干脆把位点记录从内部日志挪到了外部配置中心每次checkpoint成功都主动上报任何一位对不上就停止写入避免带病运行。5.3 切换窗口怎么控制才安全切换是迁移过程中风险最高的动作比搬数据本身危险一百倍。我们的切换流程是一个固定模板先切只读流量比如查询请求验证新库读性能没问题再切写流量从5%开始到20%到50%再到100%。每一步之间至少观察15分钟看报错率、锁等待、慢查询、以及同步延迟是否反弹。切换窗口最忌讳的是“一把梭”。我见过一个同事凌晨三点觉得延迟已经追平了直接把全部流量切过去结果新库的查询撑不住线上接口大面积超时。好在回切脚本提前演练过40秒就切回了源库。所以有一句话我每次都跟团队强调永远不要相信没有演练过的回切脚本。回切要在压测环境至少演练三次而且演练的时候故意制造一些故障——比如模拟目标库挂掉、模拟网络抖动看看回切脚本是不是还能撑住。5.4 类型映射那些“一眼看不出来”的问题类型映射看起来是最基础的问题但恰恰是生产中踩坑最多的。MySQL的tinyint(1)在Flink CDC的默认映射里会被转成boolean如果下游是一个老的Java应用用Integer去接这个字段就会直接抛异常。我们的处理方式是在Flink SQL的建表语句里显式声明字段类型不让它自动推断。还有一个常见问题是时区。MySQL的datetime不带时区timestamp带时区如果源库和目标库的时区配置不一致同步过去的数据就会差8个小时。这属于那种“看起来一切正常业务一查就不对”的典型案例。KFS的做法是在通道层统一把所有时间字段转成UTC存储到目标库写入时再按目标库时区做转换。这样无论源库设在哪个时区中间链路都不会出现歧义。大字段的校验也容易翻车。TEXT和BLOB这种字段如果用普通的方式做checksum读取的时候会消耗大量IO。我们的做法是校验时只针对大字段计算hash值不参与全字段拼接。如果hash对不上再单独抽取这一行逐字节对比。否则每次校验都要把整个大字段拉出来性能根本无法接受。最后再分享一个个人体会。做了几年数据同步我最大的感受是迁移这件事慢一点不可怕可怕的是“慢而不自知”和“快但乱掉”。KFS这套方案通篇的核心并不是什么黑科技它只是把每个环节都变成了可观测、可校验、可追溯的流程。位点有没有对齐写入是不是幂等校验有没有通过全部用数据说话。数据同步这个岗位拼到最后拼的就是对细节的较真程度。每一笔账能不能守住取决于你对位点、对幂等、对校验这套基本功有没有敬畏心。