欢迎光临
我们一直在努力

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

万亿级数据冷热双写实录: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 集群。

血泪崩塌点:

  • 单实例 CPU 解析瓶颈:Canal 的 Binlog 协议反序列化与事件过滤是单线程模型的。在遇到大促突发写入洪峰(主库 TPS 达到 8 万/s)时,单个 Canal 进程的 CPU 核心直接被吃满,Binlog 同步延迟迅速从毫秒级拉长到数小时;
  • 缺乏背压(Backpressure)缓冲:一旦下游的 ES 集群发生 GC 停顿或写入限流,Canal 内存队列在几十秒内被打爆,直接引发 JVM OOM 崩溃。
  • 第二代: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 彻底终结了过往痛点的核心杀手锏包括:

  • 并发无锁快照扫描(Chunked Snapshotting):以往迁移存量数据必须锁表或单线程读。Flink CDC 利用自适应主键分片算法,在不加任何全局锁的前提下,允许 64 个 Worker 并行扫描源库的历史存量数据,存量迁移速度提升 20 倍;
  • 两阶段提交与精确一次(Exactly-Once)语义:借助 Flink 原生的 Chandy-Lamport 分布式快照(Checkpoint)算法,将数据库 Binlog 位点与目标端(ClickHouse / 分布式库)的写入状态绑定在同一个分布式事务中,即使节点中途宕机,恢复后既不重发也不丢数据;
  • Schema 动态演进自动同步(Schema Evolution):主库执行 ALTER TABLE ADD COLUMN 时,Flink CDC 能够自动感知 DDL 变更事件并在下游异构存储中自动同步变更结构,彻底终结了过去因加字段导致同步任务大面积报错中断的运维噩梦。
  • 架构演进没有终点。看清每一代工具在极限高压下的物理短板,才能在万亿数据的惊涛骇浪中,筑起最坚韧的数据传输动脉。

    赞(0)
    未经允许不得转载:171主机测评 » 万亿级数据冷热双写实录:Canal、Kafka 与 Flink CDC 的选型血泪史
    分享到: 更多 (0)

    评论 抢沙发

    • 昵称 (必填)
    • 邮箱 (必填)
    • 网址