news 2026/9/17 21:40:16

Flink CDC 实战指南:在 Flink 1.20 上构建 MySQL 到 StarRocks 的流式 ELT 管道

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Flink CDC 实战指南:在 Flink 1.20 上构建 MySQL 到 StarRocks 的流式 ELT 管道

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 机器,然后:

  1. 下载并解压 Flink 1.20.3 发行包,得到flink-1.20.3目录,进入该目录并设置FLINK_HOME

    cd flink-1.20.3
  2. 开启 Checkpoint:向conf/config.yaml追加以下配置,使作业每 3 秒做一次 Checkpoint(增量快照读取与 at-least-once 写入的正确性都依赖 Checkpoint 周期性推进):

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

    ./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数据库以及ordersproductsshipments三张带主键的表,并插入初始记录(注意: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

  1. 下载 Flink CDC 稳定版二进制发行包(flink-cdc-x.y.z-bin.tar.gz,从 Apache 官方发布渠道获取)并解压,得到包含binliblogconf四个目录的flink-cdc-x.y.z目录。

  2. 将以下两个 Pipeline 连接器 JAR 下载到Flink CDC 主目录的lib目录(注意:是 Flink CDC Home 的 lib,不是 Flink Home 的 lib):

    • flink-cdc-pipeline-connector-mysql
    • flink-cdc-pipeline-connector-starrocks

    稳定版 JAR 可从 Maven 公共仓库获取;如需 SNAPSHOT 版本,则需要自行基于 master 或 release 分支构建。

  3. 由于 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: 2

YAML 的顶层结构由 YamlPipelineDefinitionParser 解析校验:sourcesink是必填块,routetransformpipeline为可选块,出现未知顶层键会直接抛出带允许键列表的错误提示,便于尽早发现拼写错误。逐段解读:

source 段(由 MySQL pipeline 连接器消费,参数定义见 MySqlDataSourceOptions):

参数说明
type固定为mysql,工厂据此通过 SPI 找到 MySQL Source
hostname/portMySQL 服务地址与端口(端口默认 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-zoneMySQL 会话时区;不设置时使用系统默认时区。必须与实际 MySQL 服务时区一致,否则 binlog 时间戳会解析错乱
scan.startup.mode可选,默认initial(先读全量快照再读增量),其他可选值还有earliest-offsetlatest-offsettimestampspecific-offsetsnapshot

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-bytes150 MB缓冲区写满字节数(内存缓冲为所有表共享)
sink.buffer-flush.interval-ms300000每张表的数据刷新间隔
sink.io.thread-count2不同表之间并发 Stream Load 的线程数
sink.at-least-once.use-transaction-stream-loadtrueat-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.timeout30 minStarRocks 侧 schema 变更的超时时间
unicode-char.max-bytes3CHAR/VARCHAR 字符到字节的换算系数,上游若使用 utf8mb4 建议设为 4,避免列长被低估

关于table.create.properties.replication_num的解析机制:StarRocksDataSinkFactory 在校验配置时专门放行了table.create.properties.sink.properties.两个前缀,随后 TableCreateConfig.from 会把所有带table.create.properties.前缀的键剥掉前缀、转成小写后原样拼进自动生成的CREATE TABLEDDL 中——这正是能自由透传 StarRocks 任意建表属性(如replication_numfast_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 端会实时发生对应变化:

  1. 插入一条记录:

    INSERT INTO app_db.orders (id, price) VALUES (3, 100.00);
  2. 新增一列(Schema 变更事件):

    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;

每执行一步刷新一次 DBeaver,StarRocks 中orders表的结构和数据都会实时更新。对shipmentsproducts表做同样操作,也能看到实时同步结果。

从源码结构看,这类变更之所以能“零代码”生效,是因为 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.order01app_db.order02app_db.order03等分片表会被合并同步到同一张ods_db.ods_orders表。需要提醒:目前尚不支持多张表之间存在相同主键数据的场景(该限制在当前版本中已被明确,后续版本计划支持)。

从源码看,route每条规则解析为RouteDef(见 toRouteDef):source-tablesink-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 = jsonstrip_outer_arrayignore_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/VARCHARunicode-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),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/17 21:39:41

数学建模国赛工具流:从Python到LaTeX的全流程协同方案

1. 从“会用工具”到“工具流”,差的是一次系统性思考“2026数学建模国赛工具流”这个标题,我第一眼看到就觉得特别对味。数学建模国赛拼到最后,真正拉开差距的往往不是某个单独的工具用得有多溜,而是整个团队从拿到题目到提交论文…

作者头像 李华
网站建设 2026/9/17 21:34:49

图书馆管理系统毕业设计系统流程图绘制全攻略

“系统流程图”这几个字,看着简单,真画起来能劝退一大半做毕业设计的同学。尤其是图书馆管理系统这种经典课设题目,业务线又长又绕,借书、还书、续借、预约、罚款、统计全搅在一起,很多同学对着Visio或draw.io发半天呆…

作者头像 李华