Flink CDC MySQL 到 Kafka 实战:全库同步、Schema 变更与表路由(Flink 1.20 快速入门)
【免费下载链接】flink-cdcFlink CDC is a streaming data integration tool项目地址: https://gitcode.com/GitHub_Trending/flin/flink-cdc
本文基于 Flink CDC 官方文档quickstart-for-1.20/mysql-to-kafka教程编写,面向希望在 Flink 1.20 上快速搭建「MySQL → Kafka」流式 ELT 管道的开发者。你将完整掌握:用 Docker Compose 准备 MySQL/Kafka/ZooKeeper 环境、用 Flink CDC CLI 以纯 YAML 提交同步任务(无需一行 Java/Scala 代码)、验证全库同步与 Schema 变更(schema evolution)实时同步、使用route规则做表名改写与分库分表合并、通过partition.strategy与value.format精细控制 Kafka 分区与消息格式、利用sink.tableId-to-topic.mapping实现表到 Topic 的精确映射。
环境准备
开始前需要一台安装了 Docker 的 Linux 或 macOS 计算机。整套环境由两部分构成:一个 Flink 1.20 Standalone 集群(运行 CDC 管道作业),以及一组 Docker 容器(MySQL 作为管道 Source、Kafka 作为管道 Sink、ZooKeeper 负责 Kafka 集群管理)。
准备 Flink 1.20 Standalone 集群
下载 Flink 1.20.3 发行包(
flink-1.20.3-bin-scala_2.12.tgz,可从 Apache 官方归档页获取),解压并进入目录,设置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(CDC 管道依赖 checkpoint 进行状态一致性与精确写入):execution: checkpointing: interval: 3s启动集群:
./bin/start-cluster.sh集群启动后可访问 Web UI(默认
http://localhost:8081/)确认 JobManager 与 TaskManager 状态。
如需更多 TaskManager,可多次执行start-cluster.sh。
注意:如果将 Flink 集群部署为云服务,需要在
conf/config.yaml中把rest.bind-address和rest.address配置为0.0.0.0,然后使用公网 IP 访问 Web UI。
准备 Docker Compose
创建docker-compose.yml,写入以下内容(版本以当前 Docker 环境为准,教程使用 ZooKeeper 3.7.1 + Kafka 2.8.1 + MySQL 5.7 示例镜像):
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该 Compose 文件准备三类容器:MySQL 作为管道 Source、Kafka 作为管道 Sink、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 CREATE DATABASE app_db; USE app_db; -- create orders table CREATE TABLE `orders` ( `id` INT NOT NULL, `price` DECIMAL(10,2) NOT NULL, PRIMARY KEY (`id`) ); -- insert records INSERT INTO `orders` (`id`, `price`) VALUES (1, 4.00); INSERT INTO `orders` (`id`, `price`) VALUES (2, 100.00); -- create shipments table CREATE TABLE `shipments` ( `id` INT NOT NULL, `city` VARCHAR(255) NOT NULL, PRIMARY KEY (`id`) ); -- insert records INSERT INTO `shipments` (`id`, `city`) VALUES (1, 'beijing'); INSERT INTO `shipments` (`id`, `city`) VALUES (2, 'xian'); -- create products table CREATE TABLE `products` ( `id` INT NOT NULL, `product` VARCHAR(255) NOT NULL, PRIMARY KEY (`id`) ); -- insert records 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 提交管道作业
说明:CDC 发行包仅针对稳定版本提供下载,SNAPSHOT 版本需要自行基于 master 或 release 分支构建。
下载 Flink CDC 二进制包并解压,得到包含
bin、lib、log、conf四个目录的flink-cdc-<version>目录(发行包结构可参考仓库中 flink-cdc-dist 的 assembly 定义)。将两个 Pipeline 连接器 jar 放入Flink CDC Home 的
lib目录(注意:不是 Flink Home 的 lib 目录):flink-cdc-pipeline-connector-mysqlflink-cdc-pipeline-connector-kafka
同时,由于 MySQL JDBC 驱动不再随 CDC 连接器打包,还需要将 MySQL Connector Java(如 8.0.27 版本)放入 Flink 的
lib目录,或通过 CLI 的--jar参数传入。编写 YAML 管道定义文件。下面是把 MySQL 全库同步到 Kafka 的示例
mysql-to-kafka.yaml:################################################################################ # 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要点说明:
source.tables: app_db.\.*用正则匹配该库下所有表,即"整库同步";server-id配置为范围值,需大于管道并行度,且不能与其他复制客户端冲突;- sink 段的
topic参数表示所有表事件统一写入该 Topic;properties.前缀的键值对会透传给底层 Kafka Producer(可参考 KafkaDataSinkOptions 中的PROPERTIES_PREFIX定义)。
在 Flink CDC Home 目录下用 CDC CLI 提交作业到 Standalone 集群:
bash bin/flink-cdc.sh mysql-to-kafka.yaml提交成功后控制台输出:
Pipeline has been submitted to cluster. Job ID: 04fd88ccb96c789dce2bf0b3a541d626 Job Description: MySQL to Kafka PipelineFlink Web UI 中可以看到该作业正在运行:
从源码结构看,这条命令的调用链是:CLI 入口 CliFrontend 解析命令行参数后交给 CliExecutor,后者根据部署目标(本例为remote)选择FlinkPipelineComposer.ofRemoteCluster,再由 YamlPipelineDefinitionParser 将 YAML 解析为PipelineDef(含 source/sink/route 等定义)并执行提交。
验证同步结果
订阅 sink Topic 即可观察写入 Kafka 的消息:
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等字段,示例如下:
{ "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" } }op字段中c表示 create(INSERT)。
实时同步 Schema 与数据变更
CDC 管道的核心价值之一是 schema evolution:上游 DDL 变更会被转换为 Schema 变更事件,管道在推进数据的同时自动演进下游。在 MySQL 容器中依次执行以下操作:
docker compose exec mysql mysql -uroot -p123456-- 1. 插入一条记录 INSERT INTO app_db.orders (id, price) VALUES (3, 100.00); -- 2. 给 orders 表新增一列 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;在 Kafka 消费者端可以观察到对应的更新消息(op为u,before/after分别记录变更前后的行内容):
{ "before": { "id": 1, "price": 4, "amount": null }, "after": { "id": 1, "price": 100, "amount": "100.00" }, "op": "u", "source": { "db": "app_db", "table": "orders" } }类似地,修改shipments和products表,也能在对应的 Kafka Topic 中实时看到同步结果。
从源码结构看,Kafka sink 对 Schema 变更事件的处理有明确边界:PipelineKafkaRecordSerializationSchema 的serialize方法中,遇到SchemaChangeEvent会直接返回 null,即不向 Kafka 发送 Schema 变更消息本身——因为 Kafka Topic 本身没有强 Schema 约束,"表结构演进"通过后续数据记录中自然携带新字段体现。而 sink 真正写 Kafka 的执行器由 KafkaDataSink 构建,它基于 Flink 官方的KafkaSink,支持可配置的DeliveryGuarantee(默认AT_LEAST_ONCE)。
使用 route 规则路由表结构与数据
Flink CDC 提供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中各表的 Schema 与数据变更会被同步到对应的 Kafka Topic。source-table还支持正则匹配多张表,多张表的 Schema 会被合并(merge)后同步到同一个 Kafka Topic:
route: - source-table: app_db.order\.* sink-table: kafka_ods_orders这样即可把app_db.order01、app_db.order02、app_db.order03等分片表合并同步到一个kafka_ods_ordersTopic。从源码看,YAML 中route段的解析在 YamlPipelineDefinitionParser.toRouteDef,要求每条规则必须包含source-table与sink-table字段,最终组装为RouteDef列表参与管道构图。
执行以下命令查看 Kafka 中新创建的 Topic:
docker compose exec Kafka kafka-topics.sh --bootstrap-server 0.0.0.0:9092 --list预期输出(含系统 Topic 与各规则生成的 Topic):
__consumer_offsets kafka_ods_orders kafka_ods_products kafka_ods_shipments yaml-mysql-kafka观察kafka_ods_ordersTopic 中的消息,可见记录中的source.table已是路由后的表名:
{ "before": null, "after": { "id": 1, "price": 100, "amount": "100.00" }, "op": "c", "source": { "db": null, "table": "kafka_ods_orders" } }使用 partition.strategy 控制分区写入
partition.strategy选项用于配置数据发送 Kafka 分区的策略,当前可用值见 PartitionStrategy 枚举:
all-to-zero:所有数据发送到 0 号分区(默认行为);hash-by-key:按主键的哈希值分发数据变更到不同分区。
对应地,KafkaDataSinkOptions 中该选项的默认值即为ALL_TO_ZERO;在 PipelineKafkaRecordSerializationSchema 中可以看到实现逻辑:策略为ALL_TO_ZERO时 partition 固定为 0,否则交由 Kafka 依据 key 哈希路由。
例如在 sink 段追加:
source: # ... sink: # ... topic: yaml-mysql-kafka-hash-by-key partition.strategy: hash-by-key pipeline: # ...并预创建一个 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可见不同表的数据被哈希分散到了不同分区:
// partition 0 { "before": null, "after": { "id": 1, "price": 100, "amount": "100.00" }, "op": "c", "source": { "db": "app_db", "table": "orders" } } // partition 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 的 JSON 序列化格式,当前不支持用户自定义格式。可用值在 JsonSerializationType 中定义,对应实现分别为 DebeziumJsonSerializationSchema 与 CanalJsonSerializationSchema:
debezium-json(默认):编码before、after、op、source字段为 JSON;canal-json:编码old、data、type、database、table、pkNames字段为 JSON;ts_ms字段默认不出现在输出结构中,需要配置 MySQL source 的metadata.list选项来暴露额外元数据字段。
canal-json格式的输出示例:
{ "old": null, "data": [ { "id": 1, "price": 100, "amount": "100.00" } ], "type": "INSERT", "database": "app_db", "table": "orders", "pkNames": [ "id" ] }使用 sink.tableId-to-topic.mapping 精确映射表到 Topic
sink.tableId-to-topic.mapping参数用于指定上游表到 Kafka Topic 的映射规则。与route规则的关键区别在于:表到 Topic 映射不会合并(merge)上游表的 Schema,各表的 TableId 保持原样不变,只是被分发到不同的 Topic。
规则以;分隔多条映射;每条映射由上游 Table ID(正则)与下游 Kafka Topic 名两部分组成,以:分隔。源码中两个分隔符分别定义在 KafkaDataSinkOptions 的DELIMITER_TABLE_MAPPINGS(;)与DELIMITER_SELECTOR_TOPIC(:)常量中。配置示例:
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: # ...配置后会自动创建以下 Topic:
yaml-mysql-kafka-ordersyaml-mysql-kafka-productsyaml-mysql-kafka-shipments
各 Topic 中的记录示例(注意source.table保留原始表名):
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" } }Kafka Sink 配置项速查
结合源码 KafkaDataSinkOptions 的定义,Kafka sink 的关键配置项及默认值汇总如下:
| 配置项 | 取值/默认值 | 说明 |
|---|---|---|
sink.delivery-guarantee | AT_LEAST_ONCE(默认) | 提交时的投递语义保证 |
partition.strategy | all-to-zero(默认)/hash-by-key | 数据写入 Kafka 分区的策略 |
key.format | json(默认)/csv | 编码 Kafka 消息 key 的格式 |
value.format | debezium-json(默认)/canal-json | 编码消息 value 的 JSON 格式 |
topic | 无默认值 | 配置后所有事件统一写入该 Topic |
sink.add-tableId-to-header-enabled | false(默认) | 开启后为每条 Kafka 记录添加namespace、schemaName、tableName头 |
sink.custom-header | 空(默认) | 为每条记录添加自定义头,格式key1:value1,key2:value2 |
sink.tableId-to-topic.mapping | 无默认值 | 上游表到 Kafka Topic 的映射规则,;分隔多条、:分隔规则 |
debezium-json.include-schema.enabled | false(默认) | 开启后每条 debezium 记录包含 Schema 信息,仅debezium-json格式支持 |
properties.* | 透传 Kafka Producer 配置 | 如properties.bootstrap.servers |
清理环境
教程结束后,在docker-compose.yml所在目录执行以下命令停止并移除所有容器:
docker compose down在FLINK_HOME下执行以下命令停止 Flink 集群:
./bin/stop-cluster.sh小结
本文完整复现了 Flink CDC 官方 MySQL→Kafka 快速入门教程(适用于 Flink 1.20):从 Flink Standalone 集群与 Docker 环境的搭建,到纯 YAML 提交整库同步管道,再到 Schema 变更实时同步、route表路由与分片表合并、partition.strategy分区策略、value.format输出格式以及sink.tableId-to-topic.mapping表到 Topic 映射。所有结论均可在仓库中对应的连接器源码(flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-kafka/)与 CLI 实现(flink-cdc-cli/)中逐一定位,便于进一步扩展为自己的生产管道。
【免费下载链接】flink-cdcFlink CDC is a streaming data integration tool项目地址: https://gitcode.com/GitHub_Trending/flin/flink-cdc
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考