编程 用 Flink CDC 把 MySQL 实时同步到 Doris:Pipeline 配置与特性梳理

2026-09-05 07:10:25

数据库到下游系统的实时数据同步,要同时兼顾稳定性、延迟和成本。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 集群。更多数据管道构建示例见官方文档:

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

推荐文章

程序员茄子在线接单