
Flink CDC 2.x 升级 3.x 迁移指南三步法、配置映射表与 5 个高频坑【免费下载链接】flink-cdcFlink CDC is a streaming data integration tool项目地址: https://gitcode.com/GitHub_Trending/flin/flink-cdcApache Flink CDC 是面向 MySQL、PostgreSQL、Oracle 等数据库的实时数据集成工具提供全量增量同步、表路由和 Schema 变更自动传递能力。本文面向仍运行 Flink CDC 2.x 作业的运维与数据平台工程师读完你可以按清单把现有作业改写为 3.x 的 YAML pipeline并完成切换验证与回退准备。为什么现在要动2.x 的三个实际瓶颈现象一每加一个同步链路就要写代码。2.x 的作业用 Java/Scala 构建MySqlSource再各自接 deserializer 和 sink。表多、链路多时改一个过滤条件都要重新打包、重新提交作业数量和代码仓库的 jar 成正比。3.x 怎么解整条链路写成一份 YAML用 CLI 直接提交不再打包自定义代码。链路从一个工程变成一个文件评审和备份都简单了。现象二源表加列后下游要么不动要么作业挂掉。2.x 的 DDL 不会自动传到目标端需要自己写 deserializer 处理多数团队的选择是不处理。3.x 怎么解pipeline 内置 Schema 变更事件流源端的ALTER TABLE会自动同步到 sink前提是 sink 连接器支持该能力如 Doris、StarRocks 的建表参数配合上线前实测一次即可确认。现象三依赖版本对不上就启动失败。2.x 的 groupId 是com.ververicaSQL fat jar 和 connector jar 容易混用Debezium 传递依赖冲突要靠人肉排查。3.x 怎么解统一走org.apache.flink的 pipeline connector jarsource 和 sink 分别放两个 jar 到 Flink CDC 的lib目录即可依赖管理从看代码变成看目录。新旧写法速查表与架构主线2.x 写法3.x 写法注意点依赖com.ververica:flink-connector-mysql-cdcorg.apache.flink:flink-cdc-pipeline-connector-mysql2.x 的 jar 与 3.x 不通用自写 Java 作业打包提交一份 pipeline YAML bin/flink-cdc.sh提交提交脚本在 flink-cdc-dist 发布包中源码见 flink-cdc-cli/table-name user_\.数据库与表分开tables: app_db.\.*库.表合并为一条正则通配符是正则语法*要写成\.*server-time-zone Asia/Shanghaisource 下server-time-zone与 MySQL 服务器实际时区不一致会出现 8 小时偏差自写 deserializer 处理 DDL内置 Schema 变更传递依赖 sink 连接器能力无统一路由route段source-table正则 →sink-table分表合并到单表就靠它flink run -c Main job.jarbash bin/flink-cdc.sh pipeline.yaml旧 savepoint 不能直接给新作业恢复架构主线一句话3.x 把source → 路由/转换 → sink固化为统一 pipeline 模型所有 source 连接器产出统一的变更事件流经过共享的路由与 transform 阶段再进入任意 sink。以 MySQL 同步到 Doris 并合并分表为例一份完整定义如下source: type: mysql hostname: 127.0.0.1 port: 3306 username: root password: 123456 tables: app_db.order\.* server-id: 5400-5403 server-time-zone: UTC route: - source-table: app_db.order\.* sink-table: ods_db.orders sink: type: doris fenodes: 127.0.0.1:8030 username: root password: pipeline: name: mysql-to-doris parallelism: 4版本对应关系以官方文档为准参见 pipeline-connectors/overviewCDC 3.6.x 支持 Flink 1.20.* 与 2.2.CDC 3.5.x 支持 Flink 1.19.与 1.20.*。三步迁移法盘点、转换、切换迁移流程可以压缩成三步每步有一份可勾选清单第一步盘点枚举所有 2.x 作业登记源端库/表、目标端、server-id、时区、并行度。用information_schema统计各表行数作为切换后的比对基线。确认 Flink 集群版本满足目标 CDC 版本要求如 3.6.x 需 Flink 1.20.* 或 2.2.*不满足先排升级。第二步转换按上文模板为每个作业写 YAML表名正则、route、pipeline.parallelism。server-id区间宽度不小于并行度且与同集群其他作业不重叠。把 MySQL JDBC 驱动 jarmysql-connector-java-8.0.27放进Flink CDC 的lib不是 Flink 的lib。在conf/config.yaml确认 checkpoint 间隔已开启增量快照依赖它。第三步切换提交作业bash bin/flink-cdc.sh mysql-to-doris.yaml。Flink UI 确认作业 RUNNING 且首轮 checkpoint 完成。在源端加一列做 DDL 实测确认 sink 端表结构同步更新。新旧并行运行、比对数据后再停止旧作业并保存 savepoint。旧作业的 jar 与 savepoint 归档保留用于回退。运行效果参考官方教程中的 MySQL 到 Doris 流式同步链路含 Schema 变更与分表合并演示见 mysql-to-doris 教程踩坑实录5 个最容易翻车的问题坑 1作业只读到存量数据binlog 不动现象启动后目标端有初始数据增量一直为空。 根因未启用 checkpoint。2.x/3.x 的增量快照算法都依赖 checkpoint 协调全量与增量顺序没开 checkpoint 增量阶段不会推进官方 FAQ 有同款问题记录见 faq.md。 解法在conf/config.yaml开启execution.checkpointing.interval如3s后重启作业。坑 2连接报Authentication plugin caching_sha2_password错误现象启动即失败日志指向认证插件。 根因MySQL 8.x 默认认证方式需要较新的 JDBC 驱动环境里混了旧版驱动。 解法确认 8.0.27 版mysql-connector-java已放入 Flink CDC 的lib目录并清掉同名旧 jar避免版本混杂。坑 3增量数据时间戳整体偏移 8 小时现象全量阶段正常进入 binlog 阶段后时间字段差 8 小时。 根因server-time-zone与 MySQL 服务器实际时区不一致。 解法以服务器时区为准显式配置如server-time-zone: Asia/Shanghai不要留空猜默认值。坑 4新旧作业并行期间偶发 binlog 位点错乱现象目标端偶发乱序或漏事件无明确报错。 根因MySQL 按server-id区分客户端新旧作业或多作业间区间重叠会互相干扰位点。 解法给新旧作业分配互不重叠的server-id区间切换完成后先停旧作业再让新作业独占区间。坑 5拿 2.x 的 savepoint 直接恢复 3.x 作业失败现象启动即报状态不兼容。 根因两代作业算子拓扑不同savepoint 不保证互通。 解法接受重新全量增量。提前确认源端 binlog 保留窗口大于全量耗时切换窗口内不删 binlog失败后可以无数据缺口地重来。上线前验证项与回退预案切换前逐项过一遍阈值不达标不要停旧作业数据一致性逐表比对源端与 sink 行数要求 100% 相等或差异率 0.01% 且全部可解释为在途写入。同步延迟源端写测试行sink 端 P95 可见延迟 5 秒持续观察 30 分钟无回退。checkpoint成功率 100%无连续失败失败会拖慢全量阶段并影响位点推进。DDL 传递源端ALTER TABLE加列sink 端结构实时更新且作业不重启。告警就绪对作业状态与同步延迟配置告警后再放生产流量。回退预案切换当天起旧作业的代码、jar 和最后一个 savepoint 都保留在制品库中。若新作业 24 小时内出现指标不达标或反复重启立即停止新作业用旧 savepoint 重启 2.x 作业即可它从上次位点续读不做全量重扫。兜底措施是源端 binlog 保留窗口覆盖整个验证期建议至少 3 倍验证时长保证任何时刻重新拉起都能续上。收尾迁移的全部工作量就是盘点 → 转换 → 切换核心产出是一份份 YAML。建议今天就从第一步开始先对现有 2.x 作业做一遍完整盘点产出映射表和行数基线。【免费下载链接】flink-cdcFlink CDC is a streaming data integration tool项目地址: https://gitcode.com/GitHub_Trending/flin/flink-cdc创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考