搞数据同步的人大概都有过这种经历:全量数据好不容易导完,增量延迟也追到了毫秒级,你松了一口气,正准备宣布迁移成功,业务方却丢过来一张对账报表,里面躺着几千条源头有、目标没有的记录。那一刻你就会明白,异构数据同步这件事,延迟只是最容易被看到的指标,真正难的是把每一笔账都守住。
这篇文章想聊的,就是用我们自己基于 Kafka 落地的一套同步框架 KFS(Kafka-based Sync Framework),去处理“不停机迁移”这个场景下的数据一致性问题。所谓不停机迁移,就是业务还在线上跑着,源系统不能锁、不能停,数据要完整搬到目标系统,切流量时还不能丢账、不能重账。KFS 在中间扮演的角色,不是一个简单的搬运工具,而是一套带水位线、带对账、带回放能力的同步机制。它解决的核心问题不是“快不快”,而是“准不准”。
这篇内容适合正在做数据库迁移、实时数仓建设、异构存储同步的工程师,尤其是被“延迟追平了但数据对不上”折磨过的团队。我会把 KFS 的整体设计、核心配置、切换步骤和踩坑经验都拆开讲清楚。
1. 不停机迁移的真正难点:延迟只是冰山一角
1.1 先讲一段真实经历
前两年我们做过一次 MySQL 到 ClickHouse 的迁移,源端是核心订单库,单表将近 3 亿行,每天新增大概 800 万行。业务方的要求很明确:不能停机,每个订单都不能丢,而且切换当天不能让运营感觉到任何异常。
当时我们第一版方案很朴素,用 DataX 导全量,再用 Canal 监听 binlog 把增量灌到目标端。全量跑了大概四个小时,增量也一度追到延迟只有几秒,看起来一切顺利。结果在灰度验证阶段,对账脚本一跑,发现目标端比源端少了六千多条记录,还有三百多条重复。六千多条记录对 3 亿行来说只是十万分之二,但对业务方来说,任何一条核心订单数据出问题都是大事。
后来查了很久,原因五花八门:全量导出期间源端发生了 DDL,导致导出任务部分分片失败被静默跳过;增量阶段目标端写入偶发超时后重试,但重试时走的是另一个没有幂等控制的写入通道;还有一批 update 操作因为 ClickHouse 的 MergeTree 引擎更新语义和我们预期不一致,被直接覆盖掉了。这些问题都有一个共同点:延迟监控完全正常,但数据的一致性在无人察觉的情况下悄悄破了。
从那次之后我形成了一个判断:不停机迁移的核心不是把延迟压到多少毫秒,而是建立一套机制,让每一条数据从源端到目标端的生命周期都可追踪、可校验、可回放。这就是后来 KFS 这套东西的起点。
1.2 异构数据同步到底“异”在哪
很多人理解的异构数据同步,就是两边的表结构不一样,做个字段映射就行。但真实业务里,“异构”带来的麻烦远比字段映射复杂。
首先是数据模型的差异。MySQL 里一张表可以有丰富的主键、唯一索引、自增字段,而 ClickHouse 这类分析型数据库对更新和删除的支持是受限的。你用主键去 upsert 一套逻辑,到目标端可能根本不生效,必须靠 ReplacingMergeTree 或 CollapsingMergeTree 这类引擎自己去处理去重。一旦引擎配置不对,数据就悄悄重了。
其次是类型系统和语义差异。MySQL 的 decimal 到 ClickHouse 的 Decimal 要做精度控制;MySQL 的 datetime 带不带时区,到了 ClickHouse 的 DateTime64 要怎么转换;源端 varchar 里不小心塞了一个特殊字符,JSON 序列化会不会被截断;这些都会被一并算入“异构”的范畴。
再就是事务边界不同。MySQL 的一个事务在 binlog 里可能对应多条变更记录,这些记录之间的业务完整性在源端是原子提交的。但到了迁移管道里,它们被打散成一条条独立消息,如果目标端写入失败一半、成功一半,这个事务的原子性就没了。系统引擎不知道,它只知道单条消息处理成功或失败。
KFS 在设计初期就明确了一点:它不是一个字段对字段的映射工具,而是一条带有状态管理的同步管道。每条消息进入管道之后,它的来源事务、位点、处理结果都必须有记录,这样才能在“异构”的缝隙里找到可对账的依据。
1.3 KFS 解决的是“账”的问题
标题里我说“守住每一笔账”,这里的“账”有两层意思。
一层是业务账,就是订单、流水、余额这类真实数据,一笔都不能少、不能错。另一层是同步账,就是 KFS 内部记录的数据流转账:哪一批数据从源端抽出来了,哪些增量已经消费了,消费到哪个水位了,校验结果是什么,补偿任务执行了几次。业务账是否完整,靠的就是同步账是否可追溯。
这也是 KFS 和普通同步工具最大的区别。普通同步工具是“搬完就算完”,KFS 则是搬完只是开始,它还会持续比对源端和目标端,把差异数据重新拉平,甚至支持在一定时间窗口内做历史回放。用一句话概括:KFS 把同步从“尽力而为”变成了“确认交付”。
2. KFS 核心设计:全量、增量、校验三线合一
2.1 KFS 是什么
KFS 是我们内部基于 Kafka 生态自研的一套同步框架,全称 Kafka-based Sync Framework。它由四个组件组成:KFS Agent 负责从源端捕获数据,包括全量快照和增量变更;KFS Channel 基于 Kafka 的 topic 作为数据通道;KFS Controller 负责分配任务、记录水位线、管理消费进度;KFS Verifier 负责周期性的数据校验和差异统计。
说到底 KFS 没有发明什么特别牛的技术,它更像是把一套迁移方法论固化成系统。Kafka 在这里承担的是缓冲和解耦的职责。为什么选 Kafka 不选其它消息队列?因为我们要面对的是海量 binlog 事件,Kafka 的吞吐量、消息堆积能力和消费位点管理机制都比较成熟。尤其是 Kafka 的消费者可以手动提交位点,这意味着 KFS 可以在数据处理成功之后再提交消费进度,天然适合“确认交付”的模式。
但 Kafka 也带来一个问题:它本身不保证消息不会重复。同一个消息在 rebalance 之后可能被重新消费,这就要求 KFS 在写入目标端时必须有幂等控制。这个后面会专门讲。
2.2 全量链路:导出快照时不能影响业务
不停机迁移最矛盾的地方在于:你要导全部数据,但又不能锁表。MySQL 的 mysqldump 如果不开一致性快照,导到一半数据变了,导出结果就是个逻辑混乱的集合。KFS 的做法是分片并行 + 一致性快照 + 增量补偿三步走。
先讲分片。KFS Controller 会根据主键范围把大表拆成多个分片,每个分片由一个 Agent 并发导出。分片的目的有三个:一是提高导出速度,二是降低单条连接的压力,三是如果某个分片导失败了,重跑这个分片就行,不用全量重来。分片大小怎么定?一般控制在 20 万到 50 万行之间,或者按主键 ID 的 100 万分位一段。太小了任务太多,调度开销大;太大了单任务失败后的重试成本高。我们一般建议每片导出时间不超过 5 分钟。
一致性快照这一步,如果是 MySQL 8.0,建议直接用 GTID + SHOW CREATE TABLE 先拿到一份结构快照,之后所有分片都在同一个 REPEATABLE READ 事务里做 SELECT。这样做的前提是数据库能撑住这个事务的 undo 膨胀,所以一般放在业务低峰期启动全量导出,或者从只读从库拉取快照。
但快照再一致,它也只是一个时间点的状态。从生成快照那一刻开始,业务产生的所有增量变更都还没进目标端。所以全量链路一定会接一条增量补偿链路,把快照位点之后的 binlog 变更补进去。
2.3 增量链路:用水位线代替时间戳
增量同步最常被误解的概念就是“延迟”。大家习惯看“当前时间减去最后一条同步数据的时间”,这个数字确实直观,但它其实不能精确描述同步进度,尤其在源端发生大事务或者 DDL 的时候。
举个例子,假设源端凌晨两点跑了一个大数据量批处理任务,一口气更新了 5000 万行。binlog 事件按顺序进入 Kafka,消费端一直接收到这些事件。如果用“最后一条数据的时间戳”算延迟,因为批处理任务的数据写入时间在凌晨,你会看到延迟是 0,好像已经追平了。但实际上消费端还在处理两个小时前的那批历史数据,业务实时产生的新数据排在后面,根本没有被处理。这就是典型的“假追平”。
KFS 不依赖时间戳判断进度,它用的是水位线机制。Controller 会持续记录两件事:一个是源端 Kafka 分区里最新一条消息的位点,另一个是每个消费者组实际提交到目标端的位点。两者之间的差值才是真正的消费缺口。这个数字和时间没有直接关系,它只代表“还有多少条消息没有被处理完”。
补进来之后,增量同步的延迟监控就会分成三个层次:端到端延迟(业务写入源端到目标端可见)、捕获延迟(源端产生到进入 Kafka)、消费延迟(进入 Kafka 到写入目标端)。KFS 的控制台会同时展示这三个指标,出现异常时可以快速判断瓶颈到底在哪一段。
2.4 校验链路:迁移的“审计日志”
校验不是迁移结束之后才做的事,而是从全量导出阶段就开始持续运行的。KFS Verifier 做的事情是周期性从源端和目标端各取一组数据,按分片或按主键范围做比对,记录差异数量。
校验有全量校验和抽样校验两种。全量校验适合单表数据量在千万级以下或者已经确认关键路径比较稳定的情况;抽样校验适合超大表,可以按主键区间、按时间范围、按业务分片抽出一部分来做。我们的经验是,迁移窗口内做全量校验,迁移完成之后做定期抽样校验,这样既能保证切换安全,也能在后续运行中及时发现增量链路的潜在问题。
校验数据的存储也很讲究。KFS Verifier 会把每一次校验的差异记录写入一张独立的校验审计表,内容包括源端位点、目标端主键、差异类型(源有目标无、源无目标有、字段值不一致),以及发现时间。这张审计表的价值在于,它让数据同步这件事变得可以复盘。业务方问“这条数据为什么不对”的时候,你能直接翻出校验记录,而不是靠嘴解释。
3. 实操:用 KFS 跑通一次 MySQL 到 ClickHouse 的不停机迁移
3.1 环境拓扑与版本
先说我们这次迁移的拓扑。源端是 MySQL 8.0.28,共三套分库,每个分库核心业务表 3 亿行;目标端是 ClickHouse 22.8 集群,三节点;中间是 Kafka 2.8 集群,三个 broker,核心 topic 分区设置为 24 个。KFS Controller 部署在两个节点上,通过选举保证高可用。
为什么选用 ClickHouse 作为目标?因为业务场景是数据分析平台,订单数据进了 ClickHouse 之后要做多维聚合查询。MySQL 保留近 90 天的热数据,90 天以前的数据会被归档到 ClickHouse 里做长期分析。所以这次迁移实际上是把已有历史数据一次性导入 ClickHouse,并保持后续每天的新增数据持续同步。
这个场景对强一致性要求没那么苛刻,因为 ClickHouse 本身不是交易系统,但也不能接受丢数据。业务方期望的是“数据可以稍有延迟,但最终必须完整”。这也是 KFS 最擅长的场景。
3.2 关键配置与参数选择
KFS 用一份 YAML 配置来描述同步任务。下面是一个简化的配置示例,我把关键字段都加了注释:
sync_task: name: order_center_to_clickhouse source: type: mysql hosts: - 10.0.1.10:3306 - 10.0.1.11:3306 - 10.0.1.12:3306 username: kfs_user password: "******" include_tables: - trade_order - order_item binlog: # 全量导出开始时的 binlog 文件名和位点 start_file: mysql-bin.000018 start_pos: 123456789 channel: type: kafka brokers: 10.0.2.10:9092,10.0.2.11:9092,10.0.2.12:9092 topic_prefix: sync_order_center partition_count: 24 # 消息体使用 Avro 编码,带 schema 演进能力 serializer: avro target: type: clickhouse hosts: 10.0.3.10:8123,10.0.3.11:8123,10.0.3.12:8123 database: analytics # 全量导入时的写入并发数 bulk_insert_concurrency: 12 bulk_insert_batch_size: 20000 # 目标表引擎采用 ReplicatedReplacingMergeTree merge_engine: ReplicatedReplacingMergeTree sync_mode: # 三步全开:全量、增量、校验 enable_full_dump: true enable_incremental_sync: true enable_verifier: true # 达到多少秒延迟后触发告警 latency_alert_threshold_s: 30 watermark: # 全量任务达到这个位点之后才开始追增量 full_dump_watermark_topic: sync_watermark sync_interval_ms: 5000这里有两个参数需要重点解释。
第一个是partition_count。Kafka topic 的分区数决定了消费的并行度。24 个分区对应 24 个消费线程,在我们这个规模下比较合适。分区数太少,消费速度跟不上 binlog 产生速度;分区数太多,Kafka 自身和 ClickHouse 的连接数压力会增大。如果你拿不准,可以先按“QPS 峰值 / 单分区消费能力”来粗算,单分区消费能力我们实测在 3000 到 5000 条每秒之间,然后在这个基础上留 30% 到 50% 的冗余。
第二个是bulk_insert_batch_size。ClickHouse 批量写入不是越大越好。批量太大,单次写入占用内存多,失败后重试代价也大;批量太小,ClickHouse 的合并压力会很高。我们的经验值是一批 2 万到 5 万行,具体要根据字段数量和平均行宽调整。
KFS Controller 会定期把当前消费位点写入到 Kafka 的一个专用 topic 里,这个 topic 就是水位线记录的载体。Controller 挂了之后,新的 Controller 启动时会从这个 topic 读取上次的消费进度,继续执行,不会丢位点。
3.3 切换过程:先灰度,再全量切换
不停机迁移的切换,我坚决不建议一次性把流量全部打过去。KFS 的推荐做法是至少分三步走。
第一步,保持源端单写,目标端开启同步。这个阶段全量已经导完,增量持续在跑,KFS Verifier 周期性做全量校验。你把业务查询流量按 1% 的比例灰度切到目标端,对比查询结果和响应时间。如果发现目标端数据有缺失,可以直接回滚,因为源端仍然是唯一写入方,业务没有受到任何影响。
第二步,当灰度查询连续 3 天校验通过率都是 100%,且延迟 P99 低于 5 秒时,可以进入双写阶段。双写就是同一笔业务数据同时写入源端和目标端。这个阶段的关键是写入幂等,目标端的写入必须以业务主键为准做 upsert,否则双写和 binlog 同步同时生效会制造重复数据。
第三步,双写稳定运行一段时间后,把读流量全量切到目标端,源端只保留写入能力。最后再把源端的写流量也切到目标端,源端进入只读状态,等待业务方确认无误后退出。
这里有一个很容易被忽略的细节:双写阶段切换读流量之前,一定要把增量同步的补数任务再完整跑一遍,确保目标端数据已经落后源端不超过 5 秒。否则用户刚切过去看到的可能就是上一分钟的数据。
4. 每一笔账怎么守:幂等、对账与补偿
4.1 目标端的幂等设计
不停机迁移过程中,数据重复几乎是不可避免的。Kafka 的 at-least-once 语义、消费端 rebalance、网络重试、双写与 binlog 同步并存,这些因素任何一个都能导致同一条数据被写入目标端多次。如果目标端没有幂等能力,账就乱了。
ClickHouse 的幂等策略主要靠表引擎。我们用的是 ReplicatedReplacingMergeTree,配合版本字段来实现“同主键保留最新值”。具体做法是在建表时增加一个版本列,KFS 将 binlog 中的事务时间戳或一个全局自增版本号写入这个列。ClickHouse 在后台合并时,遇到相同主键会保留版本号最大的行。
KFS 的 Agent 在写入目标端之前,会先把源端的 binlog 事件按主键做一次本地去重。比如同一个事务里对同一行数据修改了三次,KFS 只需要把最后一次修改的结果写到目标端,这样可以减少目标端的合并压力。如果 binlog 事件跨了事务,KFS 就按“版本号大的覆盖版本号小的”原则处理,最终一致性是由目标端的合并算法保证的。
对 MySQL 这类传统数据库做目标端时,幂等就更直接了,主键冲突就用INSERT ... ON DUPLICATE KEY UPDATE,基本不会出问题。真正容易出问题的是那些没有主键的表,这种表在同步时最好给 KFS 配置一个“附加主键”的规则,比如用多列拼接作为逻辑主键,否则对账和去重都无从谈起。
4.2 对账怎么做才有说服力
对账最怕的就是两边各算各的,最后数字对不上还找不到原因。KFS Verifier 的对账逻辑分三层:数量对账、抽样明细对账、热点值对账。
数量对账最简单,源端执行SELECT count(*) FROM trade_order WHERE update_time >= ?,目标端执行同样的条件统计,比较两个 count。但数量对账有个致命弱点:如果源端和目标端的数据一致但都少了一批,count 可能还是相等的。所以它只能作为快速筛查,不能作为最终结论。
抽样明细对账是真正有说服力的。KFS 会按照主键的散列值抽样,例如抽主键% 100 < 5的数据,分别从源端和目标端查出来,逐字段比对每个字段的哈希值。这个维度能捕捉到“行数相同但内容不同”的隐蔽问题。
热点值对账是针对特定业务场景设计的,比如用户 ID、订单状态、金额分布这类关键维度,单独跑一轮分组统计做比对。我们在实际迁移中发现,数量对账和抽样明细对账都通过了,但热点值对账发现了 18 条数据的目标端金额字段和源端不一致——这是源端 decimal 精度映射到 ClickHouse 时小数位被截断导致的。
4.3 补偿机制:从死信队列到人工订正
KFS 的每条消息处理都有一个状态流转:成功、失败、补偿中、已放弃。处理失败的消息不会简单被丢弃,而是进入死信队列。Controller 会定期扫描死信队列,按照配置的重试策略重新投递。
重试策略我们一般这样设置:前 5 次重试间隔分别为 1 秒、5 秒、30 秒、5 分钟、30 分钟。超过 5 次后消息进入待人工处理队列,同时告警通知到值班群。
最让工程师头疼的一类问题是“重试永远成功不了”。比如目标表结构改了,导致字段长度不够;或者源端字段里有非法字符,到目标端怎么都插不进去。遇到这种情况,靠系统重试是没用的,必须人工介入。KFS 提供了一条手动订正的 API,可以让工程师直接查看死信消息的完整 payload,修改之后重新投递,或者直接跳过并记录原因。
我们的经验是,迁移期间要安排专人盯死信队列,每两小时清一次。死信队列堆积太久,说明链路上有一个结构性错误没有解决,这种时候强推重试只会制造更多的脏数据。
5. 常见问题与排查技巧实录
5.1 问题速查表
先给你一张可以直接抄的排查表,都是我们实践中反复踩过的问题。
| 症状 | 可能原因 | 排查手段 | 解决办法 |
|---|---|---|---|
| 延迟一直追不平 | 全量导出的 binlog 位点不对,增量从旧位点开始消费 | 查看 Controller 记录的消费位点和 Kafka 最新位点 | 重置消费组位点到全量开始位点 |
| Kafka 消费组频繁 rebalance | 单条消息处理时间太久,超过了 max.poll.interval.ms | 查看消费者日志和每分钟消费条数 | 增大单次 poll 批次,提高并发线程数 |
| 对账发现目标端多了数据 | 双写阶段和 binlog 增量同时写入,幂等配置失效 | 查看目标表是否设置了版本列 | 在表引擎中增加版本字段 |
| 对账发现目标端少了数据 | 全量分片任务有失败被静默跳过 | 查看 KFS 任务列表的分片状态 | 补跑失败分片 |
| ClickHouse 查询性能下降 | 大批量 insert 触发了过多 merge | 查看 merge 队列长度和分区数 | 降低写入并发,合并小分区 |
| 源端 DDL 导致同步中断 | 目标端表结构未同步更新 | 查看错误日志中的 schema 不匹配信息 | 手动执行 DDL 变更,然后重放消息 |
5.2 延迟突然飙升怎么排查
延迟飙升是迁移期间最常遇到的“心脏病”级别问题。我自己的排查顺序是这样的。
先看消费端有没有在消费。登录 Kafka 消费者组客户端,执行kafka-consumer-groups --describe --group kfs_order_consumer,看一眼每个分区的 LOG-END-OFFSET 和 CURRENT-OFFSET。如果 CURRENT-OFFSET 很久没动,说明消费端卡住了,去查目标端日志或者看是不是有锁表。
再看是哪个环节卡住。如果消费端活跃但延迟在涨,大概率是目标端写入慢。ClickHouse 这种数据库对批量插入比较友好,但对单条插入很不友好。KFS 默认使用攒批模式,如果单条消息轻量度太高,比如一个事务里有一条大消息夹杂着几千条小消息,写入 batch 会一直等大消息完成,导致整体吞吐下降。这个场景下可以把 batch 等待时间调短一点,让消费端把大消息和小消息分开处理。
还有一种很隐蔽的情况:源端执行了大事务,binlog 产生量在短时间内暴增,Kafka 消息积压是正常的。这时候不要慌,先算一下消费能力和积压量的对比,确认能不能在目标时间内消化完。如果确实消化不完,再考虑扩容分区数和消费者实例。
5.3 校验对不上的几类原因与处理
校验对不上的原因很多,但归纳起来就三类:源端数据变了、目标端写入有问题、校验逻辑本身有漏洞。
源端数据变了一般是因为全量导出和增量同步之间存在时间窗口。比如全量快照是在 10:00:00 生成的,而一条数据在 10:00:05 被更新了,增量链路通过 binlog 捕获到了这条 update,但目标端那条数据还停留在全量快照里的旧值,直到增量消息被处理完才更新。如果校验刚好在 10:00:03 跑了一轮,就会把这个时间差判成差异。这个问题不是 bug,而是数据的时间语义问题。解决方法是校验时给数据加一个“有效时间窗口”,只比较更新时间早于某个时间点的数据,或者容忍一定时间内的差异。
目标端写入有问题通常是字段映射或者类型转换的锅。我们遇到过一个案例:源端字段是 varchar(64),但里面存了一个被截断的 JSON 字符串,目标端用 JSON 解析函数去处理它,结果解析失败,消息被丢进死信队列。这种问题靠重试解决不了,一定要从映射配置和字段内容两方向同时排查。
校验逻辑本身有漏洞往往是最隐蔽的。有一次我们连续三天校验通过率都是 100%,结果切换前一周发现目标端缺了一批数据。后来查出来,那批数据的主键恰好不在抽样规则里——我们当时用的是主键% 100 < 5抽样,但那批历史数据的主键全部是两位数,取模后落在 0 到 4 区间的比例远低于 5%,天然被抽样规则漏掉了。从那以后,我们把抽样规则改成了主键按位异或后的散列值,保证主键分布再奇怪也能均匀采样。
6. 写在这轮迁移之后:几点沉淀下来的经验
这轮迁移做完之后,我最大的感受是:不停机迁移这个事,技术方案的选择固然重要,但更关键的是把“可验证、可回退、可追溯”这三条原则贯彻到每一个环节。KFS 帮我们解决的是工具层面的问题,但在工具之上,团队能不能坚持跑完整轮校验、能不能在延迟异常时冷静定位而不是盲目加资源、能不能在业务方质疑数据准确性时拿出审计记录,这些才是决定迁移成败的底色。
我个人现在的习惯是,任何一条新的同步链路在正式上线前,都会先搭一套模拟数据环境,把全量、增量、校验三条链路完整跑一遍,用脚本模拟源端的 update、delete、大事务、DDL 变更,观察目标端在这些场景下的表现,再把对账脚本固化进发布流程。这样真正面对生产环境的时候,至少能保证“出问题的时候我们知道问题在哪”。
数据同步这个领域没有银弹,延迟追平不是终点,每一笔账都被守住才是。KFS 这套方案目前在我们内部已经跑了好几条链路,后续如果有更深入的对账策略和补偿机制优化,我再来更新。希望这篇拆解能帮你少踩几个坑。