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

发布时间:2026/9/17 21:40:49

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 CDCflink-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_HOMEcd 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_PASSWORD123456 - MYSQL_USERmysqluser - MYSQL_PASSWORDmysqlpw在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 的 libflink-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 Sourcehostname/portMySQL 服务地址与端口端口默认 3306username/password连接 MySQL 的账号密码该账号需要具备读取 binlog 的权限REPLICATION SLAVE、REPLICATION CLIENTtables需要同步的表支持正则表达式。注意点号.是库名与表名的分隔符正则中要匹配“任意字符的点”必须写成转义的\.例如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-offset、latest-offset、timestamp、specific-offset、snapshotsink 段参数定义见 StarRocksDataSinkOptions参数必填默认值说明type是-固定为starrocksname否-sink 名称会体现在作业的表描述里jdbc-url是-FE 的 MySQL 协议查询端口如jdbc:mysql://127.0.0.1:9030多 FE 用逗号分隔。用于建表、执行 schema 变更等 DDLload-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否trueat-least-once 语义下是否使用事务 Stream Loadtable.create.num-buckets否-自动建表时的分桶数StarRocks 2.5 可不设由 StarRocks 自动决定table.create.properties.*否-自动建表时追加的建表属性。本例的table.create.properties.replication_num: 1正是因为 Docker allin1 镜像只有 1 个 BE 节点副本数必须为 1table.schema-change.timeout否30 minStarRocks 侧 schema 变更的超时时间unicode-char.max-bytes否3CHAR/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-zoneStarRocks 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 price100.00, amount100.00 WHERE id1;删除一条记录DELETE FROM app_db.orders WHERE id2;每执行一步刷新一次 DBeaverStarRocks 中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见 toRouteDefsource-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),仅供参考
延伸阅读

更多相关文章

2026/9/17 21:40:49

Java面向对象编程三大特性深度解析与实践

1. 面向对象编程三大基石在Java的世界里,封装、继承和多态就像建筑房屋的三大支柱,它们共同构成了面向对象编程(OOP)的核心思想体系。作为从业十余年的Java开发者,我深刻体会到这三个概念在实际工程中的重要性远超教科书上的简单定义。记得刚…

2026/9/17 21:40:49

卡尔曼滤波与矩阵分析:RM电控状态估计的数学基础

做RM电控的兄弟应该都有这种感觉:调云台稳像或者做IMU姿态解算的时候,绕不开卡尔曼滤波这道坎。我在整理中科大RM电控合集的时候,特意把卡尔曼滤波前瞻和矩阵分析基础放到一起讲,原因是很多人看公式直接被符号劝退:F、…

2026/9/17 21:40:49

2026年入行嵌入式怎么学?从单片机到Linux的硬核路线

先说结论:如果有人跟你说,26年入行嵌入式不用学Linux、不用碰驱动、不用懂内核,先把单片机焊明白就行,那你大概率会在这个行业里多走两年弯路。嵌入式这个行当,过去几年被各种“高薪”“缺口大”的帖子吹得有点失真&am…

2026/9/17 22:20:54

Linux环境下CommVault备份与恢复Oracle数据库的实践指南

简介:在Linux服务器上使用CommVault统一备份平台保护Oracle数据库时,可参考这份PDF文档。文档面向数据库管理员与运维工程师,完整覆盖了从安装准备、软件部署到备份策略配置及灾难恢复的实操流程。安装前需重点确认CommVault版本与数据库版本…

2026/9/17 22:20:54

幂律分布实战指南:识别、验证与业务干预

1. 幂律不是“定律”,而是一种观察模式——从城市人口到微博转发,它藏在你每天刷到的数据背后你有没有注意过:全国前10大城市的人口加起来,可能只占全国总人口的不到15%,但它们贡献了近40%的GDP;你发的一条…

2026/9/17 22:20:54

海岸谜题探险:沉浸式户外解谜游戏设计

1. 项目概述:海岸谜题探险的设计初衷去年夏天我在加州1号公路自驾时,被一段废弃的沿海步道激发了灵感。这条隐藏在峭壁间的步道沿途布满了风化严重的木牌,上面模糊的文字像是某种密码。这个偶然发现让我萌生了设计"Coastal Riddle Quest…

