简介这份文档由尚硅谷研究院编写面向具备一定Flink基础的大数据研发人员聚焦实时数据同步场景系统讲解Flink CDC 3.0从MySQL捕获变更数据并同步至Doris的完整流程。内容从CDC概念切入对比基于查询与基于Binlog两种方式的差异重点说明后者低延迟、可捕获全量变更且不增加数据库压力的优势并介绍flink-cdc-connectors组件的使用方式。案例部分涵盖数据准备、DataStream与Flink SQL两种实现路径以及MySQL到Doris的Streaming ETL落地涉及环境搭建、配置文件编写、任务启动与测试等环节同时强调开启Binlog、设置检查点等关键配置对任务稳定性的影响。资源为1个docx文件压缩包约145KB结构紧凑便于查阅。目前已有1788人学习适合希望掌握CDC原理并独立完成同步任务部署的工程师参考实践。1. 从 Binlog 到 Doris为什么 Flink CDC 3.0 值得你花一个下午跑通如果你正在做数仓实时同步大概率遇到过这种局面MySQL 里数据一直在变下游 Doris 却只能靠 T1 批量导数业务方要看“刚刚下单的用户画像”你只能尴尬地说“明天早上才有”。Flink CDC 3.0 就是冲着这个场景来的——它把 MySQL 的 Binlog 变更实时捕获出来经过 Flink 做轻量 ETL直接落到 Doris整条链路可以做到秒级延迟。和早期版本最大的不同是3.0 引入了 Pipeline 连接器体系MySQL 到 Doris 的同步不再需要你手写 Java DataStream 代码一个 YAML 配置文件加一条flink-cdc.sh命令就能跑起来。这份尚硅谷的教程覆盖了 DataStream API、Flink SQL、Pipeline YAML 三种用法从环境准备到断点续传都有完整示例。适合有 Flink 基础、正在选型实时同步方案的工程师也适合想从 CanalKafka 那套架构迁移过来的团队做技术验证。2. CDC 选型与 Flink CDC 3.0 的架构变化为什么不是 Canal 加 Kafka2.1 基于查询和基于 Binlog 的本质差异在动手之前先把选型逻辑理清楚。CDC 这件事市面上就两条路基于查询的Sqoop、DataX和基于 Binlog 的Canal、Maxwell、Debezium、Flink CDC。基于查询的方式本质是定时跑SELECT靠时间戳或者自增 ID 捞增量问题是它抓不到DELETE一条记录被删了你查询侧根本看不见而且高频轮询对 MySQL 本身有压力。基于 Binlog 的方式是伪装成 MySQL 的从库订阅 Binlog 事件流插入、更新、删除全都能捕获延迟可以压到毫秒级对源库的压力也小得多。Flink CDC 属于后者底层用的是 Debezium 做 Binlog 解析但它在 Debezium 之上做了两件关键的事一是把 Binlog 读取位置信息以 Flink Checkpoint 状态的方式管理起来天然支持断点续传和 Exactly-Once二是提供了 Flink SQL Connector你可以像查普通表一样写SELECT * FROM mysql_tableFlink 自动帮你处理全量和增量的切换。2.2 Flink CDC 3.0 的 Pipeline 架构解决了什么问题2.x 时代你要做 MySQL 到 Doris 的同步得自己写 DataStream 代码建 Source、做数据转换、建 Sink、处理 Checkpoint 配置、管理重启策略。代码量不小而且每换一个源表或者目标表就得改代码重新打包。3.0 引入了 Pipeline 的概念把 Source 和 Sink 抽象成可插拔的连接器中间的数据转换用 YAML 描述。你只需要在配置文件里声明“源是 MySQL 的 test 库目标是 Doris 的 test 库”然后执行bin/flink-cdc.sh job/mysql-to-doris.yamlFlink CDC 会自动帮你生成 Flink 作业并提交到集群。这个变化带来的实际好处是同步任务从“写代码”变成了“写配置”运维成本大幅下降。但代价是灵活性——如果你需要复杂的字段映射或者自定义 UDFPipeline 模式目前支持有限还是得回到 DataStream 或 Flink SQL 方式。所以我的建议是标准的分库分表同步、整库同步用 Pipeline YAML需要复杂 ETL 逻辑的用 Flink SQL 或 DataStream。2.3 环境准备MySQL Binlog 和依赖包一个都不能少不管用哪种方式MySQL 侧的 Binlog 必须开。教程里给的配置是binlog_formatrow这是硬性要求因为 Debezium 需要解析行级别的变更。另外binlog-do-db要指定你需要捕获的库别图省事开全库Binlog 膨胀起来磁盘扛不住。# /etc/my.cnf 中添加 server-id 1 log-binmysql-bin binlog_formatrow binlog-do-dbtest binlog-do-dbtest_route改完配置必须重启 MySQL然后登录确认SHOW VARIABLES LIKE binlog_format返回ROW。这一步翻车的概率不低——很多人改了my.cnf但没重启或者 MySQL 8.0 的配置文件路径不是/etc/my.cnf而是/etc/mysql/mysql.conf.d/mysqld.cnf导致配置根本没生效。Flink 侧需要准备的是 CDC 连接器 JAR 包。DataStream 方式通过 Maven 依赖引入flink-connector-mysql-cdcPipeline 方式则是把flink-cdc-pipeline-connector-mysql-3.0.0.jar和flink-cdc-pipeline-connector-doris-3.0.0.jar放到 Flink CDC 安装目录的lib/下。注意版本要对齐Flink 1.18.0 配 CDC 3.0.0MySQL 驱动用 8.0.31Doris 连接器版本要和 Doris 集群版本兼容。提示Pipeline 模式下 Flink CDC 自带一个轻量级 Flink 运行时但生产环境建议还是提交到独立的 Flink 集群方便统一管理资源和监控。3. DataStream 与 Flink SQL 两种方式实操从建表到断点续传3.1 DataStream 方式Checkpoint 配置是核心DataStream 方式适合需要精细控制数据流处理的场景。教程里的代码结构很清晰准备环境、开 Checkpoint、建 Source、读数据、打印、执行。我重点说 Checkpoint 这块因为这是最容易出问题的地方。// 开启 Checkpoint间隔 3 秒Exactly-Once 语义 env.enableCheckpointing(3000L, CheckpointingMode.EXACTLY_ONCE); // 超时 1 分钟 env.getCheckpointConfig().setCheckpointTimeout(60 * 1000L); // 两次 Checkpoint 最小间隔 3 秒 env.getCheckpointConfig().setMinPauseBetweenCheckpoints(3000L); // 任务取消时保留 Checkpoint env.getCheckpointConfig().enableExternalizedCheckpoints( CheckpointConfig.ExternalizedCheckpointCleanup.RETAIN_ON_CANCELLATION); // 重启策略1 分钟内失败 3 次则停止 env.setRestartStrategy(RestartStrategies.failureRateRestart( 3, Time.days(1L), Time.minutes(1L))); // 状态后端和 Checkpoint 存储路径 env.setStateBackend(new HashMapStateBackend()); env.getCheckpointConfig().setCheckpointStorage(hdfs://hadoop102:8020/flinkCDC);这段配置的逻辑是Flink CDC 把 Binlog 的读取位点保存在 Checkpoint 里任务失败重启后从最近一次 Checkpoint 恢复不会丢数据也不会重复读。RETAIN_ON_CANCELLATION这个选项很关键——如果你手动 Cancel 任务时没保留 Checkpoint下次启动就只能从头全量读对于大表来说这是灾难。setCheckpointStorage指向 HDFS生产环境必须配用本地文件系统在任务重新调度后会找不到状态。Source 的构建用MySqlSource.Stringbuilder()关键参数包括hostname、port、databaseList、tableList、username、password、deserializer和startupOptions。startupOptions有五个选项initial()先全量快照再读增量earliest()从 Binlog 开头读latest()只读启动后的变更specificOffset()从指定位点读timestamp()从指定时间戳读。第一次跑用initial()后续恢复用 Savepoint 或者 Checkpoint 自动管理。3.2 Flink SQL 方式DDL 里藏着的细节Flink SQL 方式代码量少适合快速验证和简单同步。核心就是一段CREATE TABLEDDLCREATE TABLE t1 ( id STRING PRIMARY KEY NOT ENFORCED, name STRING ) WITH ( connector mysql-cdc, hostname hadoop103, port 3306, username root, password 000000, database-name test, table-name t1 );这里有几个坑要注意。第一PRIMARY KEY NOT ENFORCED是必须的Flink SQL 不强制主键约束但 CDC 连接器需要知道哪个字段是主键来生成正确的变更事件。第二database-name和table-name支持正则比如table-name t.*可以匹配多张表但正则匹配的表结构必须一致否则运行时会报 schema 不兼容。第三Flink SQL 方式下 Checkpoint 配置是在flink-conf.yaml里全局设置的不是代码里所以如果你在 IDE 里跑默认没有 Checkpoint断点续传不生效。3.3 断点续传的验证方法教程里给了一个很实用的验证流程启动程序 → 创建 Savepoint → Cancel Job → 在 MySQL 里插入新数据 → 从 Savepoint 重启 → 观察日志。这个流程能验证两件事一是 Savepoint 是否正确保存了 Binlog 位点二是重启后是否从断点继续而不是从头全量读。实际操作中创建 Savepoint 的命令是bin/flink savepoint JobId hdfs://hadoop102:8020/flinkCDC/save然后在 WebUI 里 Cancel Job再从 Savepoint 重启bin/flink run -s hdfs://.../savepoint-xxx -c com.atguigu.cdc.FlinkCDCDataStreamTest ./flink-cdc-test.jar。如果你看到日志里只输出了新增的那条数据说明断点续传生效了如果又从头打印了所有历史数据说明 Savepoint 没生效检查 Checkpoint 存储路径和RETAIN_ON_CANCELLATION配置。注意Savepoint 和 Checkpoint 的区别在于Savepoint 是手动触发的、格式稳定的、用于有计划重启的Checkpoint 是自动周期性的、用于故障恢复的。生产环境做版本升级时用 Savepoint日常故障恢复靠 Checkpoint。4. MySQL 到 Doris 的 Pipeline 同步YAML 配置与路由变更4.1 Pipeline YAML 的字段含义与参数调优Pipeline 模式是 3.0 的亮点配置文件结构分三块source、sink、pipeline。source: type: mysql hostname: node01 port: 3306 username: root password: 123 tables: test.\.* server-id: 5400-5404 server-time-zone: UTC8 sink: type: doris fenodes: node01:7030 username: root password: 000000 table.create.properties.light_schema_change: true table.create.properties.replication_num: 1 pipeline: name: Sync MySQL Database to Doris parallelism: 1source.tables用的是正则表达式test.\.*表示 test 库下所有表。server-id给的是一个范围Flink CDC 会从这个范围里选一个作为伪装从库的 ID注意不要和现有 MySQL 从库的 ID 冲突。server-time-zone设成UTC8避免时间字段出现 8 小时偏差这个坑很隐蔽——数据同步过去看着都对就是时间不对。sink侧的fenodes是 Doris FE 的地址table.create.properties是建表时的属性。light_schema_change开启轻量级 Schema Change这样 MySQL 加列时 Doris 侧能自动跟上不用手动改表。replication_num设成 1 是单副本生产环境至少设 3。pipeline.parallelism控制同步任务的并行度。源表多、数据量大就调高但别超过 Flink 集群的 slot 总数。4.2 启动任务与 Schema 变更验证启动前确认两件事Flink 集群的flink-conf.yaml里execution.checkpointing.interval已经设置比如 5000 毫秒Doris 的 FE 和 BE 都已经启动。然后在 Doris 里创建目标数据库test执行bin/flink-cdc.sh job/mysql-to-doris.yaml。验证分两步先看全量数据是否同步过去SELECT COUNT(*)对比 MySQL 和 Doris 的行数再做增量测试在 MySQL 里INSERT、UPDATE、DELETE刷新 Doris 看变化是否实时反映。教程里还特别提到了“新增列”的测试——在 MySQL 的test库某张表加一列观察 Doris 对应表是否自动增加该列。这就是light_schema_change的作用如果没开这个属性加列操作会失败同步任务报错。4.3 路由变更一张表同步到不同目标路由变更解决的是“源库多张表合并到目标库一张表”或者“源库一张表拆分到目标库多张表”的需求。教程里给的test_route示例是MySQL 的test_route库和test库都有t1、t2、t3通过路由规则把test_route.t1的数据写到 Doris 的doris_test_route.t1而不是默认的test.t1。路由配置写在 YAML 的route段里语法是source-table: sink-table。这个功能在分库分表场景下特别有用——多个 MySQL 实例的同名表可以汇聚到 Doris 的一张宽表里。但要注意路由后的表结构必须兼容如果源表 A 有name字段而源表 B 没有同步到同一张目标表时会因为 schema 不匹配报错。提示Pipeline 模式目前不支持自定义 UDF 和复杂的字段转换。如果你需要在同步过程中做数据脱敏、字段拼接、类型转换还是得用 Flink SQL 方式在 DDL 之后接一个INSERT INTO sink_table SELECT ... FROM source_table。5. 避坑与排查那些让我加班到凌晨的配置项5.1 Binlog 没开或者格式不对现象任务启动后报The MySQL server is not configured to use a row-based binlog或者直接连不上 Binlog。原因binlog_format不是ROW或者log-bin没开或者用户没有REPLICATION SLAVE和REPLICATION CLIENT权限。解决登录 MySQL 执行SHOW VARIABLES LIKE binlog_format和SHOW VARIABLES LIKE log_bin确认。权限问题用GRANT REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO root%;授权。改完配置必须重启 MySQLFLUSH PRIVILEGES不够。5.2 Checkpoint 失败导致任务反复重启现象Flink WebUI 上看到 Checkpoint 一直失败任务不断重启日志里报Checkpoint expired before completing。原因Checkpoint 存储路径不可写HDFS 权限问题或者 Checkpoint 间隔太短、超时时间不够或者状态太大导致 Checkpoint 做不完。解决确认HADOOP_USER_NAME设置正确HDFS 路径有写权限。调大setCheckpointTimeout和setMinPauseBetweenCheckpoints。如果状态确实大考虑增大 TaskManager 内存或者启用增量 CheckpointRocksDB 状态后端。5.3 Doris Sink 写入报 schema 不兼容现象Pipeline 任务启动后Doris 侧报Table column count mismatch或者Unknown column。原因MySQL 表结构和 Doris 表结构不一致或者light_schema_change没开导致加列不同步。解决手动在 Doris 里建好和目标 MySQL 一致的表结构或者在 YAML 里开启table.create.properties.light_schema_change: true。注意 Doris 的字段类型和 MySQL 不是一一对应的VARCHAR长度、DATETIME精度都可能需要调整。5.4 server-id 冲突导致 Binlog 读取中断现象任务运行一段时间后突然报A slave with the same server_uuid/server_id as this slave has connected to the master。原因Flink CDC 伪装从库的server-id和现有 MySQL 从库或者其他 CDC 任务的server-id冲突了。解决在 YAML 或 DataStream 配置里给server-id指定一个独立范围比如5400-5404确保不和现有从库重叠。如果同一个 MySQL 实例上有多个 CDC 任务每个任务的server-id范围要错开。5.5 时间字段出现 8 小时偏差现象MySQL 里的CREATE_TIME是2024-01-15 10:00:00同步到 Doris 变成2024-01-15 02:00:00。原因Flink CDC 默认用 UTC 时区解析 Binlog 时间戳而 MySQL 服务器用的是Asia/Shanghai。解决在 Source 配置里显式设置server-time-zone: UTC8。DataStream 方式是在MySqlSource.builder()里加.serverTimeZone(UTC8)Flink SQL 方式是在 DDL 的WITH里加server-time-zone UTC8。6. 进阶技巧用 Savepoint 做版本升级和无缝迁移跑通基础同步之后真正体现功力的是版本升级和集群迁移时的无缝切换。Flink CDC 的 Savepoint 机制在这里能救命——你可以在不丢数据、不重复读的前提下把任务从一个 Flink 集群迁移到另一个或者从 CDC 2.x 升级到 3.0。具体操作流程是这样的假设你有一个正在运行的 DataStream 任务用的是 CDC 2.4现在想升级到 3.0。第一步对当前任务触发 Savepointbin/flink savepoint JobId hdfs://hadoop102:8020/flinkCDC/savepoint。第二步Cancel 当前任务但不要删除 Savepoint。第三步用新的 CDC 3.0 依赖重新打包代码注意 Source 的startupOptions要改成Savepoint模式或者不设置Flink 会自动从 Savepoint 恢复。第四步从 Savepoint 启动新任务bin/flink run -s hdfs://.../savepoint-xxx -c com.atguigu.cdc.FlinkCDCDataStreamTest ./flink-cdc-3.0-test.jar。这里的关键是 Savepoint 里保存了 Binlog 的位点信息新任务从那个位点继续读不会重复也不会丢。但有一个坑CDC 2.x 和 3.0 的状态结构可能不兼容直接恢复会报State migration failed。我的经验是跨大版本升级时先停任务、记录当前 Binlog 位点SHOW MASTER STATUS然后用startupOptions.specificOffset()从指定位点启动新任务而不是依赖 Savepoint。这样虽然需要手动记录位点但兼容性最好。另一个实用技巧是配合 Doris 的light_schema_change做在线加列。MySQL 业务表加字段是常态如果每次都要停同步任务、手动改 Doris 表结构运维成本太高。开启light_schema_change后MySQL 执行ALTER TABLE ADD COLUMNFlink CDC 捕获到 DDL 变更自动在 Doris 侧执行对应的 Schema Change同步任务不中断。但要注意Doris 的 Schema Change 是异步的加列操作在 Doris 侧生效有延迟如果紧接着就往新列写数据可能会短暂报错。还有一个我踩过的坑Pipeline 模式下如果 MySQL 源表做了RENAME TABLE或者DROP TABLEFlink CDC 默认会报错停止。生产环境如果允许这类 DDL需要在 YAML 里配置schema-change-behavior参数设成LENIENT模式让任务跳过不支持的 DDL 而不是直接挂掉。从那以后我每次上线同步任务前都强制走一遍“Savepoint 创建 → Cancel → 从 Savepoint 恢复 → 验证数据不重不丢”的流程确认断点续传真的生效了才敢交给业务方用。希望帮到你。本文还有配套的精品资源点击获取