编程 Flink CDC 3.0 用 YAML 跑 MySQL 到 Doris:三层 API、连接器清单和 Checkpoint 上的坑

2026-09-11 23:40:16

Flink CDC 3.0 用 YAML 跑 MySQL 到 Doris:三层 API、连接器清单和 Checkpoint 上的坑

项目信息

  • 仓库:
  • 文档:
  • 许可:Apache-2.0;语言:Java;主分支:master。

到底是什么

Flink CDC 是一个基于流的数据集成工具,目标是给数据集成用户一套更完整的编程接口。它支持用 YAML 配置文件定义 ETL(Extract、Transform、Load)流程,由框架自动生成定制化的 Flink 算子并提交作业。

官方列出的核心能力:

  • 端到端的数据集成框架
  • 为数据集成用户提供易于构建作业的 API
  • 支持在 Source 和 Sink 中处理多个表
  • 整库同步
  • 表结构变更自动同步(Schema Evolution)

在任务提交过程中还做了优化,另外增加了数据转换(Data Transformation)、整库同步以及精确一次(Exactly-once)语义这些高级特性。

三层 API

YAML API:声明式、零代码方式定义数据管道。在 YAML 里描述 source、sink、routing、transformation、schema evolution 规则,通过 flink-cdc.sh 提交。

SQL API:与 Flink SQL 集成,用 SQL DDL 定义 CDC source。把 connector JAR 放到 FLINK_HOME/lib/,在 Flink SQL Client 里直接使用。

DataStream API:编程方式构建自定义 Flink 流应用,把对应 connector 作为 Maven 依赖加入。所有 artifact 的 group ID 为 org.apache.flink

YAML 示例,从 MySQL 捕获实时变更同步到 Apache Doris:

source:
type: mysql
hostname: localhost
port: 3306
username: root
password: 123456
tables: app_db.\\.*
server-id: 5400-5404
server-time-zone: UTC

sink:
type: doris
fenodes: 127.0.0.1:8030
username: root
password: ""
table.create.properties.light_schema_change: true
table.create.properties.replication_num: 1

pipeline:
name: Sync MySQL Database to Doris
parallelism: 2

提交方式是 flink-cdc.sh 加 YAML 文件,一个 Flink 作业会被编译并部署到指定 Flink 集群。

版本与架构

版本历程(来自阿里云 Flink Committer 任庆盛的分享):

  • 2020 年 7 月诞生。
  • 2021 年 8 月发布 2.0,首次为 MySQL CDC source 引入增量快照算法,实现全增量同步无缝切换。
  • 2022 年 11 月发布 2.3,将大多数 connector 对接至增量快照框架。
  • 2023 年 12 月推出 3.0,正式升级为实时数据集成框架,提供 YAML API,为数据同步提供端到端解决方案。

架构上,Flink CDC 基于 Flink runtime 实现,复用 Flink 的资源管理能力和跨环境部署能力。针对数据集成场景深度定制了多种自定义算子,例如 schema operator、router、transformer,并引入 composer 组件来协调组合这些算子,按用户定义的同步逻辑构建 Flink 作业。CLI 的作用是:一个脚本就把 YAML 定义交给 composer 构建成 Flink 作业并提交至指定集群。

connector 本身基于 Flink connector,简单数据转换封装后即可复用。schema 变更同步方面,定义了 Accessor 和 MetadataApplier,分别负责源端和 sink 端 schema 元信息的获取与处理。

连接器与部署

数据源连接器覆盖:MySQL、PostgreSQL、Oracle、SQL Server、MongoDB、OceanBase、TiDB、Db2、Vitess。下游可接数据湖仓,包括 Iceberg、Hudi、Doris、StarRocks 等。

部署方面,Pipeline 可以提交到 standalone、Kubernetes、YARN 等不同 Flink 集群部署模式。

实操注意点

官方教程是「使用 Flink CDC 构建实时数据湖」:MySQL 分库分表同步到 Iceberg,基于 Docker 演示,只涉及 SQL,不需要写 Java/Scala 代码。三个容易踩的点:

  1. Checkpoint 默认不开启。要让 Iceberg 能提交事务,必须显式开启 Checkpoint。
  2. mysql-cdc 在 binlog 读取阶段开始前,需要等待一个完整 checkpoint,用来避免 binlog 记录乱序。
  3. 分库分表合并时,可以定义复合主键 (database_name, table_name, id),避免不同库表 id 相同导致冲突。

未来规划集中在扩展生态(如 Iceberg 等)、支持更多变更类型,以及持续支持对异常处理方式进行自定义。

结论

如果你的场景是 MySQL 等关系库整库同步到湖仓,且希望零代码落地,3.0 起的 YAML API 是接入成本最低的一条路径;需要更细粒度控制或嵌入已有 Flink 作业时,退回 SQL API 或 DataStream API。对本站做 MySQL/运维方向的读者,优先确认三件事:Checkpoint 有没有开、binlog 乱序窗口能不能忍、分库分表的主键够不够唯一。

复制全文 生成海报 Flink CDC 数据集成 MySQL 实时同步 Java

推荐文章

程序员茄子在线接单