ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

Kafka 异构数据同步实战:KFS 框架守护迁移数据一致性

Kafka 异构数据同步实战:KFS 框架守护迁移数据一致性 做了这么多年的数据迁移我最怕听到的一句话不是“延迟涨到多少秒了”而是“有两批数据对不上了”。异构数据同步这件事讨论度最高的永远是延迟Kafka 消费延迟、目标库回放延迟、端到端延迟。但真正决定一次不停机迁移成败的从来不是延迟多低而是每一笔变更有没有被完整、有序、可追溯地搬到新系统。KFS 就是我们在实践中沉淀下来的一套基于 Kafka 的同步框架专门用来守护迁移过程中这笔“账”。如果你正准备做数据库异构迁移或者被双写、ETL 链路搞得焦头烂额这篇文章应该能帮你少踩几个坑。我会从设计思路说起再把关键机制、实操步骤和排障方法都摊开讲重点是 KFS 怎么在不停机的前提下把账守住。内容不算浅但我会尽量把每个决策背后的原因说清楚希望你读完不只拿到一套能用的方案还能理解为什么有些坑是非踩不可的。1. 先想清楚异构迁移到底在迁移什么1.1 结构不同账不能不同很多人一听到“异构数据同步”第一反应就是“把数据从 A 搬到 B”然后开始纠结同步工具选型。但真正做过迁移的人都知道异 construct 的麻烦不在“搬”而在“映射”。举个最常见的例子源库是 MySQL目标库是 HBase 或者 ClickHouse。MySQL 里一张订单表有 20 个字段目标端可能被设计成宽表字段名变了类型也变了。更复杂一点源端的一张表在目标端被拆成两张表或者反过来源端两张表 JOIN 后写进目标端一张表。这时候同步工具如果只做字段位置对应基本就是在埋雷。KFS 在项目启动前会强制做一轮“映射评审”。不是开发凭感觉写几个 SQL而是把源端每一张表的字段、类型、约束、默认值、字符集全部拉出来和目标端模型逐项比对。遇到无法直接映射的单独列一个“转换规则清单”比如时间戳从字符串转成 datetime状态码从 0/1 映射成 enabled/disabled。这个阶段不能省因为异构同步过程中出现的大多数数据不一致根源都不在工具而在映射定义模糊。还有一个容易被忽略的点唯一键。源端主键在目标端不一定存在比如 MySQL 的联合主键到了 HBase 里可能被拼接成 RowKey。如果目标端没有唯一约束同步时重复写入就很难发现。所以 KFS 的第一条原则是无论目标端是什么存储每一条同步过去的记录都必须有一个业务上可信的“唯一标识”并尽量在目标端建模时落成索引或约束。否则后面做对账、做幂等都是空谈。1.2 KFS 的核心架构与角色分工KFS 不是一个单点工具而是一条完整链路。它由四个角色组成源端捕获器、同步通道、目标端执行器、对账服务。源端捕获器负责读取源库的变更日志比如 MySQL 的 binlog、PostgreSQL 的 WAL把每一次插入、更新、删除转换成统一消息。同步通道就是 Kafka它不负责业务转换只负责把消息稳定地、按序地送到目标端。目标端执行器消费 Kafka 消息执行映射规则写入目标存储同时记录同步位点。对账服务则是 KFS 的账本定时比对齐两边的数据发现问题立刻触发告警和重放。这套分工最大的好处是解耦。源端不需要知道目标端长什么样目标端也不用关心源端怎么捕获变更。Kafka 夹在中间天然可以作为缓冲。源库如果突然来了一波大事务Kafka 可以先扛住目标端按自己的节奏消费不至于把两边都拖垮。相比双写在业务代码里塞一段“同时写源库和目标库”的逻辑KFS 不侵入业务系统。业务系统只负责自己的读写同步完全在数据链路层面完成。这样迁移期间业务代码一行都不用改风险也大幅降低。选 Kafka 还有一个实际原因消息可回放。Kafka 的消息不会消费完就删除而是根据保留策略留存一段时间。一旦目标端出现数据不一致我们可以从某个历史位点重新消费不需要让业务系统配合补数据。这个能力在不停机迁移中几乎是刚需。2. KFS 怎么把每一笔账记清楚2.1 分区与顺序账本不能乱Kafka 本身只保证分区内的消息有序不保证全局有序。KFS 对此并不纠结因为业务层面我们只需要保证“同一个实体的变更顺序不乱”。设计分区键时我会优先选业务主键或者自然键。最典型的例子是订单表按order_id做分区键同一笔订单的所有变更都会进同一个分区消费端看到的是按时间排好的顺序。这样即使在目标端落库时并发执行同一订单的更新也不会前后颠倒。但分区键选不好会出大问题。有一个项目用的是默认轮询分区策略结果订单的创建消息和支付消息被分到不同分区消费端并发处理后先更新了支付状态再插入订单记录目标端直接报主键冲突。后来我们把分区键改成order_id问题立刻消失。另外要注意热点分区。如果用user_id做分区键头部用户的大促订单可能把所有消息都压在同一个分区其他分区空闲。KFS 的解决方法是允许配置“多级分区键”比如先按order_id的哈希值粗分再把疑似热点 ID 单独抽出走独立分区。实际效果不错但需要根据业务数据分布提前测算不能上线后才发现倾斜。2.2 位点、幂等与提交顺序“每笔账”这三个字落到技术上核心是“位点”和“幂等”。KFS 在源端捕获变更时会把源库的日志位点比如 binlog 文件名和 position或自增序号写进消息头。目标端执行器每消费一批消息经过转换写入目标库后才会把这一批的位点提交给 Kafka。这个设计保证了一个基本原则消息永远不丢。哪怕执行器在写入目标库后突然宕机Kafka 知道上次提交的位点重启后从该位点继续消费。但“不丢”还不够因为还可能重复。消费端写入成功后如果还没来得及提交位点就挂了重启后会再次消费同样的消息。所以目标端写入必须幂等同一笔消息重复执行和只执行一次最终结果必须一样。具体做法不算复杂。消息体里带一个全局唯一的sync_id目标表加一列sync_id写入时用INSERT ... ON DUPLICATE KEY UPDATE或者MERGE语句。如果sync_id已经存在就认为这条消息已经处理过直接跳过或者只更新必要字段。这里要注意sync_id对应的索引必须存在否则每条消息都要全表扫描性能直接崩掉。KFS 默认不采用“先提交位点再写目标库”的策略因为那样一旦写入失败消息就永久丢失了。我们的取舍是宁可让目标库重复执行也不能让源端的更新凭空消失。毕竟重复数据可以通过sync_id去重丢数据要找回就麻烦得多。2.3 每条变更都带上下文为了让对账变得简单KFS 的消息体不是简单的“字段值集合”而是完整的变更上下文。一条消息大致长这样{ sync_id: 8f3a2f1e-9c3b-4d7b-b0a2-1c9d0e2a5f6b, source_offset: mysql-bin.000023:45678901, op: UPDATE, table: orders, key: {order_id: 10086}, before: {status: PAID, amount: 99.00}, after: {status: SHIPPED, amount: 99.00}, ts: 1710000000000 }before和after不一定都要保留但对账时非常有用。比如目标端只记录最新状态一旦发现不一致我们可以直接从消息里看出这条记录之前是什么样、现在是什么样不用再去源库翻历史。source_offset是源端的日志位点配合sync_id相当于给每一笔变更都盖了一个身份戳。后续做数据对账不需要全表比对只需要按时间区间拉取源端和目标端的sync_id集合找出差异区间再精确定位到某几条消息。很多人做数据同步只关心“最新状态”忽略历史变化。但不停机迁移的复杂性在于目标端和源端可能在很长一段时间内并行运行业务随时可能回滚或切换。没有完整变更上下文出问题之后很难复盘。3. 实操用 KFS 完成一次不停机迁移3.1 迁移前评估摸清家底开始搭链路之前我习惯先做一份“迁移清单”内容至少包括三件事对象清单、映射关系、特殊规则。对象清单就是把源端所有需要迁移的表/集合全部列出来。很多人只列业务大表忽略字典表、配置表、临时表结果迁移后业务跑着跑着发现某些基础数据缺失又得回头补。KFS 的处理方式是从源库元数据里自动拉取表清单再由 DBA 人工确认不依赖口头沟通。映射关系要细到字段级。KFS 提供一个映射配置文件支持表达式转换比如把字符串拼接、数组展开、字段重命名。每个映射规则都要有人签字确认尤其是涉及金额、状态、时间这三类字段任何误解都可能造成严重的数据错误。特殊规则指的是增量同步之外的处理逻辑。例如源端有一张表只是临时状态表业务上允许清空重建那就不需要走逐条同步直接快照覆盖。还有一些表存在逻辑删除同步到目标端时要过滤掉。这些规则提前写清楚比在同步过程中临时加逻辑要安全得多。迁移前还要生成基线快照。KFS 的做法是选择一个业务低峰期先用FLUSH TABLES WITH READ LOCK或者等价手段获取一致性的快照点同时记录源库当前位点。快照导出可以并行做但必须保证导出开始时所有参与的表都在同一个一致位点。之后增量同步从这个位点开始快照数据加上增量数据才能完整还原出源库的全貌。这个步骤如果漏了后面会面临“快照数据和增量数据重叠”或“快照数据和增量数据之间有空洞”的问题对账时极其痛苦。3.2 搭建同步通道从 CDC 到目标端假设源端是 MySQL目标端是 ClickHouse一条最简 KFS 链路的搭建分四步。第一步在 Kafka 创建 topic。分区数建议略大于目标端写入并发度但不能太多。我一般会配置为“目标端并发写入线程数 x 2”让每个消费线程都有独立分区同时留一点余量给故障转移。副本数按 Kafka 集群标准来至少 2 副本。第二步启动源端捕获器。捕获器连接 MySQL开启 binlog解析出变更消息后发送到 Kafka。为了保证同步不丢必须给捕获器加一个本地持久化队列作为缓冲。MySQL 的 binlog 如果因为网络抖动暂时发不出去捕获器不能直接丢弃消息。第三步启动目标端执行器。执行器消费 Kafka 消息按照映射配置转换数据然后以批量方式写入目标表。ClickHouse 这类系统对批量写入很敏感单次写入行数太少会严重影响性能。KFS 默认设定批量大小是 2000 行或者 2 秒攒一批这两个条件谁先到就先刷一批。配置文件大致长这样kfs: source: type: mysql host: 10.0.0.1 database: shop binlog: position: mysql-bin.000023:45678901 target: type: clickhouse table: ods_orders batch_size: 2000 flush_interval_ms: 2000 sync: partition_key: order_id idempotent: column: sync_id unique_index: uk_sync_id retry: max_attempts: 3 backoff_ms: 1000第四步启动对账服务。对账服务先做一次“存量核对”也就是把源端快照和目标端已有数据做比对确认基线没问题后再开始周期性的增量核对。增量核对不要求每次全量比对只需对比最新位点区间内的sync_id集合效率高很多。搭建完成后看两个核心指标端到端延迟和目标端写入耗时。KFS 通常会将每条消息的写入耗时记录到监控系统如果写入耗时持续上涨说明目标端出现瓶颈需要调整批量大小或增加并发。3.3 延迟控制不是“越快越好”不停机迁移期间大家都希望同步越快越好恨不得延迟压到 0。但真实场景里一味追延迟往往适得其反。目标端写入能力是有限的。Kafka 消费速度如果远大于目标库落盘速度消息会在执行器本地堆积导致内存压力增大GC 频繁最终写入耗时越来越高。KFS 专门实现了一个“滑动窗口限流器”思路很简单统计过去 5 秒的平均端到端延迟和目标端写入耗时。如果平均延迟低于设定目标值就放开消费速度如果延迟开始抬头就自动降低消费并发或拉长批量间隔。这里要说明一点KFS 不会让目标端执行器无限减速。滑动窗口是限流不是停机。比如某些团队确实遇到过“故意把消费延迟 30 分钟”的需求为了错开业务高峰。这不是不行但有两个前提一是 Kafka 的retention.ms必须大于“最大允许延迟 预留缓冲”否则消息过期被清掉同步就断了二是消费者要做心跳配置如果暂停消费时间超过了max.poll.interval.msKafka 会认为消费者已死触发 rebalance。应对办法是调大这个参数或者用 Kafka 的 pause/resume 机制来代替停止消费。延迟控制真正的目标是“不让源端和目标端的差距持续扩大”。我习惯给 KFS 配置一个动态阈值如果业务高峰期端到端延迟短期到 10 秒可以接受但如果延迟在低峰期还没有回落说明链路里存在瓶颈需要人工介入排查。3.4 切换、校验与回滚同步链路跑了一段时间应用层可以开始做切换了。但“切换”不是把域名或流量指向新系统就完事而是有三个前置条件。第一存量数据核对通过。这个核对不只是条数一致还要抽样比对关键字段。KFS 对账服务会输出差异列表每一行差异都要有明确解释要么是映射规则导致要么是源端本身数据有问题。第二增量位点追平。把目标端执行器暂停写入让 Kafka 消费位点追上源端最新位点。这时源端的写操作需要短暂暂停或者在业务侧开启只读。我一般选择低峰期做 10 到 30 秒的只读窗口业务影响很小。第三记录切换位点。切换后KFS 继续监听源端但不再写入目标端而是把增量消息全部保留在 Kafka 中。这就是天然的“回滚缓冲”。如果新系统出现问题需要切回源系统只需要把 Kafka 里积压的消息重新放给目标端执行器就能保证两边数据不丢。回滚演练一定要在正式切换前做一次。我在真实项目里见过一个尴尬场景切换后在回滚时发现 Kafka 的副本磁盘爆了消息全部丢失。所以回滚方案不是“留着 KFS 就行”而是要确认同步链路在回滚期间依然健康磁盘、内存、带宽都要有余量。4. 常见问题与排查技巧实录4.1 消费延迟突然飙升怎么办KFS 用久了最常见的告警就是“consumer lag 持续升高”。延迟攀升的原因通常不是 Kafka 本身而是目标端或者源端出了问题。我列一个快速排查表遇到延迟飙升时按顺序看现象可能原因快速处理方式目标端写入耗时上涨目标库存在慢 SQL、锁等待登录目标库查慢查询暂停批量写入等锁释放Kafka 消费线程数小于分区数消费并行度不够增加消费线程或者减少分区数源端捕获器堆积binlog 解析慢或本地队列阻塞查看捕获器 CPU 和 I/O扩大缓冲队列目标端 GC 频繁写入数据量超出内存承受范围调小批量大小降低单批内存占用升级堆内存网络抖动Kafka broker 与执行器之间的带宽不足检查监控必要时临时扩容带宽最常见的是目标库首次批量写入触发了大量索引更新导致锁竞争。有一次我把批量大小从 2000 调到 5000结果 ClickHouse 的 merge 线程突然被打满写入耗时从 5ms 飙升到 800ms。后来我把批量大小调回 2000 并用滑动窗口限流压住写入速率延迟就稳定了。不要一看到消费延迟高就加分区或加线程。先看清楚延迟是发生在“Kafka 到执行器”这段还是“执行器到目标库”这段。KFS 监控里能直接看到两个指标消费位点差和写入耗时。如果消费位点差不大但写入耗时很高问题显然在目标端如果消费位点差一直在涨说明消费能力不足。4.2 对不上账如何定位和修复数据对账发现差异时我的原则是“不要急着重新同步全量”。全量重导不仅慢还可能覆盖目标端已经修正过的新数据。先用对账服务按sync_id找出差异集合。例如源端最新 10 分钟同步过来的记录有 1200 条目标端只查到了 1198 条缺了两条。这时去 Kafka 里查对应时间段的原始消息看这两条的op是什么。如果一个是DELETE目标端已经删掉了那很可能不是缺数而是对账逻辑没有把“已删除”状态算进去。如果确认是漏写最简单的修复是从 Kafka 按位点重新消费那几条消息。KFS 支持“定向重放”你可以指定一个 topic、一个分区、一个时间范围把消息重新发给执行器。因为写入是幂等的重复执行不会产生副作用。还有一个定位技巧对比目标表的MAX(sync_id)和源端最新位点。如果位点差小于一个很小的阈值说明两侧基本一致差异可能来自映射规则。比如源端字段是字符串目标端定义成整型同步时发生了隐式转换导致某些值被截断。这种问题用 SQL 很难查出来反而是查映射配置更快。4.3 大事务和 DDL 带来的一堆坑不停机迁移里最怕两件事大事务和 DDL。大事务意味着源端一个事务里更新了几百万行CDC 捕获器会产生几百万条消息。Kafka 本身扛得住但目标端执行器如果按“攒 2000 条写一次”的默认逻辑会把一个小事务拆成很多批中间一旦有某几批失败执行器需要处理部分成功。KFS 针对这个问题做了“事务组标记”捕获器在消息头里写入tx_id执行器遇到同一个tx_id的消息时会先攒齐再一次性写入或者记录事务边界批量重放时能知道哪些消息属于同一笔事务。如果你用的同步工具没有这个特性建议至少给消息加事务 ID方便回滚和定位。DDL 更麻烦。源端做了一次ALTER TABLE ADD COLUMN目标端如果还没准备好执行器写入时就会报字段不存在。KFS 的默认策略是检测到 DDL 消息时暂停该分区的消费并发出告警由 DBA 在目标端执行对应 DDL 后手动恢复。自动化执行 DDL 太危险我不建议在生产环境开启自动改表。另一个常见问题是源端删了一个字段但目标端历史数据中还保留着。CDC 消息不会包含被删字段的历史值对账时容易误报差异。遇到这种情况需要把映射规则里的“忽略字段”配置好让对账服务跳过这些不参与比对的列。4.4 几条独家心得最后分享几个我在实践中慢慢养成的习惯不一定写进文档但很管用。第一给对账服务单独开一个“标记 topic”。主同步链路负责搬数据对账服务把每次比对的差异结果、重放记录、修复状态都发到这个独立 topic 里。这样后续排查时有完整的审计日志不会被主链路的健康检查噪音淹没。第二监控指标要拆细。不要只看一个“端到端延迟”至少要拆成捕获器产生消息的延迟、Kafka 到执行器的积压量、执行器单批写入耗时、目标端最近一次写入位点。每个指标对应一条链路的某个环节哪一段出问题一目了然。第三迁移上线前做一次“反向压测”。不只是测目标端能扛多少写入更要测“当源端突然产生大事务时Kafka 到目标端的回放能力还能不能跟上”。有一次压测模拟了平时 20 倍的写流量目标端直接 OOM这让我意识到同步链路的容量规划要按峰值算不能按均值算。结合我自己的经验迁移结束不代表同步链路马上要拆。我会建议至少保留两到三个业务周期让对账脚本持续跑着同时把 Kafka 的消息留存时间调长一些。这比任何一次切换演练都让人安心。毕竟延迟数字是给领导看的账对得上才是给自己兜底的。
返回列表