Flink CDC Pipeline 是基于 Apache Flink 的数据集成能力,允许用户通过一份 YAML 配置文件来声明一个数据同步作业,而无需编写任何 Java 代码。用户只需描述数据从何处来(source)、要写入何处(sink),以及(可选)如何对数据进行路由(route)与转换(transform),Flink CDC 便会在作业提交时自动将这些声明编译为一组 Flink 算子,并部署为一个 Flink 作业运行。
Pipeline 作业的组成
一个 Pipeline 作业由以下部分组成:
组成部分 | 是否必填 | 说明 |
source | 是 | 定义数据源,负责访问元数据并从外部系统读取变更数据。 |
sink | 是 | 定义数据目标,负责应用表结构变更并向外部系统写入变更数据。 |
pipeline | 是 | 定义作业级配置,例如作业名称、并行度、时区、表结构变更行为等。 |
route | 否 | 定义源表到目标表的路由规则,支持分库分表合并同步。 |
transform | 否 | 定义字段投影、过滤、主键/分区键重设等数据转换规则。 |
其中,source、sink 和 pipeline 是一个作业必需的三个部分,route 和 transform 为可选部分。
核心特性
整库同步:通过库表匹配规则,在一个作业中同步一个数据库实例下的多张表。
表结构演进(Schema Evolution):自动将上游 DDL 变更同步到下游系统。
数据转换(Transform):支持字段投影、过滤、计算列以及内置与自定义函数。
精确一次(Exactly-Once):作业失败重启后仍能保证端到端的数据不重不丢。
示例
下面是一个将 MySQL app_db 库下的所有表同步到 Doris 的 Pipeline 作业:
source:type: mysqlhostname: localhostport: 3306username: rootpassword: 123456tables: app_db.\\.*sink:type: dorisfenodes: 127.0.0.1:8030username: rootpassword: ""pipeline:name: Sync MySQL Database to Dorisparallelism: 2