2026/9/17 22:20:54

复杂网络统计量与级联失效仿真:定位城市交通隐性瓶颈

简介:《复杂网络环境下交通流》PPT以复杂网络理论为切入点,聚焦交通系统网络化背景下的车流分析与优化,适合交通工程、网络科学及智能交通领域的研究生、科研人员或相关课程教学使用。资源包含1个PPT演示文稿,压缩包仅3.49MB&…

2026/9/17 22:20:54

在 Claude 测试影响分析里,TaoToken Key 替代临时补丁

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/9/17 22:15:54

Windows Server 2008 R2远程桌面授权报错:120天过期与CAL配置全解析

周一早上一到公司,运维群里就有人艾特我:“服务器远程桌面连不上了,提示‘远程桌面授权模式尚未配置,远程桌面服务将在11天后停止工作’。”这种报错,凡是用过Windows Server 2008 R2的人应该都不陌生。网上关于“120天…

2026/9/16 12:52:37

拯救者Y7000黑屏故障排查与维修实战指南

1. 项目概述:一台黑屏的拯救者Y7000,到底卡在哪一步? 联想拯救者Y7000系列笔记本,从2018年第一代搭载i5-8300H开始,到后来的i7-9750H、i7-10750H、i5-11400H,再到2023年款的R7-7840HS,它始终是学…

2026/9/17 0:03:13

WiFi密码安全测试:从原理到实战的字典暴力破解指南

1. 写在前面:我为什么要研究WiFi密码这件事先交代一下背景。我身边有不少朋友,家里的WiFi密码常年是"12345678"或者"88888888",问就是"好记"。直到有一次,隔壁邻居蹭网蹭到我家路由器后台都进不去&…

2026/9/17 0:03:13

redis-py服务控制与监控函数实战:从ping到slowlog的巡检指南

我用 redis-py 写了快五年的业务代码,坦白说,真正让我觉得这个客户端“像一个成熟工具箱”的,不是 get/set 那套基本操作,而是它那批专门做服务控制与状态监控的辅助函数。日常开发里,大家把redis.Redis(host..., deco…

2026/9/17 0:03:13

SpringBoot+Vue3实现中小企业设备管理系统开发实践

1. 项目概述与核心价值中小企业设备管理系统是制造业、服务业等领域的基础信息化工具。传统设备管理往往依赖Excel表格或纸质记录,存在数据孤岛、流程混乱、维护成本高等痛点。这套基于Java SpringBootVue3MyBatis的技术方案,通过前后端分离架构实现了设…

2026/9/16 22:55:57

USB Type-C PCB布局分区设计:电源、高速信号与PD协议全攻略

做硬件这行,Type-C接口算是典型的“看着简单,做起来全坑”的东西。光引脚就24个,高低速信号、电源、控制线全部塞在一个小小的连接器里,如果PCB布局不做规划,打样回来基本就是“插上没反应”、“高速掉线”、“静电一打…

2026/9/16 22:56:09

系统编程学习原型如何补齐稳定性边界

系统编程学习原型如何补齐稳定性边界预算有限时&#xff0c;我先优化明显多余的复制&#xff0c;而不是猜测性地换容器。用借用传递只读数据通常就能减少分配&#xff1a; fn parse(line: &str) -> Result<Item, Error> { /* ... */ }用基准确认热点确实在分配&am…

2026/9/16 22:56:16

雨花区哪家财务公司代理记账比较好?

在雨花区&#xff0c;企业处理财税事务常常面临诸多挑战&#xff0c;选择一家靠谱的财务公司至关重要。湖南巨勤财务管理咨询有限公司就是本地正规实体财税服务机构&#xff0c;深耕本地工商财税行业多年&#xff0c;熟悉当地工商局、税务局最新政策与申报流程。主营公司注册、…

还想了解更多?直接咨询顾问

免费诊断 + 免费方案 + 透明报价。

全国咨询热线400-8866-253
免费获取方案
咨询二维码