恒美微站 Logo 恒美微站
  • 首页
  • 关于我们
  • 建站服务
  • 主题模板
  • 案例展示
  • 资讯中心
  • 联系我们

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

  • 首页
  • 资讯中心
  • /
  • Flink CDC 实战指南:在 Flink 1.20 上构建 MySQL 到 StarRocks 的流式 ELT 管道

相关资讯

Java面向对象编程三大特性深度解析与实践 2026/9/17 21:40:25
卡尔曼滤波与矩阵分析:RM电控状态估计的数学基础 2026/9/17 21:40:25
深入 Kona:Optimism 仓库中 OP Stack Rust 实现的全景开发指南 2026/9/17 21:35:25

最新资讯

为 MIPS 补上 TFLite Micro 构建目标:交叉编译跑通关键词唤醒
相关性分析:三大系数与Python/SPSS实战指南
MATLAB神经网络工业级案例集:43个可直接运行的建模快照
抖音去水印批量下载免费教程:douyin-downloader 完整上手指南
Open edX XBlock App Suite 技术解析:新版 XBlock Runtime 的设计目标、Learning Context 机制与 Python/REST API 实践
pytest内置神器全掌握:monkeypatch、tmp_path与capsys实用技巧清单

今日推荐

每日热评|13% 的 Agent 技能带严重漏洞,这个注册表想用“验证+签名”解决信任危机
即梦AI保姆级教程:从生图到数字人,一站式搞定AI视频创作
BERT+LLM混合架构:突破NER长尾实体抽取瓶颈的工程实践

本周热门

AI SDK Harness 依赖更新指南:掌握 harness 包 SDK 依赖的升级、桥接同步与一致性校验
Refine v5 Ant Design NumberField 组件实战:基于 Intl 的本地化数字格式化
Flutter应用改名全指南:从Android到iOS的配置与工具实践

本月精选

自研推理加速器Redwood:两周内实现PyTorch模型高效部署的实战教程
V4L2摄像头采集实战:从camera_client.rar到出图全流程解析
从“谁发明了钢琴键”到知识问答智能体:RAG与记忆工程实践

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

发布时间:2026/9/17 21:40:25
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),仅供参考

关于恒美微站

恒美微站专注于为个体商户、工作室提供极简自助建站服务,让每个人都能轻松拥有专业网站。

快速链接

  • 关于我们
  • 建站服务
  • 主题模板
  • 案例展示
  • 资讯中心

服务项目

  • 可视化建站
  • 拖拽编辑
  • 主题定制
  • SEO 优化
  • 网站托管

联系方式

  • 📍 地址:北京市朝阳区建国路 88 号
  • 📞 电话:400-888-8888
  • ✉️ 邮箱:info@hmyw.cn
  • 🕐 时间:周一至周日 9:00-18:00

© 2024 恒美微站 hmyw.cn 版权所有 | 京 ICP 备 12345678 号