Flink CDC 实战:从零搭建 MySQL 到 Kafka 的流式 ELT 管道(Flink 2.2 快速上手) Flink CDC 实战从零搭建 MySQL 到 Kafka 的流式 ELT 管道Flink 2.2 快速上手【免费下载链接】flink-cdcFlink CDC is a streaming data integration tool项目地址: https://gitcode.com/GitHub_Trending/flin/flink-cdc本文基于 Flink CDC 官方文档中的 MySQL to Kafka 快速上手教程适配 Flink 2.2 集群完整展开带你用 Docker Compose 一键拉起 MySQL、Kafka 环境通过 Flink CDC CLI 提交纯 YAML 定义的管道任务零 Java 代码实现整库同步、Schema 变更演进、分库分表合并以及多分区写入。读完本文你将掌握 MySQL CDC 到 Kafka 的完整链路配置方法并能看懂管道 Sink 侧各参数partition.strategy、value.format、sink.tableId-to-topic.mapping等在源码中的定义、默认值与生效机制。一、环境准备整套练习在一个装了 Docker 的 Linux 或 macOS 机器上进行分为两部分一个 Flink 2.2 独立集群以及由 Docker Compose 拉起的 MySQL、Kafka、Zookeeper 容器。1.1 准备 Flink 2.2 独立集群从 Apache Flink 官方归档站点下载 Flink 2.2.0 二进制包解压后进入flink-2.2.0目录tar -zxvf flink-2.2.0-bin-scala_2.12.tgz export FLINK_HOME$(pwd)/flink-2.2.0 cd flink-2.2.0打开 Checkpoint 机制。CDC 管道依赖 Flink 的 Checkpoint 保证 Exactly-Once 状态一致性因此在conf/config.yaml中追加每 3 秒一次 Checkpoint 的配置execution: checkpointing: interval: 3s启动集群./bin/start-cluster.sh启动后可访问http://localhost:8081/查看 Flink Web UI。如需更多 TaskManager重复执行start-cluster.sh即可。注意如果集群部署在云环境需要将conf/config.yaml中的rest.bind-address和rest.address配置为0.0.0.0再使用公网 IP 访问 Web UI。1.2 准备 Docker Compose 环境创建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三个容器分工明确MySQL管道数据源使用官方 CDC 示例镜像已开启 BinlogKafka管道目标Sink这里选用带 Zookeeper 的 2.8.1 镜像Zookeeper负责 Kafka 集群的协调管理。在docker-compose.yml所在目录执行docker compose up -d该命令以分离模式启动所有容器随后用docker ps确认三个容器均在运行。1.3 准备 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 提交管道任务Flink CDC 提供了独立的可执行发行包内置bin/flink-cdc.sh命令行入口对应仓库中的 flink-cdc-dist 装配目录你不需要写一行 Java/Scala 代码也不需要 IDE。2.1 下载与依赖准备注意官方下载链接仅提供稳定版本如需 SNAPSHOT 版本须自行基于 master 或 release 分支构建。从 Apache 官方发布渠道下载flink-cdc-x.x.x-bin.tar.gz并解压到flink-cdc-x.x.x目录其中包含bin、lib、log、conf四个目录。从 Maven 中央仓库下载两个 Pipeline Connector 的 jar并移动到Flink CDC Home 的lib目录注意不是 Flink Home 的 lib 目录flink-cdc-pipeline-connector-mysqlflink-cdc-pipeline-connector-kafka由于 JDBC 驱动不再随 CDC 连接器打包还需将 MySQL Connector Java8.0.x 版本放入 Flink 的lib目录或通过 CLI 的--jar参数传入。从源码结构看连接器的发现机制依赖 SPIKafka Sink 工厂 KafkaDataSinkFactory 声明了identifier() kafka正是 YAML 中sink.type: kafka的匹配依据。2.2 编写管道定义文件创建mysql-to-kafka.yaml把 MySQL 整库同步到 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几个关键配置说明tables: app_db.\.*用正则匹配app_db库下所有表实现整库同步server-id: 5400-5404MySQL 要求每个从库有唯一 server-id区间上限需覆盖任务并行度此处并行度为 1properties.bootstrap.servers所有带properties.前缀的配置会被透传给底层 Kafka Producer。这一点在 KafkaDataSinkFactory 中可以看到工厂会遍历全部配置项把以properties.为前缀的键值剥离前缀后放入Properties最终由 KafkaDataSink.getEventSinkProvider() 通过KafkaSinkBuilder.setKafkaProducerConfig(...)构建出真正的 Kafka Sink并设置投递语义默认AT_LEAST_ONCE。topic不配置sink.tableId-to-topic.mapping时所有表的事件统一写入该 Topic。2.3 提交任务并验证在 Flink CDC Home 目录下执行bash bin/flink-cdc.sh mysql-to-kafka.yaml提交成功后控制台输出Pipeline has been submitted to cluster. Job ID: 04fd88ccb96c789dce2bf0b3a541d626 Job Description: MySQL to Kafka Pipeline订阅目标 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字段示例输出{ before: null, after: { id: 1, price: 4 }, op: c, source: { db: app_db, table: orders } } // ... { before: null, after: { id: 2, city: xian }, op: c, source: { db: app_db, table: shipments } }全量阶段完成后MySQL 三张表的初始数据都以op: ccreate事件写入 Kafka且消息中source.db/source.table保留了源库表信息便于下游按表区分。三、同步 Schema 变更与数据变更管道任务运行期间在 MySQL 中依次执行增列、插入、更新、删除操作-- 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 price100.00, amount100.00 WHERE id1; -- 4. 删除一条记录 DELETE FROM app_db.orders WHERE id2;消费者侧可以实时看到对应事件。其中 UPDATE 事件的before/after会同时携带旧值与新值且新增列amount出现在字段列表中——这正是 Schema 变更演进schema change evolution能力的体现Sink 侧会感知源端的ALTER TABLE事件并同步调整输出结构。{ 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 中看到对应的实时同步结果。四、route 路由重命名与分表合并Flink CDC 的route配置可以把源表的 Schema 和数据变更路由到其他表名从而实现表名/库名替换、整库同步到自定义命名空间等场景。完整示例################################################################################ # 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注意此示例中sink没有配置topic当使用 route 规则时路由后的sink-table名称即决定数据写入的 Kafka Topic。分表合并场景source-table支持正则匹配多张表匹配到的表结构会被合并后同步到同一个 Topicroute: - source-table: app_db.order\.* sink-table: kafka_ods_orders这样app_db.order01、app_db.order02、app_db.order03等分片表的结构将被合并统一写入kafka_ods_orders一个 Topic——这是分库分表数据归并的典型用法。验证 Topic 是否创建成功docker compose exec Kafka kafka-topics.sh --bootstrap-server 0.0.0.0:9092 --list输出* __consumer_offsets * kafka_ods_orders * kafka_ods_products * kafka_ods_shipments * yaml-mysql-kafka查看kafka_ods_orders中的消息可以看到source.table已被改写为路由后的表名{ before: null, after: { id: 1, price: 100, amount: 100.00 }, op: c, source: { db: null, table: kafka_ods_orders } }五、tableId-to-topic 映射不改表名、只分 Topic与 route 不同sink.tableId-to-topic.mapping只做“上游表到 Topic”的分发不会合并或改写表结构事件中的 TableId 保持原样。映射规则以;分隔每条规则由“上游 Table ID支持正则: 目标 Topic 名”组成中间用:分隔。该选项的常量定义可见 KafkaDataSinkOptionsDELIMITER_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: # ...配置后会创建三个 Topicyaml-mysql-kafka-ordersyaml-mysql-kafka-productsyaml-mysql-kafka-shipments各 Topic 内只包含对应表的记录且source.db/source.table保持源端原值{ before: null, after: { id: 1, price: 100, amount: 100.00 }, op: c, source: { db: app_db, table: orders } }route 与 tableId-to-topic.mapping 的选型建议需要改表名、合并分表结构时用route只希望把不同表分发到不同 Topic 且保留原始表结构时用sink.tableId-to-topic.mapping。六、多分区写入partition.strategypartition.strategy用于配置数据写入 Kafka 分区的策略源码中的枚举定义见 PartitionStrategy与选项声明默认值all-to-zero见 KafkaDataSinkOptions取值行为all-to-zero默认所有数据发送到 Partition 0保证单分区内顺序hash-by-key按主键的哈希值把数据变更分发到不同分区提升并行写入吞吐例如在 sink 中追加配置sink: # ... topic: yaml-mysql-kafka-hash-by-key partition.strategy: hash-by-key并手动创建一个 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可以看到 Partition 0 收到orders表记录Partition 4 收到products和shipments表记录说明按主键哈希的分区策略生效。选择hash-by-key时需注意同一主键始终落在同一分区因此同键有序性依旧有保证而不同键可并行写入适用于对吞吐要求更高的场景。七、输出格式value.format 与 canal-jsonvalue.format决定发往 Kafka 的 JSON 序列化格式选项声明与默认值见 KafkaDataSinkOptions.VALUE_FORMAT支持两种取值格式编码字段debezium-json默认before、after、op、sourcecanal-jsonold、data、type、database、table、pkNames对应实现分别位于 DebeziumJsonSerializationSchema 与 CanalJsonSerializationSchema由 ChangeLogJsonFormatFactory 根据配置值创建。目前不支持用户自定义输出格式。canal-json格式的示例输出{ old: null, data: [ { id: 1, price: 100, amount: 100.00 } ], type: INSERT, database: app_db, table: orders, pkNames: [ id ] }另外两点值得注意输出结构中默认不包含ts_ms字段如需暴露额外元数据字段需设置 MySQL Source 的metadata.list选项若希望每条 debezium 记录附带完整 Schema 信息可开启debezium-json.include-schema.enabled默认false仅对 debezium-json 格式生效同样定义在 KafkaDataSinkOptions。除文档重点提到的选项外从 KafkaDataSinkFactory.optionalOptions() 还可以确认 Kafka Sink 完整支持的选项集合key.formatcsv/json默认json、value.format、partition.strategy、topic、sink.add-tableId-to-header-enabled在每条记录 Header 中附加namespace/schemaName/tableName、sink.custom-header自定义记录 Header格式key1:value1,key2:value2、sink.delivery-guarantee默认AT_LEAST_ONCE以及sink.tableId-to-topic.mapping。八、清理环境教程结束后在docker-compose.yml所在目录停止并移除所有容器docker compose down在FLINK_HOME下停止 Flink 集群./bin/stop-cluster.sh小结本文围绕 Flink CDC 面向 Flink 2.2 的 MySQL to Kafka 快速上手教程完整走通了Docker Compose 搭建 MySQL/Kafka 环境 → 编写 YAML 管道定义 → 通过flink-cdc.shCLI 提交整库同步任务 → 验证全量/增量数据与 Schema 变更同步 → 使用route合并分表、使用sink.tableId-to-topic.mapping分 Topic 分发 → 配置partition.strategy多分区写入 → 切换value.format输出格式。所有关键参数默认值、取值范围、前缀透传规则均可在 KafkaDataSinkOptions 与 KafkaDataSinkFactory 源码中逐一对照配合 Fluss、StarRocks、Doris 等官方快速上手示例 可扩展出更多目标端的同构管道。【免费下载链接】flink-cdcFlink CDC is a streaming data integration tool项目地址: https://gitcode.com/GitHub_Trending/flin/flink-cdc创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考