万亿级数据冷热双写实录:Canal、Kafka 与 Flink CDC 的选型血泪史

在将万亿级交易流水表从单一的单机 MySQL 拆解为“多路异构存储(CQRS 架构)”的过程中,如何实时、零丢数据、有序地将主库产生的每一笔数据变更分流并双写到 Elasticsearch、ClickHouse 和新一代分布式数据库,是数据架构演化中最艰难的攻坚战。
从 2018 年最初的开源组件拼装,到如今支撑全网大促的核心数据总线,我们团队在这条数据同步链路上踩遍了所有可能踩的雷。
回顾这段演进史,从早期的 Canal,到中期的 Canal + Kafka 消息队列,再到最终全面拥抱 Flink CDC,每一个技术栈的替换背后,都是一次深夜 P0 事故复盘换来的血泪教训。
[数据双写与变更捕获 (CDC) 架构演进三部曲]
第一代 (早期):
MySQL Master ──(单点长连接)──▶ Canal Server ──(直接写入)──▶ ES / ClickHouse
(痛点: Canal 单点解析吞吐死死卡在 3万/s, 单机 OOM 全站瘫痪)
第二代 (中期):
MySQL Master ──▶ Canal 集群 ──▶ Kafka (按 OrderID Hash 分区) ──▶ 自研消费集群 ──▶ 目标库
(痛点: Kafka 分区扩容破坏局部顺序性, 消费端无分布式事务导致脏覆写)
第三代 (现代标准):
MySQL Master ──(无锁多线程并发分块读取)──▶ Flink CDC ──(分布式流式计算)──▶ 目标库
(优势: 分布式并发拉取, 原生 Exactly-Once, Schema 变更自适应)
第一代:纯 Canal 直连的“单点猝死”
在业务初期,数据量还在数亿级别,Canal 是最直观的选择——模拟 MySQL Slave 协议向 Master 发送 COM_BINLOG_DUMP,解析出 Row-based 数据变更,直接在 Java 内存中组装并写入目标 Elasticsearch 集群。
血泪崩塌点:
第二代:Canal + Kafka 的“顺序性陷阱”
为了解决单点瓶颈与下游解耦,我们引入了业界标准的“Canal + Kafka + 消费 Worker”架构。Canal 负责解析 Binlog 并投递进 Kafka Topic,下游消费者并发消费并写入 ClickHouse 和新库。
这套架构支撑了我们两年的高速增长,但很快在大促的扩容阶段暴露出一个极其致命的深水区缺陷——消息顺序性被破坏(Message Out-of-Order):
— 业务在 1 毫秒内先后发起的两条更新:
UPDATE t_trade_order SET order_status = 2 WHERE order_id = 8848201; — 状态变为: 已支付
UPDATE t_trade_order SET order_status = 4 WHERE order_id = 8848201; — 状态变为: 申请退款
在 Kafka 中,为了保证同一个 order_id 的所有变更严格按顺序落入同一个 Partition,我们采用了 hash(order_id) % num_partitions 的分区路由策略。然而,在大促前夕,为了应对翻倍的流量,运维团队在未停机的情况下将 Kafka Topic 的 Partition 数量从 64 个动态扩容到了 128 个!
- 扩容前:hash(8848201) % 64 = 12(写入 Partition 12);
- 扩容后:hash(8848201) % 128 = 76(写入 Partition 76)。
同一个订单的两条更新分别落入了不同的 Partition,被下游两个独立的并发消费线程同时拉取。由于网络抖动,第 1 条“已支付”的更新反而比第 2 条“申请退款”更晚写入新库,把最新的终态直接覆盖回了旧状态,引发了严重的资金对账差异!
第三代:全面拥抱 Flink CDC 的终极破局
为了从根本上消除单点瓶颈与分区顺序性隐患,我们在万亿迁移战役中全量切换为了 Flink CDC 3.0 体系:
// Flink CDC 无锁并发读取 MySQL Binlog 的标准 Job 配置
MySqlSource<String> mySqlSource = MySqlSource.<String>builder()
.hostname("10.20.18.10")
.port(3306)
.databaseList("production_trade")
.tableList("production_trade.t_trade_order")
.username("flink_cdc_user")
.password("Secret2026")
.deserializer(new JsonDebeziumDeserializationSchema())
// 开启并行快照扫描:支持将存量万亿大表无锁切分为数百个并行 Chunk 并发拉取
.splitSize(8096)
.distributionFactorUpper(10.0)
.startupOptions(StartupOptions.initial())
.build();
Flink CDC 彻底终结了过往痛点的核心杀手锏包括:
架构演进没有终点。看清每一代工具在极限高压下的物理短板,才能在万亿数据的惊涛骇浪中,筑起最坚韧的数据传输动脉。





