ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

数据同步不只是延迟:KFS如何用Kafka和Flink CDC守住每一笔账

数据同步不只是延迟:KFS如何用Kafka和Flink CDC守住每一笔账 有一回做账务系统的不停机迁移业务负责人每隔几分钟就问一次延迟多少了我理解他的焦虑延迟从两小时追到两分钟谁都想知道离切换还有多远。但说实话在同步链路里跑过一段时间后我的判断标准变了延迟只是血压计真正要盯的是流水号能不能对得上。那次迁移的后期延迟明明已经追到秒级异构数据同步的对账程序却报出十几笔差额其中几笔是 update 乱序导致的目标端最终值回退几笔是目标端唯一键冲突被 sink 跳过但位点继续前进还有一笔来自源端一张无主键的关联表压根没法用幂等键去重。这些都不是“延迟”两个字能反映的问题而是账目一致性问题。这篇文章想聊的就是我用 KFS我们内部对这套基于 Kafka Flink CDC 搭建的同步引擎的简称处理这些问题的过程和思考。1. 不停机迁移的真正难点延迟只是表象一致性和幂等才是关键1.1 为什么“追平延迟”只是中间目标不停机迁移的标准流程是所有做过迁移的人都熟悉的先做一次全量快照导入目标库同时启动增量同步采集源库的 binlog然后等增量追上源库的实时写入延迟降到几秒甚至毫秒级最后把读写流量切到目标库。整个过程里“延迟多少秒”是最直观的进度条因为它直接决定了你还有多久能切换。但这里有个很微妙的误区延迟归零只代表目标端消费完了当前所有已采集的变更事件不代表目标端的数据状态等于源端的数据状态。举一个我在实战中遇到的例子源库一条 update 把订单状态从“待支付”改成“已支付”如果这条事件在目标端执行失败被跳过Kafka lag 依然会掉到 0延迟监控一片绿。只有当你去对账发现目标库那条订单还停在“待支付”问题才会暴露。所以我在项目里一直跟队友强调延迟是健康指标不是完成指标。它告诉你管道通了、消息在流动但它不告诉你消息有没有被正确处理。追平延迟只是“能切换”的必要条件而不是充分条件。1.2 异构数据同步里最隐蔽的三种“账目事故”把延迟当成唯一指标最直接的后果就是忽略下面三类问题。事故类型表面现象根因为什么延迟看不出来丢账对账发现源端有数据目标端没有sink 执行失败被 catch 后跳过或位点提交先于数据落库lag 为 0消费已完成失败事件被吞掉重账目标端出现重复记录消费重放导致同一事件执行多次目标端没有幂等键兜底lag 正常重复数据靠行数和唯一索引才能发现错账目标端与源端最终值不一致同主键的 update 乱序到达旧值覆盖新值或类型转换截断lag 正常且大多数情况下数据“看着没问题”这三类问题量子力学一点说就是“观测即干扰”你越晚去校验越难定位是哪条消息在哪个环节出了问题。所以一个成熟的异构数据同步方案必须把一致性和幂等做进链路设计里而不是等对账时靠人工捞数据。1.3 KFS 到底解决什么问题先解释一下 KFS 是什么。它不是什么开源文件系统也不是某个商业产品的代号而是我们内部对“Kafka-Flink-Sync”这套同步引擎的简称。准确说它是基于 Kafka 做消息管道、基于 Flink CDC 做日志采集与回放计算再叠加自研的对账和补偿模块的一套完整工程化方案。KFS 的核心目标可以概括成一句话在源端和目标端表结构不完全一致的现实约束下做到增量事件可采集、可传输、可回放、可对账、可补偿。它不是重新发明轮子而是把 Canal、Debezium、Flink CDC、Kafka、Flink 这些组件组合成一辆能稳定跑长途的车并且配上仪表盘和备胎——仪表盘是对账备胎是补偿。后面所有章节的讨论都围绕这条链路展开源端日志解析、Kafka 管道、目标端幂等回放以及最后那道“每一笔账都算清楚”的校验线。2. KFS 同步链路是怎么设计的采集、管道、回放、对账四层2.1 为什么选日志采集而不是双写或轮询设计同步链路时第一个要拍板的问题就是数据从哪来。业内常见做法有三种业务双写、定时轮询、日志采集CDC。业务双写是最直观的但代价也最高。它要求所有写源库的代码同时写目标库这在老系统里几乎等于重写数据访问层。更麻烦的是双写一旦中途失败源和目标瞬间就不一致了你得自己处理分布式事务或者接受一段不一致窗口并靠后续补偿修复。定时轮询实现简单但有两个硬伤一是延迟由轮询频率决定做不到准实时二是它很难捕获删除操作除非你要求业务做逻辑删除否则目标库只会越积越多“幽灵数据”。另外高频轮询对源库的压力也不小迁移期间本来就在跑全量读取再加轮询源库容易先顶不住。日志采集是这三种方案里最符合“不停机迁移”诉求的。binlog 或者 redo log 是数据库原生产生的采集端只扮演一个 slave 的角色对源库几乎没有额外压力事件天然包含事务边界和变更顺序删除、更新、插入都能完整拿到。KFS 在源端用 Flink CDC 做解析就是看中这一点。2.2 传输层为什么放一个 Kafka日志解析出来之后下一步是发到 Kafka。有人会问直接从 Flink CDC 到 Flink sink中间夹一个 Kafka 不是多此一举吗起初我也这么想后来被现实教育了。Kafka 在这条链路里至少承担四个职责。第一是削峰填谷源库凌晨跑批时 binlog 量可能是白天的几十倍如果直接打到目标端目标库写入压力会被瞬间拉满Kafka 可以把这些消息暂存起来让 sink 按目标库能承受的速度消费。第二是给“多下游”提供可能同一个变更事件同步端要消费对账端要消费审计端也要消费如果都在 Flink 内部做任务的血缘关系会非常重放到 Kafka 后就变成了各取所需。第三是可回溯Kafka 消息的 retention 只要配置得足够长就能支持按位点重新消费这是后面讲补偿机制的物理基础。第四是解耦源端采集任务和目标任务的生命周期不绑在一起Kafka 宕机最多让消费端积压不会让源端采集中断。Topic 的设计也要讲究。KFS 的默认策略是“按表路由 按主键哈希分区”。一个表对应一个 topic同一行的所有变更事件按主键哈希进同一个分区这样能保证单行事件在分区内有序。分区数一般按表的写入并发来定不建议太小否则消费并行度上不去也不建议盲目设大分区越多Kafka 的元数据和分片管理开销越大。2.3 回放端的核心职责把事件变成目标库里的正确状态回放端是 KFS 里最容易被低估的一层。很多人觉得消费 Kafka 消息然后写目标库不是很简单吗其实回放端的核心职责有两个字状态。它要把一串串变更事件在目标库里还原成源库对应的最终状态。这意味着 sink 不能简单地执行 insert。比如源库对同一行做了三次 updateKafka 里是三条消息中间任何一条执行失败目标库的最终状态都可能不对。Flink 的 checkpoint 机制可以在失败时回滚但回滚后消息会重放如果目标库没有幂等键重放就会造成重复数据。所以回放端在真正执行写库之前需要做三件事确认目标库方言和写入策略根据源库主键或业务唯一键生成幂等键在内存或本地状态里维护最近一段时间的事件序列号用来识别乱序。下面是一个基于 Flink CDC 的源端采集配置示例实际生产里我们会在 deserializer 之后直接接入 Kafka producerMySqlSourceString source MySqlSource.Stringbuilder() .hostname(10.0.12.10) .port(3306) .databaseList(account_db) .tableList(account_db.t_account_flow) .username(cdc_user) .password(******) .serverTimeZone(Asia/Shanghai) .deserializer(new JsonDebeziumDeserializationSchema()) .startupOptions(StartupOptions.latest()) .build(); DataStreamSourceString stream env.addSource(source); stream.addSink(new KafkaSink(kfs-account-flow-topic));从源端到回放端整条链路的顺序是binlog - Flink CDC - Kafka - Flink 消费 - 目标库。 每一层都要有记录每一层都能定位到某一条消息从哪里来、到哪里去后面做对账和排错才不至于大海捞针。3. 守住每一笔账的三个底层机制位点、幂等、补偿3.1 位点管理同步任务重启后不能多跑也不能少跑KFS 的位点管理本质上是把三个“光标”对齐Flink 的 checkpoint、Kafka 消费组的 offset、源端 binlog 的 position。它们之间的关系决定了同步任务在重启后从哪里继续。Flink CDC 本身会记录 binlog positioncheckpoint 会把 position 保存到状态后端。当任务从 checkpoint 恢复时理论上能做到“从哪里停从哪里继续”。但要注意这里有一个所有分布式系统的经典问题at-least-once 和 exactly-once 的取舍。端到端的 exactly-once 在异构同步场景里极其昂贵。源库和目标库不是一个存储系统Flink 没法用两阶段提交同时控制两边强行做只会让吞吐掉一个量级。KFS 的工程取舍是传输和计算层尽量做成 exactly-once 语义但数据库写入层采用“at-least-once 幂等兜底”的组合对外表现等效于 exactly-once。如果你看到某个同步任务把 checkpoint 间隔设成几秒一次重启后 lag 恢复很慢那多半是把 checkpoint 做得太频繁了。我们常用的配置是间隔 30 秒到 1 分钟同时开启 RocksDB 状态后端不然状态一大checkpoint 本身就会成为瓶颈execution.checkpointing.interval: 30s execution.checkpointing.mode: EXACTLY_ONCE state.backend.type: rocksdb state.checkpoints.dir: hdfs:///flink/checkpoints还有一个特别容易被忽略的坑Kafka offset 提交的时机。如果 sink 先提交 offset再写目标库目标库写入失败会直接丢数据如果 sink 先写目标库再提交 offset任务重启后会有重复消费。KFS 的做法是 sink 写库成功之后才提交 offset宁可重放不可丢账因为重放可以由目标端的幂等键兜住丢账就只能靠人力捞了。3.2 幂等回放为什么必须设计同步业务主键很多系统的表虽然有主键但主键在目标库不一定唯一或者源库根本没有主键。这时候就必须给回放端一个“幂等键”也就是目标库在做 upsert 时用来定位记录的键。幂等键的第一候选是业务唯一键比如订单号、流水号、用户 ID 时间戳组合。它不要求与源库物理主键相同只要能精确唯一标识一行业务数据就行。设计时有一个原则永远不要用自增 ID 做跨库幂等键因为你永远不知道两边插入顺序是否一致。目标端的 upsert 语句核心就是“同一把钥匙开两扇门”INSERT INTO target_order ( id, order_no, user_id, amount, status, update_time ) VALUES (?, ?, ?, ?, ?, ?) ON DUPLICATE KEY UPDATE amount VALUES(amount), status VALUES(status), update_time VALUES(update_time);有人会问如果目标库是 PostgreSQL、ClickHouse 或者别的存储怎么办语法不同但思想相同要么用ON CONFLICT要么用REPLACE INTO要么先查后写。真正的难点不是语法而是你是否在每条消息里都带上了足够定位的字段。KFS 在采集端就会做一次“幂等键补齐”如果表没有唯一键就把所有非 null 字段拼起来做 hash如果 hash 仍然有碰撞就把这张表标记为低置信度表进入对账重点名单。3.3 补偿机制对账发现差额之后怎么“找回来”再完善的链路也挡不住实际环境的各种意外。所以 KFS 必须有最后一层兜底补偿。补偿的第一种方式是 Kafka 位点回放。对账任务发现某张表在某个时间段内有多余或缺失的记录时可以定位到那条 Kafka 消息的偏移量从该偏移量开始重新消费把事件再过一遍。这正是前面强调 Kafka retention 要足够长的原因。如果消息已经被清理位点回放这条路就断了。补偿的第二种方式是从源库直接抽取指定主键范围的数据做一次局部全量刷新。这种方式适合补偿那些 Kafka 里已不可寻的旧事件。但注意直接抽源库数据覆盖目标库本质上也是一次写库必须走和正常回放一样的幂等逻辑否则补偿操作本身可能制造新账目差异。我在实际中最常遇到的问题是补偿任务和对账任务之间的时序竞争。对账发现差额补偿任务去补数据但补完的那一刻源库又产生了新的变更结果对账又报出差额。解决这个问题不能靠“等所有任务都停了再校验”而是要把对账窗口设计成带边界的滑动窗口比如只对 5 分钟前到现在的数据进行校验给增量同步留出追平的时间。4. 异构表结构下的坑字段映射、类型转换、DDL 同步4.1 字段映射的第一原则业务字段先对齐物理类型再对齐异构迁移里“异构”两个字往往被低估。很多人以为异构只是“从 MySQL 迁到另一套 MySQL”实际上真正的异构是源端一张订单表有 30 个字段目标端因为新业务设计拆成了 3 张表源端的order_status是varchar(20)目标端改成了tinyint枚举。KFS 的字段映射第一原则是业务语义对齐不是物理类型对齐。也就是说order_no在源库是字符串在目标库也是字符串那它们映射源库有remark目标库觉得不需要这个字段那这条字段就允许丢弃但要记录在映射配置里方便后续审计。我们通常维护一份 JSON 映射配置每个字段明确写出源字段、目标字段、转换规则、是否必传{ table: account_db.t_account_flow, targetTable: ods_account_flow, fields: [ { source: id, target: id, required: true }, { source: account_no, target: account_no, required: true }, { source: amount, target: amount, convert: decimal(12,2) }, { source: create_time, target: create_time, convert: epoch_millis } ], ignoreFields: [internal_remark] }这份配置不只是给人看的它还会被回放端加载决定一条消息在写目标库之前要做哪些 transformation。把映射配置外置是异构同步后期最省心的事因为业务改字段不用动代码改配置就行。4.2 类型转换里容易翻车的几个点类型转换的坑通常不会在你跑 1000 条测试数据时出现而是在你跑完 1 亿条真实数据后突然冒出来。第一个是 decimal 精度。源端amount decimal(10,2)目标端定义成decimal(10,2)看起来没问题但一旦某条金额超过 99999999.99插入就会报错。KFS 的做法是在映射配置里统一把金额类字段放大到目标端能接受的最大精度并且在迁移前扫描一遍源库最大值别等写了 8000 万条才炸。第二个是时间字段的时区。很多系统的源库用datetime默认是东八区目标库如果用timestamp并且连接参数写成 UTC那所有时间都会差 8 小时。这个错账非常隐蔽业务侧不会立刻发现直到对账按天粒度统计时发现每天凌晨的数据都跑到前一天去了。第三个是字符集。源库是utf8mb4目标库是utf8遇到 emoji、生僻字写入直接报错或变问号。这类问题在迁移前很难暴露因为测试数据太“干净”了。建议迁移前先做一个字段值长度和字符集扫描把超长和特殊字符样本提前揪出来。第四个是大字段。一个text字段塞了几百 KB 的 JSON一条 binlog 事件消息可能变得非常大Kafka 单条消息过大会触发 producer 的 max.request.size 限制sink 端内存也会被撑爆。KFS 对这类字段的策略是拆到独立 topic 或者降级为“仅同步元数据 定期全量刷新”总之别让大字段拖垮整条链路。4.3 DDL 同步异构迁移里最容易被忽略的“增量”迁移期间最怕什么怕业务方突然加字段。源库加了一个new_column目标库还没加增量同步的事件里已经带了新字段sink 写目标库时因为字段不存在直接报错任务卡住。如果这个问题发生在凌晨 2 点而值班的人 5 点才看到告警你今晚的“切换窗口”就基本泡汤了。DDL 同步有两种策略。一种是自动解析源库的 DDL转换成目标库方言并自动执行。听起来很美好但实际上 DDL 的语义在不同数据库之间差距很大自动转换很容易出错比如目标库不允许在几亿行的大表上直接加索引。另一种是“半自动 审批”KFS 检测到源库 DDL 后生成差异报告推送给 DBADBA 在目标库手动执行执行完成后恢复同步任务。我自己更推荐后者。异构迁移期间DDL 本来就应该是受控变更业务方不会天天改表结构重点是第一时间发现而不是第一时间自动执行。KFS 在这里做了一件事把 DDL 事件作为一条特殊消息发到 Kafka 里的一个“ddl-topic”对账中心单独监控谁改了表、改了哪个字段、目标库有没有跟上全程留痕。5. 延迟追平阶段最容易“掉账”的场景和完整排查过程5.1 源端大事务拆分后的事件边界问题有一次同步任务跑到一半对账开始报差数。差数集中在一个时间点大概几千条涉及的都是同一张订单明细表。第一反应是 Kafka 丢消息了但查了 topic 的 offset、生产端指标、消费者 lag全都正常。后来把源库 binlog 翻出来发现那个时间点有一个超级大事务一条 update 影响了十几万行。Debezium 解析 binlog 时会把一个事务拆成多条独立事件逐条发到 Kafka。Flink 消费端逐条回放目标库在某个时刻只看到了这个事务的前半部分如果这时候下游有个统计任务查了目标库读到的就是中间状态。这个问题的本质是事务边界在传输过程中被磨平了。KFS 的解决办法是引入事务级缓冲当单条变更事件所属的事务影响行数超过阈值比如 5000 行消费者会把整个事务的事件在本地攒齐排序后再批量写入目标库。代价是延迟会微升但换来的是目标库不会看到“半个事务”账目一致性才有保障。5.2 目标端唯一键冲突导致回放失败但位点还在前进还有一个经典的“假健康”场景。某张对账表的目标库有一个唯一索引但源库历史数据里本来就有两条业务主键完全相同的孤儿数据。正常情况下源库唯一索引能挡住重复但因为是历史脏数据这两条在源库的 binlog 里是两条独立的 insert。sink 写入第一条成功第二条冲突抛异常。如果我们在 sink 里做了“捕获异常并跳过”那么这条消息就被吞了但 Kafka offset 继续提交lag 显示为 0监控一片绿。直到对账按主键逐条比对才发现目标库少了一条。这类问题的排查链路我建议按这个顺序走先看 sink 日志里有没有 “Duplicate entry”、“Conflict”、“Skip” 之类的关键字再看死信队列或跳过计数有没有增长最后再结合对账差额反推是哪个主键范围出了问题。KFS 对这类事件的处理是绝不静默跳过要么写死信队列要么打告警并且回到 KFS 的管理页面能看到“已跳事件明细”而不是只留一个计数。5.3 无主键表在同步中的“不可追踪”问题有些系统里确实存在没有主键的表业务上叫“流水表”只负责往后面追加很少有 update 和 delete。这类表做全量迁移不难难的是增量同步时你没法精确告诉目标库“我要更新的到底是哪一行”。有一次对账发现一张无主键的登录日志表行数完全一致但 checksum 不一致。为什么因为这张表里有几行数据被 update 过源库 binlog 里的 update 附带的是主键信息但表没有主键Debezium 只能把整行都塞进去。目标库按整行替换的方式更新但表里存在重复值替换时匹配到了错误的那一行时间字段就错位了。KFS 的处理分几种情况如果表确实没有任何字段能构成唯一组合那就不做增量 update 同步改成周期全量同步或者只在迁移窗口内允许追加如果表有若干可组合字段就生成一个内部 hash 键作为回放主键。这个决策一定要在迁移前做别等到增量同步跑起来才发现某张表“无药可救”。5.4 时区或精度问题导致的对账永远差一口对账程序用的是按update_time取最近 10 分钟的数据再和源库比对。跑了一整天数据都一致但每到整点附近总会多出或缺少几十条记录。排查了很久发现源库update_time是datetime(3)精度到毫秒目标库对应字段只定义到秒。源库的一条更新如果发生在 14:30:00.800目标库落库后变成 14:30:00对账SQL再按“源库 update_time 14:20:00 AND update_time 14:30:00”这个左闭右开区间去取数就会把那条 14:30:00.800 的记录漏掉或者把它算到下一个区间里导致每次对账都差一口。解决方法是两件事一起做一是把对账的窗口边界加 padding比如上下各扩 1 分钟多取出来的数据按主键精确比对二是把这类时间字段的精度统一KFS 在字段映射阶段就把所有时间字段转成毫秒时间戳目标库不存在字段精度不一致的问题。对账的粒度可以粗但比对的粒度必须细到主键。6. 上线前怎么验证每一笔账对账、灰度切流与演练6.1 对账不是比行数而是比“可解释的差异”上线前的验证阶段很多人做对账就喜欢比行数源库 1 亿行目标库 1 亿行好了没问题。但行数一致能说明的问题非常有限。设想一下一条 update 被错误地应用到了另一行行数不会变一张表里两条记录被互相覆盖行数也不会变。KFS 的对账设计了三个层级。第一层是行数级按表、按分片、按时间窗口对比行数用最快的速度发现大规模丢数据。第二层是 hash 级对每张表的主键和关键业务字段做 hash然后对比源和目标的分片 hash能发现“行数相同但内容不同”的问题。第三层是明细级对重点账户、重点订单做逐行比对把每一笔账的每个字段都拉出来对比。对账结果在 KFS 里会被分成四类“完全一致”、“源有目标无”、“目标有源无”、“值不一致”。前两类基本都是缺失问题走补偿第三类要格外小心可能是源库的删除还没同步过来也可能是目标库多了不该有的数据第四类最常见一般是乱序、类型转换或者漏更新需要结合具体字段定位。6.2 灰度切流的正确顺序只读、双写、切写、回退切换不是一次开关拨过去就完了我见过太多迁移在切流这一步翻车。KFS 配合的切流策略顺序是强制性的先切只读流量到目标库再切一部分写流量做双写最后全部切写并且保留一段时间的回退窗口。只读切流的意义是让业务侧先在目标库上验证查询性能和结果正确性这时候源库还在承担所有写入即使目标库查询有问题也不会污染数据。双写阶段同一笔业务请求同时写源库和目标库源库的 binlog 又会通过 KFS 同步到目标库所以目标库可能被写两次——一次来自业务双写一次来自同步链路。这就非常依赖前面讲的幂等键没有幂等兜底双写阶段目标库分分钟涨出大量重复数据。切写之后不要急着销毁同步任务。回退窗口期内业务一旦发现目标库有问题要切回源库源库必须能拿到目标库这段时间产生的新数据也就是要支持反向同步。KFS 在设计时把正向和反向的采集、回放、对账都做成了可配置的同一套框架回退不只是改个 switch而是把整条链路由“源到目标”换成“目标到源”。6.3 演练清单kill 任务、断网、目标库故障上线之前有一件事比写代码更重要故障演练。KFS 这个方案到底靠不靠谱需要在可控环境下先出一次事故试一试。我每次迁移前都会逼团队做三轮演练。第一轮直接 kill -9 掉 Flink 同步任务看它从 checkpoint 恢复后Kafka 消费位点是否和数据状态一致目标库会不会出现缺失或重复。第二轮把 Kafka 集群里的某个 broker 停掉制造分区 leader 切换观察生产者和消费者的重试逻辑是否正常工作积压会不会暴涨恢复后能不能从正确位点继续消费。第三轮把目标库连接池耗尽或者把目标库磁盘写满看 sink 是无限重试还是进入死信队列告警是否能在 5 分钟内被收到。下面是迁移前检查清单我每次都会过一遍检查项通过标准全量数据导入行数、hash、抽样明细三层校验全通过增量同步位点任务重启后从 checkpoint 恢复无丢账无重账幂等键覆盖所有需要增量同步的表都有明确幂等键对账告警模拟差数能触发告警且能在管理端定位到具体表DDL 变更流程源库执行测试 DDL目标库能收到差异报告回退演练切回源库后新增数据能反向回流无业务报错这套流程走完我才会在切换单上签“可以”两个字。项目收尾的时候我回过头看这套 KFS 方案最大的体会是它真正的价值不是把延迟从分钟压到秒而是在延迟压到秒之后还能让你放心地拍胸脯说“每一笔账都对得上”。位点、幂等、补偿、对账这几个机制单拎出来都不难理解难的是把它们串成一条链路并且在实际迁移的混乱中坚持不偷懒。最后分享一个很多人会忽略的小技巧对账的 SQL 脚本和查询链路一定要独立于同步链路不要用同一个系统的同一个连接池去验证另一个系统否则你等于用一个可能出错的系统去验证另一个也可能出错的地方。数据迁移这件事越到后面越考验耐心你愿意花的校验时间最后都会变成切换时的底气。
返回列表