news 2026/9/17 2:16:26

基于 Flink CDC 构建 MySQL 到 Kafka 的 Streaming ELT 整库同步管道

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
基于 Flink CDC 构建 MySQL 到 Kafka 的 Streaming ELT 整库同步管道

基于 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 集群

  1. 下载 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
  2. conf/config.yaml配置文件中追加以下参数开启检查点,每隔 3 秒执行一次 checkpoint。检查点是 Flink 故障恢复与端到端一致性的基础,对于持续运行的 CDC 作业至关重要:

    execution: checkpointing: interval: 3s
  3. 启动 Flink 集群:

    ./bin/start-cluster.sh

启动成功后,即可在 http://localhost:8081/ 访问到 Flink Web UI。多次执行start-cluster.sh可以拉起多个 TaskManager 以提升并行处理能力。

注意:如果 Flink 搭建在云服务器上,需要将conf/config.yaml中的rest.bind-addressrest.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 数据库中准备数据

  1. 进入 MySQL 容器:

    docker compose exec MySQL mysql -uroot -p123456
  2. 创建数据库app_db以及ordersproductsshipments三张表,并插入初始数据:

    -- 创建数据库 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

  1. 下载flink-cdc-<version>-bin.tar.gz二进制压缩包并解压,得到目录flink-cdc-<version>,该目录下包含binliblogconf四个目录。

  2. 下载以下 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)放入 Flinklib目录,或通过--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 范围(如54005400-5408)。开启增量快照时推荐使用范围语法,且每个 ID 必须在当前 MySQL 集群的所有运行进程中唯一,连接器会以此身份加入集群读取 binlog。默认在 5400~6400 之间随机生成,官方建议显式指定随机生成
source.server-time-zone数据库服务器的会话时区,不设置时使用系统默认时区(ZoneId.systemDefault()系统默认时区
source.portMySQL 端口3306
sink.topic配置后所有事件都将发送到该 Topic无默认值
sink.properties.bootstrap.serversKafka 集群地址,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格式,每条消息包含beforeafteropsource字段。其中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 中的消息将实时更新:

  1. 插入一条数据:

    INSERT INTO app_db.orders (id, price) VALUES (3, 100.00);
  2. 增加一个字段(表结构变更同步):

    ALTER TABLE app_db.orders ADD amount varchar(100) NULL;
  3. 更新一条数据:

    UPDATE app_db.orders SET price=100.00, amount=100.00 WHERE id=1;
  4. 删除一条数据:

    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 相比传统"定时抽数"的关键优势之一。同样地,修改shipmentsproducts表,也能在 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_orderskafka_ods_shipmentskafka_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_offsets
  • kafka_ods_orders
  • kafka_ods_products
  • kafka_ods_shipments
  • yaml-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.yamlsink块中增加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:包含olddatatypedatabasetablepkNames字段。

目前不支持用户自定义输出格式。另外,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-orders
  • yaml-mysql-kafka-products
  • yaml-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),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/17 2:15:35

MATLAB读取SAC地震数据:rdsac.m脚本实现与实战

简介&#xff1a;一个用于处理 SAC 格式地震数据的 MATLAB 脚本包&#xff0c;面向地震学、地球物理学领域的科研人员与技术人员&#xff0c;解决在 MATLAB 环境中直接读取和分析 SAC 文件的需求。压缩包内共 1 个文件&#xff0c;为 rdsac.m 脚本&#xff0c;体积仅 1KB。该脚…

作者头像 李华
网站建设 2026/9/17 2:14:38

嵌入式面试I2C/SPI高频考点:从协议原理到调试实战

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/9/17 2:11:36

Hadoop实战:基于MapReduce的豆瓣电影数据分析与可视化

简介&#xff1a;Hadoop豆瓣电影分析可视化源码是一套面向大数据课程实验的完整项目参考&#xff0c;围绕豆瓣电影Top250榜单数据&#xff0c;模拟真实大数据分析场景&#xff0c;适用于本科及高职大数据专业的课程设计、Hive案例实践与毕业设计参考。项目需要搭建Hadoop集群&a…

作者头像 李华
网站建设 2026/9/17 2:07:07

从零构建第一个机器学习模型:Scikit-learn完整实操指南

这几周后台一直有人在问&#xff0c;说想系统学机器学习&#xff0c;但看到各种深度学习框架的入门教程就头皮发麻&#xff0c;问我有没有更温和的切入点。其实答案一直都很明确&#xff1a;从Scikit-learn开始&#xff0c;用它构建你的第一个机器学习模型。这个库足够简单、足…

作者头像 李华