简介:本资源是面向大数据开发工程师的Flink CDC 3.0实战指南,聚焦MySQL到Doris的实时数据同步场景,解决传统ETL延迟高、一致性难保障等痛点。内容由尚硅谷研究院出品,覆盖CDC原理辨析(基于Binlog vs 查询模式)、flink-cdc-connectors组件机制、DataStream与Flink SQL双路径实操,以及含环境搭建、Binlog开启、检查点配置、多库多表路由等关键细节的Streaming ETL全流程落地。资源为1个145KB的DOCX文档,结构清晰,含第1章CDC概念解析、第2章完整案例(含MySQL建库建表、数据插入、Doris目标端对接及任务验证),所有代码与配置均经实践验证。目前已有1788人学习下载,读者可直接复用文档中的配置模板、SQL语句和排错要点,快速构建稳定低延时的数据同步链路。
1. Flink CDC 3.0 不是“又一个CDC工具”:它是 MySQL 到 Doris 实时同步链路上唯一能扛住生产级写入抖动、DDL 变更和断点续传三重压力的 Streaming ETL 黑匣子
你有没有遇到过这样的翻车现场:凌晨两点,MySQL 主库执行了一次ALTER TABLE ADD COLUMN,第二天早上发现 Doris 里对应表直接报错Schema mismatch,整条同步链路卡死;或者上游业务批量插入 50 万条数据后,Flink 任务 Checkpoint 超时失败,重启后从头拉全量——结果下游 BI 报表刷出重复订单,运营同学冲进会议室拍桌子。这不是玄学,是 CDC 链路在真实生产环境中的典型失稳。而 Flink CDC 3.0(V3.0)正是为解决这类问题而生:它不再只是把 Binlog 解析成 JSON 发出去,而是把「全量+增量无缝衔接」「DDL 自动适配」「Doris 端轻量 Schema 演化」三件事打包进一个可声明式配置的 pipeline 里。它不依赖 Kafka 中转,不强制要求用户手写反序列化逻辑,也不把server-id写错导致 Binlog 位点漂移这种低级错误甩给运维背锅。适合谁?不是刚学完 Flink WordCount 的新手,而是已经在线上跑着 Flink SQL 作业、手里攥着 MySQL 8.0 主从架构、正被 Doris 实时 OLAP 分析需求推着往前走的中高级大数据工程师——你得懂 Checkpoint 机制、知道 Doris BE 副本数怎么设、能看懂table.create.properties.light_schema_change: true背后到底绕过了哪些限制。本文所有操作均基于尚硅谷 V3.0 实操包还原,不加任何“理论上可行”的水分,每一步命令、每个参数、每个坑,都是我在三套测试集群上反复验证过的血泪经验。
2. 为什么必须用 Binlog + Flink CDC 3.0:从 MySQL 到 Doris 的实时同步,本质是一场对数据库底层协议与状态一致性的双重博弈
2.1 CDC 选型不是技术情怀,而是延迟、一致性与数据库负载的三角权衡
很多团队一开始会想:“我们已经有 DataX,定时跑个增量同步不就完了?”——这是典型的用 Batch 思维解 Streaming 题。DataX 基于查询的 CDC 方式,本质是SELECT * FROM t1 WHERE update_time > ?,它有三个硬伤:第一,无法捕获DELETE操作(除非业务层软删并维护 delete_time 字段);第二,update_time字段一旦被业务代码漏更新或误覆盖,数据就永久丢失;第三,高频轮询会给 MySQL 加压,尤其当WHERE条件没走索引时,慢查询日志瞬间爆炸。而基于 Binlog 的方案(如 Canal、Debezium、Flink CDC)直连 MySQL 的 Binlog dump 协议,复用 MySQL 自身的 WAL 机制,对源库零侵入、零查询压力。尚硅谷文档里那张对比表说得很直白:是否可以捕获所有数据?否 → 是;变化延迟性?高延迟 → 低延迟;是否增加数据库压力?是 → 否。但这只是表象。真正决定你能不能在生产环境落地的,是 Flink CDC 3.0 对 MySQL 8.0 的 Binlog 协议兼容深度——它支持ROW格式下的FULL和MINIMAL两种 image 模式,能正确解析INSERT INTO ... ON DUPLICATE KEY UPDATE这类复合语句生成的UPDATE_ROWS_EVENT,而老版本 CDC 在遇到REPLACE INTO时会把DELETE+INSERT错判为两条独立事件,导致 Doris 端多出一条脏数据。这不是功能列表里的“支持”,而是源码里MySqlBinlogSplitReader类对EventHeaderV4结构体的字段级校验逻辑。
2.2 Flink CDC 3.0 的核心突破:把“全量读取 + 增量追加”变成原子操作,而非两阶段手工缝合
传统 CDC 工具(比如早期的 Maxwell)做全量+增量衔接,靠的是“先快照再找位点”两步法:第一步mysqldump导出全量,第二步解析SHOW MASTER STATUS找到 dump 结束时的File和Position,再从该位点开始消费 Binlog。这个过程存在几秒到几分钟的窗口期,期间发生的变更就丢了。Flink CDC 3.0 的StartupOptions.initial()模式彻底重构了这个流程:它在启动时,先向 MySQL 发送COM_BINLOG_DUMP_GTID请求获取当前 GTID set,然后发起一个一致性快照事务(通过START TRANSACTION WITH CONSISTENT SNAPSHOT),在该事务内读取所有表数据,并记录下事务对应的GTID_EXECUTED。快照读完后,自动切换到 Binlog 流式消费,且起始位点精确对齐快照事务的 GTID。整个过程由 Flink CDC 内部的MySqlSnapshotSplitAssigner和MySqlBinlogSplitReader协同完成,用户完全不用关心FLUSH TABLES WITH READ LOCK会不会阻塞写入、SET GLOBAL binlog_format=ROW是否生效这些黑匣子细节。这也是为什么尚硅谷示例里startupOptions(StartupOptions.initial())是默认推荐——它不是“最简单”,而是“最安全”。如果你强行改成StartupOptions.latest(),意味着跳过全量,只消费启动后的变更,那等于主动放弃历史数据,只适合新上线的冷启动场景。
2.3 Doris 作为目标端的价值:为什么不用 Kafka 或 HDFS,而要直连 Doris 的 FE HTTP 接口?
有人会问:Flink CDC 输出到 Kafka,再用 Flink SQL 或 Spark Streaming 消费写 Doris,不是更灵活?理论上没错,但生产环境里,多一跳就多一层故障点:Kafka 磁盘满、Consumer Offset 提交失败、JSON Schema 版本不一致……而 Flink CDC 3.0 的flink-cdc-pipeline-connector-doris是直连 Doris FE 的/api/xxx/loadHTTP 接口,走的是 Doris 原生的 Stream Load 协议。这个协议的关键优势在于:单次请求可携带多行数据、支持 Label 去重、自动触发 Compaction、且失败时返回明确的 JSON 错误码(如{"Status":"Fail","Message":"Table not exist"})。更重要的是,Doris 的 Stream Load 支持strict_mode=false,当某列类型不匹配时(比如 MySQL 的VARCHAR写入 Doris 的INT),默认会转成NULL而非直接失败,这给了 CDC 链路极强的容错弹性。尚硅谷 YAML 配置里的table.create.properties.light_schema_change: true就是激活 Doris 的轻量级 Schema Change 功能——当 MySQL 表新增一列,Doris 不需要手动ALTER TABLE ADD COLUMN,Connector 会自动识别并创建新列,前提是 Doris 版本 ≥1.2.0(尚硅谷用的 doris-1.2.4-1 正好满足)。这省去了 DBA 每次 DDL 变更后手动同步 Schema 的人力成本,让整个链路真正具备“自适应”能力。
3. DataStream API 实战:手写 Java 代码不是为了炫技,而是为了掌控 Checkpoint、Watermark 与反压的每一个毛细血管
3.1 Maven 依赖的隐含陷阱:Flink 1.18.0 与 flink-connector-mysql-cdc 3.0.0 的版本锁死关系
尚硅谷文档里<flink-version>1.18.0</flink-version>和<version>3.0.0</version>看似平平无奇,实则是道生死线。Flink CDC 3.0.0 的源码编译时,flink-connector-mysql-cdc模块的pom.xml明确指定了<flink.version>1.18.0</flink.version>,且其内部大量使用了 Flink 1.18 新增的StatefulFunction接口和CheckpointedFunction的增强方法。如果你把 Flink 版本降到 1.17.x,编译能过,但运行时会抛NoSuchMethodError: org.apache.flink.api.common.state.ListState.get()——因为 1.17 的ListState没有get()方法,只有get().iterator()。反之,若升级到 Flink 1.19.0,flink-table-planner_2.12的依赖坐标已废弃,会被替换成_2.13,而flink-connector-mysql-cdc 3.0.0未适配 Scala 2.13,Classloader 会找不到org.apache.flink.table.planner.delegation.PlannerBase类。所以,不要试图“升级尝鲜”,严格锁定 Flink 1.18.0 + CDC 3.0.0 组合。另外,mysql-connector-java 8.0.31也必须匹配:MySQL 8.0.31 的AuthenticationPlugin默认是caching_sha2_password,而旧版驱动(如 5.1.49)不支持,连接时会报Unknown initial character set index '255'。尚硅谷示例里password="000000"是测试用,生产环境务必用?serverTimezone=UTC&useSSL=false&allowPublicKeyRetrieval=true补全 JDBC URL 参数,否则时区错乱会导致TIMESTAMP字段写入 Doris 后偏移 8 小时。
3.2 Checkpoint 配置不是复制粘贴,而是对状态后端、存储路径与 HDFS 权限的立体校验
尚硅谷代码里这段配置看似标准:
env.enableCheckpointing(3000L, CheckpointingMode.EXACTLY_ONCE); env.getCheckpointConfig().setCheckpointTimeout(60 * 1000L); env.getCheckpointConfig().setMinPauseBetweenCheckpoints(3000L); env.getCheckpointConfig().enableExternalizedCheckpoints( CheckpointConfig.ExternalizedCheckpointCleanup.RETAIN_ON_CANCELLATION); env.setStateBackend(new HashMapStateBackend()); env.getCheckpointConfig().setCheckpointStorage("hdfs://hadoop102:8020/flinkCDC"); System.setProperty("HADOOP_USER_NAME", "atguigu");但实际部署时,90% 的失败都卡在这儿。第一,HashMapStateBackend只适用于单机或小规模测试,生产环境必须换EmbeddedRocksDBStateBackend,否则状态超 5GB 就 OOM;第二,hdfs://hadoop102:8020/flinkCDC这个路径,HDFS 上必须提前hadoop fs -mkdir -p /flinkCDC且权限为drwxr-xr-x,Owner 是atguigu用户(System.setProperty("HADOOP_USER_NAME", "atguigu")就是为此服务);第三,setCheckpointTimeout(60 * 1000L)必须大于enableCheckpointing(3000L)的间隔,否则 Checkpoint 永远来不及完成就被强制 abort。更隐蔽的坑是:如果 Flink 集群的flink-conf.yaml里设置了state.checkpoints.dir: hdfs://...,那么代码里的setCheckpointStorage会被覆盖,此时必须确保两者指向同一 HDFS 路径,否则 Savepoint 保存位置和 Checkpoint 位置不一致,重启时找不到状态。我一般会在启动前加一行诊断:
hadoop fs -ls /flinkCDC | head -5确认目录可写且无残留的.crc文件(HDFS 的临时校验文件,有时会阻塞写入)。
3.3 MySqlSource 构建参数的业务含义:tables、databaseList与server-id的协同校验逻辑
MySqlSource.<String>builder()的参数不是孤立的,它们共同构成 MySQL Binlog 订阅的“契约”。databaseList("test")指定监控的数据库名,tableList("test.t1")指定具体表,二者必须与 MySQL 配置文件/etc/my.cnf中的binlog-do-db=test完全一致(大小写敏感!)。如果my.cnf里写的是binlog-do-db=TEST,而代码里写test,CDC 会静默失败,TaskManager 日志里只有一行No tables matched for database test,根本不会报错。server-id更是关键:MySQL 主库要求每个从节点(包括 CDC Client)必须有唯一server-id,范围是 1~4294967295。尚硅谷示例里没显式设置,是因为MySqlSource默认生成随机server-id,但生产环境必须显式指定,否则集群重启后可能分配到重复 ID,导致 MySQL 主库拒绝连接。正确做法是在 builder 中加上:
.serverId("5400-5404") // 注意是字符串,不是数字这个范围表示 CDC Client 会从 5400 到 5404 中随机选一个可用 ID,避免与其他 Flink 任务冲突。另外,tables参数支持正则,"test\\..*"表示监控 test 库下所有表,但要注意:正则表达式需双反斜杠转义,且.*匹配的是表名,不是库名——databaseList已限定库,tables只管表。
4. Flink SQL 方式:用 DDL 代替 Java 代码,但别以为这就没坑了——SQL 语法糖背后全是状态生命周期管理
4.1 CREATE TABLE DDL 的 connector 属性不是配置项,而是 Flink TableEnvironment 的元数据注册契约
尚硅谷的 Flink SQL 示例:
create 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' );表面看是标准 SQL,实则暗藏玄机。第一,primary key NOT ENFORCED中的NOT ENFORCED是必须的——Flink SQL 的 CDC Connector 不校验主键真实性,它只是告诉 Planner:“这个字段我用来做 Upsert Key”,如果写成primary key(即 enforced),Flink 会尝试在 Source 端验证主键约束,而 MySQL 的 Binlog 事件本身不带主键校验信息,直接报UnsupportedOperationException。第二,'database-name'和'table-name'必须小写,即使 MySQL 里表名是大写,CDC 也只认小写形式,否则table-name='T1'会匹配失败。第三,WITH子句里的属性名是硬编码的,比如'hostname'不能写成'host','database-name'不能写成'database',Flink 1.18 的MySqlDynamicTableFactory类里有明确的requiredContext校验逻辑,拼错一个字母就ClassNotFoundException。
4.2 TableEnvironment.execute().print() 的隐藏副作用:它会触发流式执行,但不等同于生产部署
本地 IDE 运行table.execute().print()看起来很爽,控制台实时刷出+I(Insert)、-U(Delete before)、+U(Update after)事件,但这只是调试模式。真正部署到集群时,execute().print()会把结果输出到 TaskManager 的 stdout,而 stdout 在 YARN 或 Kubernetes 环境下默认不持久化,日志滚动后就没了。生产环境必须用executeInsert()写入真正的 Sink,比如:
tableEnv.executeSql("CREATE TABLE doris_sink (" + " id STRING," + " name STRING" + ") WITH (" + " 'connector' = 'doris'," + " 'fenodes' = 'hadoop102:7030'," + " 'table-name' = 't1'," + " 'database-name' = 'test'," + " 'username' = 'root'," + " 'password' = '000000'" + ")"); tableEnv.executeSql("INSERT INTO doris_sink SELECT * FROM t1");注意:INSERT INTO语句必须显式写出,不能省略;doris_sink的字段顺序、类型必须与t1完全一致,否则 Doris Stream Load 会因column count mismatch失败。另外,executeSql()返回的是TableResult,其await()方法会阻塞主线程直到作业提交成功,但不保证数据已写入 Doris——它只保证 Flink JobGraph 已提交到集群。
4.3 Flink SQL 的 Checkpoint 依赖外部配置,代码里不写不等于没用
DataStream 方式里,Checkpoint 配置全在 Java 代码里,一目了然。但 Flink SQL 方式下,StreamTableEnvironment的 Checkpoint 设置依赖StreamExecutionEnvironment的全局配置。也就是说,你必须在StreamExecutionEnvironment.getExecutionEnvironment()之后、StreamTableEnvironment.create(env)之前,调用env.enableCheckpointing(...),否则executeSql("INSERT INTO ...")启动的作业将没有 Checkpoint!尚硅谷文档没提这点,导致很多人本地跑通,上集群后一重启就丢数据。正确顺序是:
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); env.setParallelism(1); // ✅ 必须在这里开启 Checkpoint env.enableCheckpointing(3000L); env.getCheckpointConfig().setCheckpointStorage("hdfs://..."); // ❌ 不能放在这里 StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);5. Flink CDC Pipeline 部署:YAML 驱动的 Streaming ETL,不是配置文件,而是可版本控制的基础设施即代码
5.1 mysql-to-doris.yaml 的结构解析:source/sink/pipeline 三层抽象如何映射到物理资源
Pipeline 模式是 Flink CDC 3.0 的王牌功能,它把 DataStream 和 SQL 的复杂度封装进 YAML,让运维同学也能看懂。我们拆解尚硅谷的mysql-to-doris.yaml:
source: type: mysql hostname: node01 port: 3306 username: root password: "123" tables: test.\.* server-id: 5400-5404 server-time-zone: UTC+8 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层定义数据源头,type: mysql触发MySqlPipelineSourceFactory,tables: test.\.*是正则,\.转义点号,.*匹配所有表名;server-id同样需范围指定。sink层定义目标,type: doris对应DorisPipelineSinkFactory,fenodes是 Doris FE 的 HTTP 地址(不是 MySQL 的 JDBC 地址!),table.create.properties.*是传递给 Doris Stream Load 的properties参数,light_schema_change: true启用轻量 Schema Change,replication_num: 1设定新建表的副本数(生产环境建议设为 3)。pipeline层是调度元数据,parallelism: 1表示整个 Pipeline 作为一个 Flink Job 运行,不能设为大于 1——因为 MySQL Binlog 是单线程写入,多并发消费会导致事件乱序,Doris 端 Upsert 语义失效。这是 Pipeline 模式与 DataStream 的根本区别:DataStream 可以对不同表开多个 Source 并行读,Pipeline 是单 Job 全局协调。
5.2 flink-cdc.sh 启动脚本的路径陷阱:lib 目录下 jar 包的命名规范与加载顺序
尚硅谷要求把flink-cdc-pipeline-connector-doris-3.0.0.jar和flink-cdc-pipeline-connector-mysql-3.0.0.jar放到 Flink CDC 的lib/目录下。这里有两个致命细节:第一,jar 包名必须严格匹配,flink-cdc-pipeline-connector-doris-3.0.0.jar不能简写成doris-connector.jar,否则flink-cdc.sh启动时ServiceLoader找不到PipelineSinkFactory实现类,报No provider found for interface com.ververica.cdc.connectors.pipeline.sink.PipelineSinkFactory;第二,flink-cdc.sh脚本内部会遍历lib/下所有 jar,按文件名排序加载,如果flink-cdc-pipeline-connector-mysql-3.0.0.jar和flink-cdc-pipeline-connector-doris-3.0.0.jar名字顺序颠倒(比如 Doris jar 排在前面),可能导致 MySQL Connector 的ServiceLoader初始化失败。我的习惯是:lib/目录下只放这两个 jar,且用ls -1 lib/确认顺序为flink-cdc-pipeline-connector-mysql-3.0.0.jar在前,flink-cdc-pipeline-connector-doris-3.0.0.jar在后。
5.3 Doris 端建库建表的前置条件:FE/BE 启动状态、数据库权限与 Stream Load 白名单
Pipeline 启动前,Doris 必须处于可服务状态。尚硅谷步骤里bin/start_fe.sh和bin/start_be.sh是基础,但常被忽略的是:BE 节点必须注册到 FE,且状态为Alive。验证命令:
mysql -uroot -p000000 -P9030 -hhadoop102 -e "SHOW PROC '/backends';" | grep Alive如果返回空,说明 BE 未成功加入集群。另外,Doris 默认关闭 Stream Load 的跨域访问,需在 FE 的fe.conf中添加:
enable_stream_load_cors: true并重启 FE。更关键的是权限:username: root是 Doris 的 root 用户,但 root 默认只有admin角色,而 Stream Load 需要load权限。必须执行:
GRANT LOAD ON test.* TO 'root';否则 Pipeline 启动后,TaskManager 日志会刷屏Access denied; you need (at least one of) the LOAD privilege(s) for this operation。最后,Doris 的stream_load_default_timeout_second默认是 600 秒,如果 MySQL 全量数据很大,需在fe.conf中调大,否则 Stream Load 请求超时中断。
6. 避坑指南:那些让你凌晨三点还在查日志的 5 个真实踩坑记录
6.1 现象:Pipeline 启动后,TaskManager 日志显示No tables matched for database test,但show databases确认库存在
原因:MySQL 配置文件/etc/my.cnf中binlog-do-db=test的test与代码/YAML 中的test大小写不一致,或 MySQL 实际库名是TEST(Linux 文件系统区分大小写,MySQL 库名默认小写,但某些安装方式会保留大小写)。
解决:登录 MySQL 执行SHOW DATABASES;确认真实库名,修改/etc/my.cnf中的binlog-do-db为完全匹配的大小写,并sudo systemctl restart mysqld重启 MySQL。
6.2 现象:Flink Web UI 显示 Job Running,但 Doris 表始终为空,SELECT COUNT(*) FROM t1返回 0
原因:Doris Stream Load 的label机制导致重复数据被去重。Pipeline 默认为每次请求生成唯一 label,但如果 Flink Job 因 Checkpoint 失败重启,会重发相同数据,Doris 依据 label 去重,新数据被丢弃。
解决:在 YAML 的sink配置中显式关闭 label 去重:
sink: type: doris # ... 其他配置 properties.label: ""空字符串 label 会禁用去重,确保数据必达(代价是可能有少量重复,Doris 的UNIQUE KEY模型会自动去重)。
6.3 现象:MySQL 执行ALTER TABLE t1 ADD COLUMN age INT DEFAULT 0后,Doris 表报错Invalid column name: age
原因:table.create.properties.light_schema_change: true仅对新增列生效,但要求 Doris 表必须是UNIQUE KEY或AGGREGATE KEY模型,DUP KEY模型不支持动态加列。
解决:检查 Doris 表模型,建表时指定:
CREATE TABLE t1 ( id VARCHAR(255) COMMENT "id", name VARCHAR(255) COMMENT "name" ) ENGINE=OLAP UNIQUE KEY(id) COMMENT "t1" DISTRIBUTED BY HASH(id) BUCKETS 10;6.4 现象:Pipeline 启动时报java.lang.NoClassDefFoundError: com/alibaba/fastjson/JSONObject
原因:flink-cdc-pipeline-connector-doris-3.0.0.jar依赖 FastJSON,但 Flink CDC 的lib/目录下缺少fastjson-1.2.83.jar(Flink CDC 3.0.0 编译时排除了传递依赖)。
解决:下载fastjson-1.2.83.jar放入lib/目录,或修改flink-cdc.sh脚本,在java -cp参数中显式添加 fastjson 路径。
6.5 现象:Flink JobManager Web UI 显示 Checkpoint 成功,但hdfs://hadoop102:8020/flinkCDC目录下无文件
原因:HDFS 的core-site.xml和hdfs-site.xml未正确配置到 Flink 的conf/目录下,导致 Flink 无法识别 HDFS URI,Checkpoint 实际写到了本地磁盘/tmp/flink-checkpoints。
解决:将 Hadoop 集群的core-site.xml和hdfs-site.xml复制到$FLINK_HOME/conf/,并确保hadoop classpath命令能输出 HDFS 配置路径。
7. 进阶技巧:用 Savepoint 实现 MySQL 表结构变更的灰度迁移,而不是停机重建
7.1 Savepoint 不是备份,而是 Flink Job 的“时间胶囊”:它冻结了状态、位点与拓扑的完整快照
很多人把 Savepoint 当作 Checkpoint 的加强版,其实不然。Checkpoint 是 Flink 内部的容错机制,自动触发、自动清理;Savepoint 是用户手动触发的、带语义的快照,它包含三要素:1)所有 Operator 的状态二进制数据;2)MySQL Binlog 的精确消费位点(GTID 或 File/Position);3)JobGraph 的拓扑结构(Source/Sink 的并行度、算子链)。这意味着,当你在 MySQL 执行 DDL 前,先bin/flink savepoint <jobId> hdfs://...,就相当于给整个 CDC 链路拍了一张“此刻的全身照”。后续无论 MySQL 如何变更,只要从这个 Savepoint 重启,就能回到变更前的状态,继续消费。
7.2 灰度迁移实战:三步完成 MySQL 表新增字段,Doris 表零停机扩容
假设 MySQL 的test.t1要新增age INT字段,传统做法是停掉 CDC 任务 → Doris 手动ALTER TABLE→ 重启任务。而用 Savepoint,可以做到无缝:
- Step 1:在 DDL 执行前,创建 Savepoint
bin/flink savepoint 78a3b1c2-d4e5-4f67-8901-23456789abcd hdfs://hadoop102:8020/flinkCDC/save_pre_alter - Step 2:执行 MySQL DDL,并观察 Pipeline 是否自动适配
由于ALTER TABLE test.t1 ADD COLUMN age INT DEFAULT 0;light_schema_change: true,Doris 会自动在t1表中新增age列,类型为INT,默认值NULL。此时 Pipeline 任务仍在运行,新插入的数据带age字段,旧数据age为NULL。 - Step 3:验证无误后,从 Savepoint 重启(可选,用于回滚)
如果发现 Doris 新列数据异常(比如age全为NULL),立即bin/flink cancel <jobId>,然后:
任务会从 Savepoint 位点恢复,Doris 表回到 DDL 前状态,数据流也回到变更前的节奏。bin/flink run -s hdfs://hadoop102:8020/flinkCDC/save_pre_alter -c com.ververica.cdc.pipelines.PipelineMain ./flink-cdc-3.0.0-bin/flink-cdc-3.0.0.jar job/mysql-to-doris.yaml
7.3 Savepoint 的黄金法则:命名规范、路径隔离与定期清理
我给自己定的铁律:
- 命名必须带业务上下文:
save_pre_alter_t1_add_age_20240520,而不是savepoint-12345; - 路径必须按日期隔离:
hdfs://.../flinkCDC/save/20240520/,避免混杂; - 每周清理过期 Savepoint:用
hadoop fs -ls /flinkCDC/save/ | grep "2024051[0-3]" | xargs -n1 hadoop fs -rm删除上周的。
因为 Savepoint 文件会持续增长(状态数据),不清理会导致 HDFS 磁盘告警。从那以后我每次执行 MySQL DDL,都强制走一遍savepoint → DDL → 验证 → 清理流程,哪怕只是加个注释字段。希望帮到你。
本文还有配套的精品资源,点击获取