Flink CDC 实战指南:在 Flink 1.20 上构建 MySQL 到 StarRocks 的流式 ELT 管道
【免费下载链接】flink-cdcFlink CDC is a streaming data integration tool项目地址: https://gitcode.com/GitHub_Trending/flin/flink-cdc
本篇以 Flink CDC(flink-cdc)官方教程为骨架,完整演示如何在 Flink 1.20 集群上,仅用 YAML 配置和 Flink CDC CLI,从零搭建一条 MySQL 到 StarRocks 的 Streaming ELT 管道:整库同步、实时 Schema 变更同步、分库分表路由合并。读完并动手跑通后,你将掌握 CDC Pipeline YAML 的完整结构、source/sink/route 各配置段的含义与源码中的参数定义、以及如何验证数据与表结构在两端实时一致。
一、环境准备
1. 准备 Flink 1.20 Standalone 集群
该教程适用于 Flink 1.20.x 运行时(对应仓库中的 flink-cdc-flink1-compat 兼容模块;仓库同时提供 Flink 2.2 的快速入门文档 quickstart-for-2.2)。准备一台已安装 Docker 的 Linux 或 macOS 机器,然后:
下载并解压 Flink 1.20.3 发行包,得到
flink-1.20.3目录,进入该目录并设置FLINK_HOME:cd flink-1.20.3开启 Checkpoint:向
conf/config.yaml追加以下配置,使作业每 3 秒做一次 Checkpoint(增量快照读取与 at-least-once 写入的正确性都依赖 Checkpoint 周期性推进):execution: checkpointing: interval: 3s启动集群:
./bin/start-cluster.sh启动成功后访问
http://localhost:8081/可看到 Flink Web UI。重复执行start-cluster.sh可以启动多个TaskManager,增加并行处理能力。
2. 用 Docker Compose 准备 MySQL 与 StarRocks
创建docker-compose.yml,包含两个服务:MySQL(内含app_db库)与 StarRocks(存储同步过来的表):
version: '2.1' services: StarRocks: image: starrocks/allin1-ubuntu:3.5.10 ports: - "8080:8080" - "9030:9030" MySQL: image: debezium/example-mysql:1.1 ports: - "3306:3306" environment: - MYSQL_ROOT_PASSWORD=123456 - MYSQL_USER=mysqluser - MYSQL_PASSWORD=mysqlpw在docker-compose.yml所在目录启动容器:
docker-compose up -d用docker ps确认容器状态;访问http://localhost:8030/可以确认 StarRocks 的 Web 界面(allin1 镜像将 FE/BE 合并在单容器内,8080 为 FE HTTP 端口、9030 为 FE 查询端口)。
3. 准备 MySQL 初始数据
进入 MySQL 容器:
docker-compose exec MySQL mysql -uroot -p123456创建app_db数据库以及orders、products、shipments三张带主键的表,并插入初始记录(注意:StarRocks sink 只支持主键表,所以源表必须带主键,这一点在 StarRocks 连接器文档的 Usage Notes 中有明确说明):
-- 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 提交管道作业
1. 准备 Flink CDC 发行包与连接器 JAR
下载 Flink CDC 稳定版二进制发行包(
flink-cdc-x.y.z-bin.tar.gz,从 Apache 官方发布渠道获取)并解压,得到包含bin、lib、log、conf四个目录的flink-cdc-x.y.z目录。将以下两个 Pipeline 连接器 JAR 下载到Flink CDC 主目录的
lib目录(注意:是 Flink CDC Home 的 lib,不是 Flink Home 的 lib):flink-cdc-pipeline-connector-mysqlflink-cdc-pipeline-connector-starrocks
稳定版 JAR 可从 Maven 公共仓库获取;如需 SNAPSHOT 版本,则需要自行基于 master 或 release 分支构建。
由于 MySQL JDBC 驱动不再随 CDC 连接器打包,还需将 MySQL Connector/J 8.x 的驱动 JAR 放入 Flink 的
lib目录,或通过提交时的--jar参数传入。
2. 编写 Pipeline 定义 YAML
以下是整库同步app_db到 StarRocks 的完整配置示例mysql-to-starrocks.yaml:
################################################################################ # Description: Sync MySQL all tables to StarRocks ################################################################################ source: type: mysql hostname: localhost port: 3306 username: root password: 123456 tables: app_db.\.* server-id: 5400-5404 server-time-zone: UTC sink: type: starrocks name: StarRocks Sink jdbc-url: jdbc:mysql://127.0.0.1:9030 load-url: 127.0.0.1:8080 username: root password: "" table.create.properties.replication_num: 1 pipeline: name: Sync MySQL Database to StarRocks parallelism: 2YAML 的顶层结构由 YamlPipelineDefinitionParser 解析校验:source与sink是必填块,route、transform、pipeline为可选块,出现未知顶层键会直接抛出带允许键列表的错误提示,便于尽早发现拼写错误。逐段解读:
source 段(由 MySQL pipeline 连接器消费,参数定义见 MySqlDataSourceOptions):
| 参数 | 说明 |
|---|---|
type | 固定为mysql,工厂据此通过 SPI 找到 MySQL Source |
hostname/port | MySQL 服务地址与端口(端口默认 3306) |
username/password | 连接 MySQL 的账号密码;该账号需要具备读取 binlog 的权限(REPLICATION SLAVE、REPLICATION CLIENT) |
tables | 需要同步的表,支持正则表达式。注意点号.是库名与表名的分隔符,正则中要匹配“任意字符的点”必须写成转义的\.,例如app_db.\.*表示同步app_db库下的全部表;多表可用逗号分隔,如db1.user_table_[0-9]+ |
server-id | 本作业伪装成的 MySQL 从库 ID。支持单值(5400)或范围(5400-5404),范围写法在增量快照模式下推荐,且必须与集群中其他正在运行的库进程互不重叠。源码注释也建议显式指定而不是用随机值 |
server-time-zone | MySQL 会话时区;不设置时使用系统默认时区。必须与实际 MySQL 服务时区一致,否则 binlog 时间戳会解析错乱 |
scan.startup.mode | 可选,默认initial(先读全量快照再读增量),其他可选值还有earliest-offset、latest-offset、timestamp、specific-offset、snapshot |
sink 段(参数定义见 StarRocksDataSinkOptions):
| 参数 | 必填 | 默认值 | 说明 |
|---|---|---|---|
type | 是 | - | 固定为starrocks |
name | 否 | - | sink 名称,会体现在作业的表描述里 |
jdbc-url | 是 | - | FE 的 MySQL 协议查询端口,如jdbc:mysql://127.0.0.1:9030;多 FE 用逗号分隔。用于建表、执行 schema 变更等 DDL |
load-url | 是 | - | FE 的 HTTP 端口,即 Stream Load 入口,如127.0.0.1:8080;多 FE 用分号分隔。用于批量写入数据 |
username/password | 是 | - | StarRocks 账号;本例 allin1 镜像 root 无密码 |
sink.buffer-flush.max-bytes | 否 | 150 MB | 缓冲区写满字节数(内存缓冲为所有表共享) |
sink.buffer-flush.interval-ms | 否 | 300000 | 每张表的数据刷新间隔 |
sink.io.thread-count | 否 | 2 | 不同表之间并发 Stream Load 的线程数 |
sink.at-least-once.use-transaction-stream-load | 否 | true | at-least-once 语义下是否使用事务 Stream Load |
table.create.num-buckets | 否 | - | 自动建表时的分桶数;StarRocks 2.5+ 可不设,由 StarRocks 自动决定 |
table.create.properties.* | 否 | - | 自动建表时追加的建表属性。本例的table.create.properties.replication_num: 1正是因为 Docker allin1 镜像只有 1 个 BE 节点,副本数必须为 1 |
table.schema-change.timeout | 否 | 30 min | StarRocks 侧 schema 变更的超时时间 |
unicode-char.max-bytes | 否 | 3 | CHAR/VARCHAR 字符到字节的换算系数,上游若使用 utf8mb4 建议设为 4,避免列长被低估 |
关于table.create.properties.replication_num的解析机制:StarRocksDataSinkFactory 在校验配置时专门放行了table.create.properties.与sink.properties.两个前缀,随后 TableCreateConfig.from 会把所有带table.create.properties.前缀的键剥掉前缀、转成小写后原样拼进自动生成的CREATE TABLEDDL 中——这正是能自由透传 StarRocks 任意建表属性(如replication_num、fast_schema_evolution)的底层原因。
pipeline 段:name是作业名,parallelism为作业并行度;此外还可在此写其他 Flink 运行时参数(如local-time-zone,StarRocks sink 会用它做TIMESTAMP_LTZ的时区转换)。
3. 提交作业
在 Flink CDC 主目录执行:
bash bin/flink-cdc.sh mysql-to-starrocks.yaml提交成功后输出:
Pipeline has been submitted to cluster. Job ID: 02a31c92f0e7bc9a1f4c0051980088a0 Job Description: Sync MySQL Database to StarRocks此时在 Flink Web UI 中可以找到名为Sync MySQL Database to StarRocks的运行中作业:
用 DBeaver 等工具通过mysql://127.0.0.1:9030连接 StarRocks,即可看到app_db下的三张表已被自动创建并写入了全量数据。
三、同步 Schema 与数据变更
管道运行起来后,进入 MySQL 容器:
docker-compose exec mysql mysql -uroot -p123456依次对orders表做四种操作,StarRocks 端会实时发生对应变化:
插入一条记录:
INSERT INTO app_db.orders (id, price) VALUES (3, 100.00);新增一列(Schema 变更事件):
ALTER TABLE app_db.orders ADD amount varchar(100) NULL;更新一条记录:
UPDATE app_db.orders SET price=100.00, amount=100.00 WHERE id=1;删除一条记录:
DELETE FROM app_db.orders WHERE id=2;
每执行一步刷新一次 DBeaver,StarRocks 中orders表的结构和数据都会实时更新。对shipments、products表做同样操作,也能看到实时同步结果。
从源码结构看,这类变更之所以能“零代码”生效,是因为 StarRocks sink 实现了MetadataApplier接口:StarRocksDataSink.getMetadataApplier 返回StarRocksMetadataApplier,运行时收到AddColumnEvent等 Schema 变更事件后,会通过StarRocksEnrichedCatalog走 JDBC 通道向 StarRocks 执行对应的ALTER TABLE。按 StarRocks 连接器文档说明,目前支持的 DDL 同步包括建表/删表/清表、增列/删列/改列名/改列类型,且新列总是追加到表末尾。
四、Route:把源表路由到目标表名
Flink CDC 的route配置可以把源表的结构和数据路由到其他表名,从而实现库名/表名替换与整库迁移。在上面配置基础上追加route段:
################################################################################ # Description: Sync MySQL all tables to StarRocks ################################################################################ source: type: mysql hostname: localhost port: 3306 username: root password: 123456 tables: app_db.\.* server-id: 5400-5404 server-time-zone: UTC sink: type: starrocks jdbc-url: jdbc:mysql://127.0.0.1:9030 load-url: 127.0.0.1:8080 username: root password: "" table.create.properties.replication_num: 1 route: - source-table: app_db.orders sink-table: ods_db.ods_orders - source-table: app_db.shipments sink-table: ods_db.ods_shipments - source-table: app_db.products sink-table: ods_db.ods_products pipeline: name: Sync MySQL Database to StarRocks parallelism: 2应用上述route配置后,app_db.orders的表结构与数据会被同步到ods_db.ods_orders,相当于完成了一次“库迁移”。
source-table支持正则匹配多张表,这是合并分片表的常用手段:
route: - source-table: app_db.order\.* sink-table: ods_db.ods_orders这样app_db.order01、app_db.order02、app_db.order03等分片表会被合并同步到同一张ods_db.ods_orders表。需要提醒:目前尚不支持多张表之间存在相同主键数据的场景(该限制在当前版本中已被明确,后续版本计划支持)。
从源码看,route每条规则解析为RouteDef(见 toRouteDef):source-table与sink-table必填,另支持可选的replace-symbol(用符号替换表名实现批量重命名,如app_db.user_[0-9]+->app_db.${replaceSymbol}_user)和description(规则备注)。
五、关键设计点与注意事项
- 数据语义是 at-least-once:从 StarRocksDataSinkFactory 可以看到,CDC 框架下该 sink 固定设置
sink.semantic = at-least-once,并且 Stream Load 强制使用 JSON 格式(sink.properties.format = json、strip_outer_array、ignore_json_size均自动注入)。幂等性来自“主键表 + 主键去重”:重复写入同一主键的行会被覆盖,而不是报错或重复。因此源表必须有主键。 - 自动建表的规则:sink 会在目标库不存在表时自动建表,主键与分布键相同、不创建分区;
table.create.properties.*可透传任意建表属性。如果 StarRocks 为 3.2+ 且希望加速后续 Schema 变更,可以加table.create.properties.fast_schema_evolution: true。 - 类型映射要心里有数:
DECIMAL(p,s)、INT等直接对应;TIME映射为VARCHAR(按HH:mm:ss存字符串);CHAR/VARCHAR按unicode-char.max-bytes从字符长度换算为 StarRocks 的字节长度,超过上限或主键列会退化为VARCHAR,完整对照表见 StarRocks 连接器文档的类型映射章节。 - Checkpoint 间隔影响可见延迟:教程将 Checkpoint 设为 3 秒,增量数据随 Stream Load 缓冲(默认 150 MB 或 5 分钟刷新)写入 StarRocks,端到端可见延迟由两者共同决定。
六、清理环境
教程结束后,在docker-compose.yml所在目录停止容器:
docker-compose down在 Flink 主目录停止集群:
./bin/stop-cluster.sh七、小结
本教程完整走通了 Flink CDC 面向 Flink 1.20 的 MySQL → StarRocks 流式 ELT 全流程:Docker 拉起 MySQL 与 StarRocks、用纯 YAML 定义 source/sink/pipeline、CLI 一行命令提交、实时验证数据与 Schema 双向一致性,并用route演示了库迁移与分片表合并。所有行为都有仓库源码可查证——YAML 解析在 flink-cdc-cli 的 YamlPipelineDefinitionParser,参数语义在 MySQL source 选项 与 StarRocks sink 选项/工厂,端到端回归测试可参考 flink-cdc-pipeline-e2e-tests。掌握这套“YAML + CLI”的模式后,将其中的 sink 替换为 Kafka、Doris、Paimon 等其它 pipeline 连接器(见 pipeline 连接器总览)即可快速搭建不同的实时数据链路。
【免费下载链接】flink-cdcFlink CDC is a streaming data integration tool项目地址: https://gitcode.com/GitHub_Trending/flin/flink-cdc
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考