基于 Flink CDC 构建 MySQL 到 Kafka 的 Streaming ELT 整库同步管道
【免费下载链接】flink-cdcFlink CDC is a streaming data integration tool项目地址: https://gitcode.com/GitHub_Trending/flin/flink-cdc
本教程以 Apache Flink CDC 项目中的 MySQL → Kafka 快速入门指南为主体,演示如何使用 Flink CDC CLI 以零代码的方式构建 Streaming ELT 作业,实现 MySQL 整库同步、表结构变更(DDL)同步、路由分发、多分区写入与多格式输出。读完本文,你将能够独立完成从环境搭建、任务配置、提交运行到下游数据验证的完整链路,并掌握背后由本仓库源码支撑的关键参数原理。
本文对应仓库内文档:docs/content.zh/docs/get-started/quickstart-for-1.20/mysql-to-kafka.md,所有配置参数说明均可在对应 connector 源码中进一步核验。
能力概览:Flink CDC 的 Streaming ELT 与传统 DataStream 的区别
Flink CDC 是一款流式数据集成工具,在 MySQL → Kafka 场景下,它允许用户直接在 YAML 中声明式地描述"数据从哪里来、到哪里去、如何转换",由 Flink CDC CLI 负责解析并编译为可提交的 Flink 作业。相比传统需要编写 Java/Scala 代码的 DataStream API 方案,本文的演示全部在 Flink CDC CLI 中完成,无需安装 IDE,也无需编写任何业务代码。
一个完整的 Streaming ELT 作业具备以下能力,也是本文后续各节将要逐一验证的:
- 整库同步:通过正则匹配一次性同步一个数据库下的所有表;
- 表结构变更同步:上游 DDL(如
ALTER TABLE)自动透传到下游消息; - 分库分表同步:通过路由(route)规则将多张源表合并到同一个下游 Topic;
- 灵活的输出编排:支持分区策略、输出格式、表名到 Topic 映射等多种参数。
准备阶段
开始前需要一台安装了 Docker 的 Linux 或 macOS 电脑,并依次准备 Flink 集群与数据管道所需的容器组件。
准备 Flink Standalone 集群
下载 Flink 1.20.3(本文以 Flink 1.20.x 版本线为准),解压后进入目录并设置
FLINK_HOME:tar -zxvf flink-1.20.3-bin-scala_2.12.tgz export FLINK_HOME=$(pwd)/flink-1.20.3 cd flink-1.20.3在
conf/config.yaml配置文件中追加以下参数开启检查点,每隔 3 秒执行一次 checkpoint。检查点是 Flink 故障恢复与端到端一致性的基础,对于持续运行的 CDC 作业至关重要:execution: checkpointing: interval: 3s启动 Flink 集群:
./bin/start-cluster.sh
启动成功后,即可在 http://localhost:8081/ 访问到 Flink Web UI。多次执行start-cluster.sh可以拉起多个 TaskManager 以提升并行处理能力。
注意:如果 Flink 搭建在云服务器上,需要将
conf/config.yaml中的rest.bind-address和rest.address修改为0.0.0.0,再使用公网 IP 访问 Web UI。
准备 Docker 环境
创建一个docker-compose.yml文件并写入以下内容,用于拉起本教程所需的三个组件:
version: '2.1' services: Zookeeper: image: zookeeper:3.7.1 ports: - "2181:2181" environment: - ALLOW_ANONYMOUS_LOGIN=yes Kafka: image: bitnami/kafka:2.8.1 ports: - "9092:9092" - "9093:9093" environment: - ALLOW_PLAINTEXT_LISTENER=yes - KAFKA_LISTENERS=PLAINTEXT://:9092 - KAFKA_ADVERTISED_LISTENERS=PLAINTEXT://Kafka:9092 - KAFKA_ZOOKEEPER_CONNECT=Zookeeper:2181 MySQL: image: debezium/example-mysql:1.1 ports: - "3306:3306" environment: - MYSQL_ROOT_PASSWORD=123456 - MYSQL_USER=mysqluser - MYSQL_PASSWORD=mysqlpw说明:上述配置中的容器内网 IP 可通过
ifconfig命令确认,需要与容器网络环境对应。
该 Docker Compose 将拉起以下三个容器:
- MySQL:数据管道的源头,
debezium/example-mysql镜像默认开启 binlog,这是 CDC 读取变更日志的前提; - Kafka:数据管道的下游,用于接收并存储变更消息;
- Zookeeper:用于 Kafka 集群管理与协调。
在docker-compose.yml所在目录执行以下命令启动全部容器:
docker compose up -d该命令以 detached 模式后台启动所有容器,可通过docker ps观察容器是否正常启动。
在 MySQL 数据库中准备数据
进入 MySQL 容器:
docker compose exec MySQL mysql -uroot -p123456创建数据库
app_db以及orders、products、shipments三张表,并插入初始数据:-- 创建数据库 CREATE DATABASE app_db; USE app_db; -- 创建 orders 表 CREATE TABLE `orders` ( `id` INT NOT NULL, `price` DECIMAL(10,2) NOT NULL, PRIMARY KEY (`id`) ); -- 插入数据 INSERT INTO `orders` (`id`, `price`) VALUES (1, 4.00); INSERT INTO `orders` (`id`, `price`) VALUES (2, 100.00); -- 创建 shipments 表 CREATE TABLE `shipments` ( `id` INT NOT NULL, `city` VARCHAR(255) NOT NULL, PRIMARY KEY (`id`) ); -- 插入数据 INSERT INTO `shipments` (`id`, `city`) VALUES (1, 'beijing'); INSERT INTO `shipments` (`id`, `city`) VALUES (2, 'xian'); -- 创建 products 表 CREATE TABLE `products` ( `id` INT NOT NULL, `product` VARCHAR(255) NOT NULL, PRIMARY KEY (`id`) ); -- 插入数据 INSERT INTO `products` (`id`, `product`) VALUES (1, 'Beer'); INSERT INTO `products` (`id`, `product`) VALUES (2, 'Cap'); INSERT INTO `products` (`id`, `product`) VALUES (3, 'Peanut');
通过 Flink CDC CLI 提交任务
Flink CDC CLI 是整条链路的"编排器":它读取 YAML 任务配置,将其编译为可执行的 Flink 作业并提交到集群,其核心提交逻辑可在 CliExecutor.java 与 CliFrontend.java 中查看。
下载链接只对已发布的版本有效,SNAPSHOT 版本需要基于 master 或 release- 分支本地编译。
下载并安装 Flink CDC
下载
flink-cdc-<version>-bin.tar.gz二进制压缩包并解压,得到目录flink-cdc-<version>,该目录下包含bin、lib、log、conf四个目录。下载以下 connector 包并移动到Flink CDC Home 的
lib目录(注意:不是 Flink Home 的lib目录):flink-cdc-pipeline-connector-mysql:MySQL 数据源 connector;flink-cdc-pipeline-connector-kafka:Kafka 数据汇 connector。
此外,还需要将 MySQL Connector Java(JDBC Driver)放入 Flink
lib目录,或通过--jar参数传入 Flink CDC CLI,因为 CDC Connectors 本身不再内置这些 JDBC Drivers。--jar参数的定义见 CliFrontendOptions.java。
编写任务配置 YAML
下面是一个整库同步的示例文件mysql-to-kafka.yaml,它把app_db库下所有表同步到一个 Kafka Topic:
################################################################################ # Description: Sync MySQL all tables to Kafka ################################################################################ source: type: mysql hostname: 0.0.0.0 port: 3306 username: root password: 123456 tables: app_db.\.* server-id: 5400-5404 server-time-zone: UTC sink: type: kafka name: Kafka Sink properties.bootstrap.servers: 0.0.0.0:9092 topic: yaml-mysql-kafka pipeline: name: MySQL to Kafka Pipeline parallelism: 1关键参数说明(可对照 MySqlDataSourceOptions.java 与 KafkaDataSinkOptions.java 源码核验):
| 配置项 | 说明 | 默认值/取值 |
|---|---|---|
source.tables | 待监听的 MySQL 表名,支持正则表达式。app_db.\.*即匹配app_db下所有表;.在正则中需要以\.转义,因为它同时被用作库名与表名的分隔符 | 无默认值 |
source.server-id | 数据库客户端的数字 ID 或 ID 范围(如5400或5400-5408)。开启增量快照时推荐使用范围语法,且每个 ID 必须在当前 MySQL 集群的所有运行进程中唯一,连接器会以此身份加入集群读取 binlog。默认在 5400~6400 之间随机生成,官方建议显式指定 | 随机生成 |
source.server-time-zone | 数据库服务器的会话时区,不设置时使用系统默认时区(ZoneId.systemDefault()) | 系统默认时区 |
source.port | MySQL 端口 | 3306 |
sink.topic | 配置后所有事件都将发送到该 Topic | 无默认值 |
sink.properties.bootstrap.servers | Kafka 集群地址,properties.前缀用于透传 Kafka producer 原生配置 | — |
pipeline.parallelism | 管道整体并行度 | — |
提交任务到集群
编写完 YAML 后,执行以下命令将任务提交到 Flink Standalone 集群:
bash bin/flink-cdc.sh mysql-to-kafka.yaml如果作业提交成功,将打印以下信息(提交成功提示的输出逻辑见 CliFrontend.java 中的Pipeline has been submitted to cluster.):
Pipeline has been submitted to cluster. Job ID: 04fd88ccb96c789dce2bf0b3a541d626 Job Description: MySQL to Kafka Pipeline在 Flink Web UI 中,可以看到一个名为Sync MySQL Database to Kafka的任务正在运行:
验证 Kafka 中的消息
使用 Kafka 命令行工具查看 Topic 中的消息:
docker compose exec Kafka kafka-console-consumer.sh --bootstrap-server 0.0.0.0:9092 --topic yaml-mysql-kafka --from-beginning默认输出为debezium-json格式,每条消息包含before、after、op、source字段。其中op的取值在 DebeziumJsonSerializationSchema.java 中定义为:c(insert)、u(update)、d(delete)。示例消息如下:
{ "before": null, "after": { "id": 1, "price": 4 }, "op": "c", "source": { "db": "app_db", "table": "orders" } } // ... { "before": null, "after": { "id": 1, "product": "Beer" }, "op": "c", "source": { "db": "app_db", "table": "products" } } // ... { "before": null, "after": { "id": 2, "city": "xian" }, "op": "c", "source": { "db": "app_db", "table": "shipments" } }可以看到,三张表的变更数据都被实时同步到了同一个 Topic 中,source字段记录了消息来自哪个库、哪张表。
同步变更:DML 与 DDL 实时透传
使用以下命令进入 MySQL 容器:
docker compose exec mysql mysql -uroot -p123456接下来对orders表依次执行插入、加列(DDL)、更新、删除操作,Kafka 中的消息将实时更新:
插入一条数据:
INSERT INTO app_db.orders (id, price) VALUES (3, 100.00);增加一个字段(表结构变更同步):
ALTER TABLE app_db.orders ADD amount varchar(100) NULL;更新一条数据:
UPDATE app_db.orders SET price=100.00, amount=100.00 WHERE id=1;删除一条数据:
DELETE FROM app_db.orders WHERE id=2;
使用上一节提到的消费者命令观察下游 Topic,可以看到实时的变更消息记录。例如 UPDATE 操作对应的消息为:
{ "before": { "id": 1, "price": 4, "amount": null }, "after": { "id": 1, "price": 100, "amount": "100.00" }, "op": "u", "source": { "db": "app_db", "table": "orders" } }before记录了变更前的完整行数据,after记录了变更后的数据,op: "u"表明这是一次更新。值得注意的是,after中出现了amount字段,说明此前ALTER TABLE引起的表结构变更也已经被同步——这正是 Streaming ELT 相比传统"定时抽数"的关键优势之一。同样地,修改shipments、products表,也能在 Kafka 中实时看到同步变更结果。
路由变更:表名/库名替换与分库分表同步
Flink CDC 提供了将源表的表结构/数据路由到其他表名的能力,借助 route 规则可以实现表名、库名替换以及整库分发。路由定义在 composer 层的 RouteDef.java 中,由 FlinkPipelineComposer.java 在编译阶段应用。
下面是一个带 route 规则的示例配置:
################################################################################ # Description: Sync MySQL all tables to Kafka ################################################################################ source: type: mysql hostname: localhost port: 3306 username: root password: 123456 tables: app_db.\.* server-id: 5400-5404 server-time-zone: UTC sink: type: kafka name: Kafka Sink properties.bootstrap.servers: 0.0.0.0:9092 pipeline: name: MySQL to Kafka Pipeline parallelism: 1 route: - source-table: app_db.orders sink-table: kafka_ods_orders - source-table: app_db.shipments sink-table: kafka_ods_shipments - source-table: app_db.products sink-table: kafka_ods_products以上三条 route 规则将app_db中三张表的结构和数据分别分发到三个不同的下游 Topic(kafka_ods_orders、kafka_ods_shipments、kafka_ods_products)。
特别地,source-table支持正则表达式匹配多表,从而实现分库分表同步。例如下面的规则会将app_db库中所有以order开头的表合并成一张表,统一发送到kafka_ods_ordersTopic:
route: - source-table: app_db.order\.* sink-table: kafka_ods_orders使用 Kafka 命令查看新建的 Topic:
docker compose exec Kafka kafka-topics.sh --bootstrap-server 0.0.0.0:9092 --list可以看到新创建的 Kafka Topic 列表:
__consumer_offsetskafka_ods_orderskafka_ods_productskafka_ods_shipmentsyaml-mysql-kafka
选取kafka_ods_ordersTopic 查询,返回数据示例如下。注意此时source.table已经变为路由后的目标表名kafka_ods_orders:
{ "before": null, "after": { "id": 1, "price": 100, "amount": "100.00" }, "op": "c", "source": { "db": null, "table": "kafka_ods_orders" } }写入多个分区:partition.strategy 分区策略
Kafka 的吞吐优势依赖于多分区并行消费。partition.strategy参数用于定义消息发送到 Kafka 分区的策略,可选项在 PartitionStrategy.java 中定义,其默认值声明于 KafkaDataSinkOptions.java:
all-to-zero:将所有数据发送到 0 号分区(默认值);hash-by-key:所有数据根据主键的哈希值分发,可保证同一主键的变更消息落在同一分区,便于下游按主键顺序消费。
在mysql-to-kafka.yaml的sink块中增加partition.strategy: hash-by-key配置:
source: # ... sink: # ... topic: yaml-mysql-kafka-hash-by-key partition.strategy: hash-by-key pipeline: # ...同时在 Kafka 中新建一个 12 分区的 Topic:
docker compose exec Kafka kafka-topics.sh --create --topic yaml-mysql-kafka-hash-by-key --bootstrap-server 0.0.0.0:9092 --partitions 12查看指定分区的数据:
docker compose exec Kafka kafka-console-consumer.sh --bootstrap-server=0.0.0.0:9092 --topic yaml-mysql-kafka-hash-by-key --partition 0 --from-beginning部分分区数据详情如下,可见不同主键的数据被分散到了不同分区(0 号分区、4 号分区):
// 分区 0 { "before": null, "after": { "id": 1, "price": 100, "amount": "100.00" }, "op": "c", "source": { "db": "app_db", "table": "orders" } } // 分区 4 { "before": null, "after": { "id": 2, "product": "Cap" }, "op": "c", "source": { "db": "app_db", "table": "products" } } { "before": null, "after": { "id": 1, "city": "beijing" }, "op": "c", "source": { "db": "app_db", "table": "shipments" } }输出格式:value.format 参数
value.format参数用于指定 Kafka 消息 value 部分的序列化格式,可选值定义在 JsonSerializationType.java 中:
debezium-json(默认值):包含before(变更前数据)、after(变更后数据)、op(变更类型)、source(元数据)字段;canal-json:包含old、data、type、database、table、pkNames字段。
目前不支持用户自定义输出格式。另外,ts_ms字段默认不会包含在输出结构中(需要与 MySQL Source 的metadata.list参数配合使用)。
在 YAML 的 sink 块中添加value.format: canal-json指定输出为 Canal JSON:
source: # ... sink: # ... topic: yaml-mysql-kafka-canal value.format: canal-json pipeline: # ...查询对应 Topic 的数据,返回示例如下:
{ "old": null, "data": [ { "id": 1, "price": 100, "amount": "100.00" } ], "type": "INSERT", "database": "app_db", "table": "orders", "pkNames": [ "id" ] }两种格式面向不同的下游生态:debezium-json与 Debezium 生态兼容,canal-json与 Canal 生态兼容,可按下游消费方的解析能力选择。
表名到 Topic 的映射关系:sink.tableId-to-topic.mapping
partition.strategy控制的是分区维度,而sink.tableId-to-topic.mapping控制的是表到 Topic的维度。该参数用于指定上游表名到下游 Kafka Topic 名的映射关系,无需使用 route 配置。
与 route 方式的关键区别在于:配置该参数可以在保留上游各源表的表名和表结构的同时,将来自特定表的数据派发到相应的 Kafka Topic。
映射语法:每组映射关系由;分割,上游表的 TableId(支持正则表达式匹配)与下游 Kafka Topic 名由:分割。该解析逻辑与分隔符常量;、:均定义在 KafkaDataSinkOptions.java 中。
在前面的 YAML 文件中增加sink.tableId-to-topic.mapping配置:
source: # ... sink: # ... sink.tableId-to-topic.mapping: app_db.orders:yaml-mysql-kafka-orders;app_db.shipments:yaml-mysql-kafka-shipments;app_db.products:yaml-mysql-kafka-products pipeline: # ...运行后,Kafka 中将会创建以下 Topic:
yaml-mysql-kafka-ordersyaml-mysql-kafka-productsyaml-mysql-kafka-shipments
各 Topic 中的消息示例:
yaml-mysql-kafka-orders{ "before": null, "after": { "id": 1, "price": 100, "amount": "100.00" }, "op": "c", "source": { "db": "app_db", "table": "orders" } }yaml-mysql-kafka-products{ "before": null, "after": { "id": 2, "product": "Cap" }, "op": "c", "source": { "db": "app_db", "table": "products" } }yaml-mysql-kafka-shipments{ "before": null, "after": { "id": 2, "city": "xian" }, "op": "c", "source": { "db": "app_db", "table": "shipments" } }
与 route 方式不同,这里的source.table仍然保留原始表名(如orders),source.db也保留了库名信息,适合下游希望同时拿到"原始表标识 + 独立 Topic"的场景。
环境清理
实验结束后,在docker-compose.yml文件所在目录执行以下命令停止所有容器:
docker compose down在 Flink 所在目录flink-1.20.3下执行以下命令停止 Flink 集群:
./bin/stop-cluster.sh小结与源码速查
本文完整演示了基于 Flink CDC 的 MySQL → Kafka Streaming ELT 管道:从 Flink Standalone 集群与 Docker 环境准备,到 YAML 声明式任务配置、CLI 零代码提交,再到整库同步、DDL/DML 实时透传、路由分发、分区策略、输出格式与表名映射。整套链路无需编写一行业务代码,非常适合作为实时数仓 ODS 层建设的起步方案。
如果需要深入源码验证或二次开发,可在当前仓库中重点关注以下文件:
- 参数定义:KafkaDataSinkOptions.java、MySqlDataSourceOptions.java
- 分区策略:PartitionStrategy.java
- 输出格式:JsonSerializationType.java、DebeziumJsonSerializationSchema.java
- 路由与编译:RouteDef.java、FlinkPipelineComposer.java
- 任务提交:CliExecutor.java、CliFrontend.java
【免费下载链接】flink-cdcFlink CDC is a streaming data integration tool项目地址: https://gitcode.com/GitHub_Trending/flin/flink-cdc
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考