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

发布时间:2026/9/17 2:16:32
基于 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 集群下载 Flink 1.20.3本文以 Flink 1.20.x 版本线为准解压后进入目录并设置FLINK_HOMEtar -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_LOGINyes Kafka: image: bitnami/kafka:2.8.1 ports: - 9092:9092 - 9093:9093 environment: - ALLOW_PLAINTEXT_LISTENERyes - KAFKA_LISTENERSPLAINTEXT://:9092 - KAFKA_ADVERTISED_LISTENERSPLAINTEXT://Kafka:9092 - KAFKA_ZOOKEEPER_CONNECTZookeeper:2181 MySQL: image: debezium/example-mysql:1.1 ports: - 3306:3306 environment: - MYSQL_ROOT_PASSWORD123456 - MYSQL_USERmysqluser - MYSQL_PASSWORDmysqlpw说明上述配置中的容器内网 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-mysqlMySQL 数据源 connectorflink-cdc-pipeline-connector-kafkaKafka 数据汇 connector。此外还需要将 MySQL Connector JavaJDBC 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 范围如5400或5400-5408。开启增量快照时推荐使用范围语法且每个 ID 必须在当前 MySQL 集群的所有运行进程中唯一连接器会以此身份加入集群读取 binlog。默认在 5400~6400 之间随机生成官方建议显式指定随机生成source.server-time-zone数据库服务器的会话时区不设置时使用系统默认时区ZoneId.systemDefault()系统默认时区source.portMySQL 端口3306sink.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格式每条消息包含before、after、op、source字段。其中op的取值在 DebeziumJsonSerializationSchema.java 中定义为cinsert、uupdate、ddelete。示例消息如下{ 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 price100.00, amount100.00 WHERE id1;删除一条数据DELETE FROM app_db.orders WHERE id2;使用上一节提到的消费者命令观察下游 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中三张表的结构和数据分别分发到三个不同的下游 Topickafka_ods_orders、kafka_ods_shipments、kafka_ods_products。特别地source-table支持正则表达式匹配多表从而实现分库分表同步。例如下面的规则会将app_db库中所有以order开头的表合并成一张表统一发送到kafka_ods_ordersTopicroute: - source-table: app_db.order\.* sink-table: kafka_ods_orders使用 Kafka 命令查看新建的 Topicdocker 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.javaall-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 分区的 Topicdocker 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-server0.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 JSONsource: # ... 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.mappingpartition.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 中将会创建以下 Topicyaml-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仍然保留原始表名如orderssource.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),仅供参考