精选·Flink

Flink CDC 的原理与典型用法

91学AI·2026/7/13·9 阅读

考察点

CDC 是实时数仓的入口技术,几乎家家都在用。面试官想确认你理解它和「定时 select 轮询」的本质区别、增量快照算法解决了什么问题,以及实际项目里怎么落地。追问会往全量+增量切换、schema 变更、整库同步方案走。

参考答案

什么是 CDC,为什么用 binlog

CDC(Change Data Capture)是把数据库的变更(insert/update/delete)实时捕获成数据流。两条技术路线:查询型(定时 select where update_time > 上次位点)侵入业务库、抓不到删除、有轮询延迟,只能算玩具;日志型(读 MySQL binlog、Postgres WAL)对业务库零侵入、变更完整、毫秒级延迟,是正经方案。Flink CDC 就是日志型路线:source 端基于 Debezium 的解析能力(2.x 后自研了读取层)把 binlog 解析成 RowData 变更流,下游直接进 Flink 计算或写出。

核心难点:全量 + 增量怎么无缝衔接

CDC 作业启动时表里已有几亿行历史数据,必须先全量快照再切增量 binlog。朴素的「锁表拍快照再读 binlog」会长时间锁住业务库,不可接受。Flink CDC 2.0 引入了无锁快照算法(基于 Netflix DBLog 论文的思路):

把表按主键范围切成若干 chunk,并行读取每个 chunk,读 chunk 前后各记一个 binlog 位点,chunk 读取期间发生的 binlog 变更单独记下来,读完做「增量合并」——把 chunk 数据加上期间的变更,得到该 chunk 在一致位点上的状态。全程不拿全局锁,业务无感知。切完所有 chunk 后从统一的低位点开始持续读 binlog 进入纯增量阶段。

这个算法同时带来三个工程红利:并行快照(大表秒级到分钟级)、断点续传(chunk 粒度进 Checkpoint,失败了重读个别 chunk 而不是整表重来)、无锁。

典型用法

ODS 层实时入湖/入仓:MySqlSource 或 SQL 的 mysql-cdc connector 直接接表,写入 Kafka / Paimon / Iceberg / StarRocks。配合 2PC 或幂等 sink 保证端到端一致。SQL 里声明主键后,CDC 流的 update/delete 会以 changelog 语义(+I/-U/+U/-D)正确传播,下游 upsert 表能处理回撤。

整库同步:CDAS(Create Database As)语法或 DataStream 的整库模式,一个作业同步几百张表,新增表自动纳入,省得每表一个作业。分库分表场景支持正则匹配合并同步。

直连计算:跳过 Kafka,CDC source 直接接窗口聚合输出结果,链路最短,适合轻量场景;严肃数仓还是建议先落 Kafka/湖表分层,便于回溯和复用。

工程注意点

binlog 保留时间必须覆盖作业停机窗口,否则恢复时位点已失效只能重跑快照。大表首次快照对 MySQL 从库有读压力,建议接从库并限速。Schema 变更(DDL)是痛点:binlog 里的结构变化要下游能接住,入湖场景选支持 schema evolution 的表格式(Paimon/Iceberg),直连 Kafka 场景要提前定好变更流程。source 并行度通常就是 chunk 读取并发,增量阶段单 binlog 流是单点的,这也是已知限制。

可能的追问

  • 增量阶段为什么并行度上不去?——binlog 是全局有序的单流,增量读取只能单并发保证顺序;吞吐瓶颈通常不在这,binlog 解析能到几万到十几万行/秒,真扛不住要考虑分库。
  • 作业从 Savepoint 恢复,binlog 位点怎么续?——位点(gtid/offset)存在 Checkpoint/Savepoint 的算子状态里,恢复后从记录的位点继续读,配合 sink 幂等保证一致。
  • 和 Canal + Kafka 的老方案比优势在哪?——少一层 Kafka 中转、全量增量一体、schema 感知、和 Flink 计算天然一体;Canal 方案胜在生态成熟、解耦彻底,两者看团队存量。

评论 (0)

暂无评论,快来抢沙发吧!

91学AI

© 2026 91学AI · 按岗位学 AI 与大数据. All rights reserved.