ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

Flink CDC 3.5.0 从理论到实践 —— 第 2 章 Flink CDC 核心原理深入

Flink CDC 3.5.0 从理论到实践 —— 第 2 章 Flink CDC 核心原理深入 Flink CDC 3.5.0 从理论到实践 —— 第 2 章 Flink CDC 核心原理深入课程定位本系列教程以MySQL 为唯一数据源Sink 覆盖Doris / Paimon / Kafka三大目标从原理到生产落地全链路实战。版本基线Flink CDC 3.5.0 Flink 1.20.x MySQL 8.0/8.4 Doris 4.1 Paimon 1.4.2 Kafka 3.x章节导读2.1 增量快照算法Incremental Snapshot2.2 Schema Evolution 机制2.3 Pipeline 连接器架构2.4 数据类型映射体系2.5 Exactly-Once 落地实现2.6 本章小结2.1 增量快照算法Incremental Snapshot增量快照算法是 Flink CDC 2.0 引入、3.x 持续优化的核心基石它解决了全量初始化阶段既要无锁、又要并发、还要断点续传的难题。理解它就理解了 Flink CDC 区别于 Debezium / Canal 的本质。2.1.1 为什么要增量快照在 Flink CDC 1.x 时代全量初始化采用 Debezium 原生的INITIAL模式存在三大痛点痛点表现影响加锁执行FLUSH TABLES WITH READ LOCK锁库阻塞在线业务读写长事务可能拖垮数据库单线程全量阶段单线程串行读取大表亿级行初始化耗时数小时甚至数天无断点续传全量阶段故障需从头重来大表恢复代价高稳定性差增量快照算法的目标全量阶段无锁、并发、可断点续传且与增量阶段无缝衔接。2.1.2 算法核心思想增量快照算法把全表扫描拆解为多个 Chunk 的并发扫描 Binlog 位点对齐整体流程如下┌─────────────────────────────────────────────────────────────┐ │ 增量快照算法执行流程 │ └─────────────────────────────────────────────────────────────┘ ┌─ 阶段 1: Chunk 切分 ─┐ │ 按主键范围将表切为多个 Chunk │ (Chunk1: [Min, 100], Chunk2: [100, 200], ...) └──────────────────────┘ │ ▼ ┌─ 阶段 2: 并发快照读取 ─────────────────────────────┐ │ 1) 记录当前 Binlog 位点 (highWatermark) │ 2) 多线程并发执行 SELECT ... WHERE pk BETWEEN ? AND ? │ 3) 读取期间不持有任何数据库锁 │ 4) 每个 Chunk 读完写入下游 └──────────────────────────────────────────────────┘ │ ▼ ┌─ 阶段 3: Binlog 补齐 ─────────────────────────────┐ │ 对每个 Chunk, 回放 [lowWatermark, highWatermark] │ 之间的 Binlog 变更, 合并到该 Chunk 的数据上 │ 保证 Chunk 读取期间发生的 UPDATE 不丢失 └──────────────────────────────────────────────────┘ │ ▼ ┌─ 阶段 4: 切换到增量消费 ──────────────────────────┐ │ 所有 Chunk 完成后, 切换为纯 Binlog 消费 │ 从最后一个 Chunk 的 highWatermark 继续 └──────────────────────────────────────────────────┘关键设计点无锁Chunk 读取基于SELECT ... WHERE pk BETWEEN ? AND ?不持有数据库锁并发多个 Chunk 由不同 Subtask 并行处理吞吐随并行度线性扩展位点对齐每个 Chunk 记录读取开始时的 Binlog 位点lowWatermark和结束时的位点highWatermark通过回放这段 Binlog 补齐 Chunk 读取期间发生的变更断点续传Checkpoint 记录已完成 Chunk 的状态故障恢复从断点继续未完成的 Chunk 重读。2.1.3 Chunk 切分策略Chunk 是增量快照的最小工作单元切分策略直接影响并发度与负载均衡。单列主键最常见场景按主键范围切分-- 假设主键为 id (BIGINT), 范围 [1, 1000000], chunk.size 100000-- 切分结果:-- Chunk1: id ∈ [1, 100000]-- Chunk2: id ∈ [100001, 200000]-- ...-- Chunk10: id ∈ [900001, 1000000]每个 Chunk 的边界通过SELECT MAX(pk) FROM (SELECT pk FROM t WHERE pk ? LIMIT chunk.size)动态计算。复合主键默认取主键第一列作为 Chunk Key。若第一列基数低如tenant_id只有几个值会导致切分不均可通过scan.incremental.snapshot.chunk.key-column指定其他非空列。无主键表Flink CDC 3.x 支持无主键表需显式指定chunk.key-column为一个非空列source:type:mysqltables:app.no_pk_tablescan.incremental.snapshot.chunk.key-column:app.no_pk_table:uuid_col注意无主键表无法保证 Exactly-Once 去重建议业务表必须有主键。Chunk 大小调优参数默认调优建议scan.incremental.snapshot.chunk.size8096大表调大如 50000减少 Chunk 数量小表调小避免单 Chunk 过慢scan.snapshot.fetch.size1024单次 SELECT 拉取行数大行可调小小行可调大Chunk 数量经验值建议总 Chunk 数为并行度的 2-4 倍保证负载均衡且有足够并发度。2.1.4 无锁设计的实现传统 CDC 全量初始化加锁的原因需要在快照期间获得一个一致的 Binlog 位点以保证快照与增量无缝衔接。Flink CDC 的无锁设计采用Chunk 级位点对齐替代全局锁不持有全局锁每个 Chunk 独立执行SELECT期间数据库可正常读写Chunk 级位点每个 Chunk 读取前后分别记录 Binlog 位点low/high watermark而非全局快照点Binlog 回放补齐Chunk 读取期间发生的变更通过回放[low, high]区间 Binlog 补齐到快照数据上最终一致性所有 Chunk 完成后整体数据等价于在某 Binlog 位点的一致性快照。对比 Debezium 原生快照维度Debezium INITIALFlink CDC 增量快照加锁FLUSH TABLES WITH READ LOCK无锁并发单线程多线程并发一致性点全局单一快照点Chunk 级位点对齐故障恢复从头开始断点续传2.1.5 全量到增量的无缝切换切换的关键是Binlog 位点对齐保证已读 Chunk 数据 未读 Binlog 变更覆盖全部数据无遗漏无重复。时间轴 ──────────────────────────────────────────► Binlog 位点: L1 L2 L3 L4 │ │ │ │ ▼ ▼ ▼ ▼ ┌─────────────────────────────────────────────────────────┐ │ Chunk1 读取区间 [L1, L2] │ │ 回放 Binlog [L1, L2] 补齐 │ ├─────────────────────────────────────────────────────────┤ │ Chunk2 读取区间 [L2, L3] │ │ 回放 Binlog [L2, L3] 补齐 │ ├─────────────────────────────────────────────────────────┤ │ ... │ ├─────────────────────────────────────────────────────────┤ │ ChunkN 读取区间 [LN-1, LN] │ │ 回放 Binlog [LN-1, LN] 补齐 │ └─────────────────────────────────────────────────────────┘ │ ▼ 切换为纯 Binlog 消费, 从 LN 继续读取增量切换过程所有 Chunk 快照读取完成取所有 Chunk 中最大的 highWatermark作为切换点Source 算子状态切换为 Binlog 消费模式从该位点继续读取后续 Binlog 变更。切换是自动的无需人工干预且切换瞬间下游无感知。2.1.6 断点续传与 Checkpoint增量快照的断点续传依赖 Flink 的Checkpoint 机制快照阶段Checkpoint 记录已完成 Chunk 列表 当前 Binlog 位点故障恢复从最近一次 Checkpoint 恢复未完成的 Chunk 重新读取已完成的跳过增量阶段Checkpoint 记录 Binlog 消费位点恢复后从该位点继续。Checkpoint 1 Checkpoint 2 Checkpoint 3 │ │ │ ▼ ▼ ▼ ┌────────────┐ ┌────────────┐ ┌────────────┐ │ Done: C1,C2│ │ Done: C1-C5│ │ Done: C1-C8│ │ Binlog: L1 │ │ Binlog: L2 │ │ Binlog: L3 │ └────────────┘ └────────────┘ └────────────┘ 故障 ↓ 全量完成 ↓ 恢复: 重读 C3 起 切换 Binlog, 从 L3 消费前提必须开启 Checkpoint且建议EXACTLY_ONCE模式。详见第 10 章。2.2 Schema Evolution 机制Schema Evolution表结构演进是 Flink CDC 3.x 的标志性能力让上游 DDL 自动同步到下游免除了人工维护表结构的负担。2.2.1 问题背景传统 CDC 方案在表结构变更时存在痛点场景传统方案痛点上游加列需手动在下游 ALTER TABLE容易遗忘导致数据丢失上游改类型下游类型不匹配写入报错需停作业改表新增表需手动建表并配置同步运维成本高多表同步每个表单独维护 DDL规模化不可行Schema Evolution 的目标上游 DDL 自动解析、传递、应用到下游。2.2.2 工作链路┌──────────┐ DDL 事件 ┌──────────────┐ 转换 ┌──────────────┐ 执行 ┌──────────┐ │ MySQL │ ───────────► │ Source 解析 │ ────────► │ Pipeline 传递│ ────────►│ Sink │ │ (上游) │ ALTER TABLE │ Binlog DDL │ CDC DDL │ 路由分发 │ DDL 执行│ (下游) │ └──────────┘ └──────────────┘ └──────────────┘ └──────────┘三步流程Source 解析MySQL CDC Source 从 Binlog 解析 DDL 事件ALTER TABLE ADD COLUMN等转为 CDC 内部的 Schema Change 事件Pipeline 传递Schema Change 事件通过 Flink 算子链传递到 SinkSink 执行Sink 连接器根据事件类型执行对应的 DDLALTER TABLE/CREATE TABLE。2.2.3schema-change.enabled参数控制是否启用 Schema Evolutionsource:type:mysqlschema-change.enabled:true# 默认 true取值行为适用场景true默认自动同步 DDL 到下游生产推荐下游表结构自动跟随false不同步 DDL仅同步 DML下游表结构需人工维护或 DDL 需评审注意关闭后上游加列下游表无对应列写入会报错或丢字段。2.2.4 自动建表当tables配置的表在下游不存在时Flink CDC 会根据上游表结构自动推断并创建下游表source:type:mysqltables:app.\.*# app 库所有表sink:type:dorisfenodes:127.0.0.1:8030table.create.properties.replication_num:3# 透传建表属性自动建表流程Source 读取上游表元数据列名、类型、主键、注释根据数据类型映射见 2.4 节转换为下游类型生成CREATE TABLE IF NOT EXISTS ...语句执行表属性可通过table.create.properties.*透传。2.2.5 三种 Sink 的 Schema Evolution 支持差异不同下游对 DDL 的支持能力不同是选型的重要考量DDL 类型DorisPaimonKafka自动建表✅✅N/A无表概念ADD COLUMN✅✅字段变化反映在消息中MODIFY COLUMN改类型⚠️ 部分✅依赖消息格式容错RENAME COLUMN❌✅❌DROP COLUMN❌✅❌TRUNCATE❌❌❌Doris 限制说明Doris 的ALTER TABLE支持加列但改类型、删列受限不支持的 DDL 会跳过并告警不会中断作业。Paimon 优势Paimon 原生支持 Schema 演进加列/改名/改类型/删列均支持是 Schema Evolution 友好度最高的 Sink。Kafka 处理Kafka 无表结构概念字段变化直接反映在消息 payload 中下游消费者需自行容错如ignore-parse-errors配合 Schema RegistryAvro 格式可做到结构化演进。2.2.6 不支持的 DDL 与处理策略不支持的 DDL处理策略删列Doris跳过 日志告警需人工处理改主键不支持需重建表 数据迁移TRUNCATE TABLE不支持需人工 TRUNCATE 下游DROP TABLE不支持下游表保留重命名表部分支持建议保持表名稳定生产建议高风险 DDL改主键、TRUNCATE走变更流程先停作业、人工处理、再恢复。2.3 Pipeline 连接器架构Flink CDC 3.x 引入的Pipeline 架构是走向生产易用化的关键用一段 YAML 替代几十段 SQL。2.3.1 三大模块一个 Pipeline 由Source Transform可选 Sink三部分组成# Source: 数据源 source:type:mysqlname:MySQL Sourcehostname:127.0.0.1port:3306username:adminpassword:passtables:app.\.*server-id:5401-5404# Transform: 数据转换可选transform:-source-table:app.ordersprojection:id, customer_id, amount, statusfilter:status paid AND amount 100description:仅同步已支付且金额大于100的订单# Sink: 数据目标 sink:type:dorisname:Doris Sinkfenodes:127.0.0.1:8030username:rootpassword:passtable.create.properties.replication_num:3# Pipeline: 管道控制 pipeline:name:MySQL to Dorisparallelism:42.3.2 架构流转┌─────────────────────────────────────────────────────────────┐ │ Flink CDC Pipeline 架构 │ └─────────────────────────────────────────────────────────────┘ ┌──────────────┐ ┌──────────────┐ ┌──────────────┐ │ Source │ │ Transform │ │ Sink │ │ (MySQL CDC) │────►│ (ETL 算子) │────►│ (Doris/ │ │ │ │ │ │ Paimon/ │ │ - 全量增量 │ │ - 投影 │ │ Kafka) │ │ - Chunk 切分 │ │ - 计算列 │ │ - 自动建表 │ │ - Binlog 消费│ │ - 过滤 │ │ - 批量写入 │ │ - DDL 解析 │ │ - 路由 │ │ - DDL 同步 │ └──────────────┘ └──────────────┘ └──────────────┘ │ │ │ Flink Runtime │ │ ┌─────────────────────────────────┐ │ └─►│ Checkpoint / 状态后端 / 重启策略 │◄────┘ └─────────────────────────────────┘2.3.3 多 Sink 分发一个 Source 可同时写多个 Sink实现一次采集多端分发route:-source-table:app.orderssink-table:ods_doris.orders# 写 Doris-source-table:app.orderssink-table:ods_paimon.orders# 同时写 Paimon-source-table:app.orderssink-table:kafka.orders_topic# 同时写 Kafkasink:-ref:doris_sink-ref:paimon_sink-ref:kafka_sink详见第 9 章 Data Transformation 与 ETL以及第 12 章综合实战。2.3.4 Pipeline vs 传统 SQL 方式维度传统 SQL 方式2.xPipeline YAML 方式3.x整库同步需写几十段INSERT INTO一段 YAMLtables: app.\.*自动建表不支持需手动建支持类型自动映射DDL 同步不支持支持多表路由难以表达route模块清晰维护成本高低2.4 数据类型映射体系Flink CDC 内部使用统一的CDC Type 体系Source 端从源库类型映射为 CDC TypeSink 端从 CDC Type 映射为目标库类型。理解映射规则是排查类型不匹配报错的关键。2.4.1 CDC Type 体系CDC Type 是 Flink CDC 内部的逻辑类型与具体数据库解耦CDC Type说明TINYINT / SMALLINT / INT / BIGINT整数DECIMAL(p, s)精确小数FLOAT / DOUBLE浮点BOOLEAN布尔DATE日期TIMESTAMP§时间戳不带时区TIMESTAMP_LTZ§时间戳带本地时区CHAR(n) / VARCHAR(n)定长/变长字符串BINARY(n) / VARBINARY(n)定长/变长二进制STRING大文本2.4.2 MySQL → CDC Type 映射MySQL 类型CDC Type备注TINYINT(1)BOOLEANMySQL 布尔约定TINYINTTINYINTSMALLINTSMALLINTINTINTBIGINTBIGINTDECIMAL(p,s)DECIMAL(p,s)FLOATFLOATDOUBLEDOUBLEDATEDATEDATETIMETIMESTAMP(0)TIMESTAMPTIMESTAMP_LTZ(0)带时区CHAR(n)CHAR(n)VARCHAR(n)VARCHAR(n)TEXT / MEDIUMTEXT / LONGTEXTSTRINGBINARY / VARBINARYBINARY / VARBINARYBLOB / LONGBLOBBYTES2.4.3 CDC Type → Doris 映射CDC TypeDoris Type关键点TINYINTTINYINTSMALLINTSMALLINTINTINTBIGINTBIGINTDECIMALDECIMALFLOATFLOATDOUBLEDOUBLEBOOLEANBOOLEANDATEDATETIMESTAMP§DATETIME§TIMESTAMP_LTZ§DATETIME§时区转换CHAR(n)CHAR(n×3)长度 ×3UTF-8 中文 3 字节VARCHAR(n)VARCHAR(n×3)长度 ×3BINARY(n)STRINGSTRINGSTRING重点Doris 用 UTF-8 存储CHAR/VARCHAR 长度会自动 ×3。如 MySQLVARCHAR(50)在 Doris 变成VARCHAR(150)超过 65533 自动转 STRING。2.4.4 CDC Type → Paimon 映射CDC TypePaimon Type备注数值类同名完整保留DATE / TIMESTAMP同名CHAR(n) / VARCHAR(n)VARCHAR(n)Paimon 统一变长BINARY / VARBINARYBYTESSTRINGSTRINGPaimon 类型与 CDC Type 几乎一一对应是映射最直观的 Sink。2.4.5 CDC Type → Kafka 消息格式Kafka Sink 不直接映射类型而是将 CDC 事件序列化为消息格式格式说明适用场景canal-jsonCanal 协议 JSON国内主流下游为 Canal 消费者debezium-jsonDebezium 协议 JSON国际通用maxwell-jsonMaxwell 协议 JSON旧系统兼容avroAvro Schema Registry强类型Schema 演进友好消息示例canal-json{type:INSERT,database:app,table:orders,data:[{id:1,amount:99.9,status:paid}],ts:1713312000}2.4.6 类型不匹配的常见报错报错原因解决Cannot cast DECIMAL to INT上游改类型下游未同步启用 Schema Evolution 或手动 ALTERString length exceedsDoris CHAR/VARCHAR 超 65533自动转 STRING无需处理Unsupported type TIMESTAMP_LTZ下游不支持带时区时间戳转为 TIMESTAMPValue out of range数值溢出检查上游 DECIMAL 精度2.5 Exactly-Once 落地实现Exactly-Once精确一次是流处理最难保证的语义Flink CDC 通过Source 端 Sink 端 Checkpoint三层协同实现。2.5.1 什么不是 Exactly-Once先澄清几个易混淆概念语义含义实现难度At-Most-Once至多一次可能丢数据最低At-Least-Once至少一次可能重复中Exactly-Once精确一次不丢不重最高关键认知Exactly-Once不等于数据只处理一次。底层可能处理多次但最终结果等价于处理一次通过去重/事务保证。2.5.2 Source 端增量快照 CheckpointSource 端的 Exactly-Once 依赖增量快照算法2.1 节快照阶段Checkpoint 记录已完成 Chunk故障恢复跳过已完成 Chunk未完成的重新读取增量阶段Checkpoint 记录 Binlog 位点文件名 偏移量 / GTID故障恢复从该位点继续位点原子性Binlog 位点作为 Checkpoint 状态的一部分与下游 Sink 写入原子对齐。Checkpoint 成功条件: 1. Source 已保存 Binlog 位点 L 2. Sink 已写入 L 之前的所有数据 3. 状态后端持久化完成 三者原子, 任一失败则 Checkpoint 失败, 作业回滚到上一个 Checkpoint2.5.3 Sink 端三种实现路径Sink 端的 Exactly-Once 因目标库能力而异路径一幂等写入Doris Unique 主键 / Paimon 主键表最常见路径依赖目标库主键去重DorisUnique Key 模型相同主键的多次写入会被合并为最新值PaimonPrimary Key Table主键相同的数据自动 Merge机制故障恢复后重放数据主键保证最终结果一致。-- Doris Unique Key 表, 自动去重CREATETABLEorders(order_idBIGINT,amountDECIMAL(10,2),statusVARCHAR(20))UNIQUEKEY(order_id)DISTRIBUTEDBYHASH(order_id)BUCKETS10;路径二两阶段提交Kafka 事务Kafka 0.11 支持事务Flink CDC 可通过事务保证 Exactly-Once预提交Checkpoint 时将消息写入事务未提交正式提交Checkpmaster 协调成功后提交事务消息对下游可见回滚Checkpoint 失败则回滚事务消息丢弃。sink:type:kafkadelivery.guarantee:exactly-once# 启用事务transactional.id.prefix:cdc-tx-# 事务 ID 前缀注意事务模式会增加延迟事务超时时间且下游消费者需设isolation.levelread_committed。路径三At-Least-Once 下游去重部分场景放宽为 At-Least-Once依赖下游去重DorisAggregate 模型非 UniqueKafka普通消息 下游业务去重适用对延迟敏感、可容忍少量重复的场景。2.5.4 三种 Sink 的 Exactly-Once 对照Sink实现路径前置条件风险Doris幂等Unique Key主键模型非 Unique 模型会重复Paimon幂等Primary Key主键表Append 表会重复Kafka事务2PCKafka 事务 下游 read_committed延迟增加2.5.5 Checkpoint 配置要点Exactly-Once 的前提是正确配置 Checkpointenv.enableCheckpointing(60000);// 间隔 1minenv.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);// 精确一次env.getCheckpointConfig().setMinPauseBetweenCheckpoints(30000);// 最小间隔env.getCheckpointConfig().setCheckpointTimeout(600000);// 超时env.getCheckpointConfig().setTolerableCheckpointFailureNumber(3);// 容忍失败env.getCheckpointConfig().setExternalizedCheckpointCleanup(CheckpointConfig.ExternalizedCheckpointCleanup.RETAIN_ON_CANCELLATION);env.getCheckpointConfig().setCheckpointStorage(hdfs:///flink/checkpoints);详见第 10 章作业运维与监控。2.6 本章小结本章深入剖析了 Flink CDC 3.5.0 的五大核心原理回答了无锁、不丢不重是如何实现的增量快照算法Chunk 切分 并发读取 Binlog 位点对齐 断点续传实现无锁全量与无缝增量切换。理解 Chunk 是最小工作单元、schema-change.enabled控制 DDL 同步。Schema EvolutionDDL 事件从 Binlog 解析 → Pipeline 传递 → Sink 执行三种 Sink 支持度不同Paimon 最强Doris 部分支持Kafka 靠消息格式。Pipeline 架构Source Transform Sink 三段式 YAML一段配置完成整库同步是 3.x 区别于 2.x SQL 方式的核心。数据类型映射CDC Type 作为中间桥梁MySQL → CDC Type → Doris/Paimon/Kafka重点注意 Doris CHAR/VARCHAR 长度 ×3。Exactly-OnceSource 端靠增量快照 CheckpointSink 端靠幂等主键Doris/Paimon或事务 2PCKafka三层协同保证不丢不重。核心要点速记增量快照 Chunk 切分 无锁 SELECT Binlog 位点对齐 Checkpoint 断点续传 Schema Evolution Binlog DDL → Pipeline 传递 → Sink 执行 Pipeline Source Transform Sink一段 YAML 替代几十段 SQL 类型映射 MySQL Type → CDC Type → 目标库 TypeDoris 字符串 ×3 Exactly-Once Source 状态 Sink 幂等/事务 Checkpoint 原子对齐下一章预告第 3 章《环境准备与版本依赖》将搭建课程实验环境包括 MySQL Binlog 配置、Flink 集群部署、CDC 驱动包获取、Docker Compose 一键套件为后续实战章节做好环境准备。参考资料Flink CDC 3.5.0 官方文档https://nightlies.apache.org/flink/flink-cdc-docs-release-3.5/Flink Checkpoint 文档https://nightlies.apache.org/flink/flink-docs-release-1.20/docs/dev/datastream/fault-tolerance/checkpointing/Doris Unique 模型https://doris.apache.org/docs/data-table/data-model/uniquePaimon Primary Key Tablehttps://paimon.apache.org/docs/master/concepts/primary-key-table/Kafka Exactly-Oncehttps://kafka.apache.org/documentation/#eos本文是《Flink CDC 3.5.0 从理论到实践》系列教程的第 2 章后续章节将持续更新欢迎关注收藏。
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进