用 Flink CDC 把 MySQL 实时同步到 Doris:Pipeline 配置与特性梳理
数据库到下游系统的实时数据同步,要同时兼顾稳定性、延迟和成本。Canal + Kafka + 自研消费链路复杂,离线同步 + 定时调度又做不到实时。Flink CDC 是另一种选择。
Flink CDC 是基于 Apache Flink 的流式数据集成平台,实时捕获外部数据源中的 INSERT、UPDATE、DELETE 以及 DDL 变更,以流的方式同步到消息队列、数据库、数据仓库或数据湖。项目由 Apache Flink 深度集成并驱动,源码托管在 GitHub:
功能特性
- 基于日志的 CDC:通过解析数据库操作日志(MySQL Binlog、PostgreSQL WAL、Oracle Redo Log、MongoDB OpLog)实现变更捕获,亚秒级延迟、低侵入、无锁,性能和一致性优于查询型 CDC。
- 基于 YAML 的数据管道:Pipeline 模型通过配置文件自动生成并提交 Flink 作业,基本无需编写 Java 或 SQL 代码,降低了使用门槛。
- 丰富的连接器:Source 端支持 MySQL、Oracle、PostgreSQL、Db2、MongoDB、SQL Server、TiDB、Vitess;Sink 端支持 Apache Doris、Elasticsearch、Fluss、Hudi、Iceberg、Kafka、MaxCompute、OceanBase、Paimon、StarRocks。
- 表结构自动同步:Schema Evolution 将上游 DDL 变更事件同步到下游,覆盖新建表、添加列、重命名列、更改列类型、删除列、截断和删除表等操作。
- 全量增量一体化:增量快照框架在启动时自动做一致性全量快照,完成后无缝切换增量日志。
- 精确一次语义:Exactly-Once 保证每条数据变更在整个链路中只被处理一次,只产生一次结果。
- 数据转换:Transform 模块支持字段级或行级加工,包括字段选择、重命名、派生字段与函数计算。
- 路由规则:Route 规则匹配一个或多个源表并映射到目标表,常用于合并子库子表,将多个上游源表路由到同一个目标表。
配置示例
以下 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 集群。更多数据管道构建示例见官方文档: