Flink CDC MySQL 到 Kafka 实战:全库同步、Schema 变更与表路由(Flink 1.20 快速入门)
2026/9/17 2:53:02 网站建设 项目流程

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.strategyvalue.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 集群

  1. 下载 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
  2. conf/config.yaml末尾追加以下配置,启用每 3 秒一次 checkpoint(CDC 管道依赖 checkpoint 进行状态一致性与精确写入):

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

    ./bin/start-cluster.sh

    集群启动后可访问 Web UI(默认http://localhost:8081/)确认 JobManager 与 TaskManager 状态。

如需更多 TaskManager,可多次执行start-cluster.sh

注意:如果将 Flink 集群部署为云服务,需要在conf/config.yaml中把rest.bind-addressrest.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 测试数据

  1. 进入 MySQL 容器:

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

    -- 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 分支构建。

  1. 下载 Flink CDC 二进制包并解压,得到包含binliblogconf四个目录的flink-cdc-<version>目录(发行包结构可参考仓库中 flink-cdc-dist 的 assembly 定义)。

  2. 将两个 Pipeline 连接器 jar 放入Flink CDC Home 的lib目录(注意:不是 Flink Home 的 lib 目录):

    • flink-cdc-pipeline-connector-mysql
    • flink-cdc-pipeline-connector-kafka

    同时,由于 MySQL JDBC 驱动不再随 CDC 连接器打包,还需要将 MySQL Connector Java(如 8.0.27 版本)放入 Flink 的lib目录,或通过 CLI 的--jar参数传入。

  3. 编写 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定义)。
  4. 在 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 Pipeline

    Flink 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 格式会编码beforeafteropsource等字段,示例如下:

{ "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 消费者端可以观察到对应的更新消息(opubefore/after分别记录变更前后的行内容):

{ "before": { "id": 1, "price": 4, "amount": null }, "after": { "id": 1, "price": 100, "amount": "100.00" }, "op": "u", "source": { "db": "app_db", "table": "orders" } }

类似地,修改shipmentsproducts表,也能在对应的 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.order01app_db.order02app_db.order03等分片表合并同步到一个kafka_ods_ordersTopic。从源码看,YAML 中route段的解析在 YamlPipelineDefinitionParser.toRouteDef,要求每条规则必须包含source-tablesink-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(默认):编码beforeafteropsource字段为 JSON;
  • canal-json:编码olddatatypedatabasetablepkNames字段为 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-orders
  • yaml-mysql-kafka-products
  • yaml-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-guaranteeAT_LEAST_ONCE(默认)提交时的投递语义保证
partition.strategyall-to-zero(默认)/hash-by-key数据写入 Kafka 分区的策略
key.formatjson(默认)/csv编码 Kafka 消息 key 的格式
value.formatdebezium-json(默认)/canal-json编码消息 value 的 JSON 格式
topic无默认值配置后所有事件统一写入该 Topic
sink.add-tableId-to-header-enabledfalse(默认)开启后为每条 Kafka 记录添加namespaceschemaNametableName
sink.custom-header空(默认)为每条记录添加自定义头,格式key1:value1,key2:value2
sink.tableId-to-topic.mapping无默认值上游表到 Kafka Topic 的映射规则,;分隔多条、:分隔规则
debezium-json.include-schema.enabledfalse(默认)开启后每条 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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询