
Flink CDC 实战Streaming ELT 同步 MySQL 到 StarRocks 全流程指南【免费下载链接】flink-cdcFlink CDC is a streaming data integration tool项目地址: https://gitcode.com/GitHub_Trending/flin/flink-cdc本文是一篇基于 Flink CDC 的端到端实战指南讲解如何使用 Flink CDC Pipeline 模式与 Flink CDC CLI在零 Java/Scala 代码、无需 IDE 的前提下快速构建 MySQL 到 StarRocks 的 Streaming ELT 作业。文中将完整覆盖环境准备、Flink 2.2 集群搭建、Docker 组件启动、任务配置 YAML 编写、作业提交以及整库同步、表结构变更同步、路由与分库分表同步三大核心能力的实操验证并结合当前仓库源码解析底层实现原理。适用前提本教程面向 Flink CDC 2.xQuickstart for 2.2版本的 Pipeline API 与 Flink CDC CLI需要一台已安装 Docker 的 Linux 或 macOS 电脑。全部演示在命令行中完成不需要编写任何 Java/Scala 代码。一、方案总览什么是 Flink CDC 的 Streaming ELT传统的数据同步方案往往需要编写 SQL 建表语句、配置 ETL 任务再处理增量与全量数据衔接。Flink CDC 的 Pipeline 模式将这一切收敛为一个 YAML 配置声明 source、sink 与 pipeline 三大部分即可由 Flink CDC CLI 解析并编译成可提交到 Flink Standalone 集群的流式作业。本教程演示的 MySQL → StarRocks 同步作业具备三项核心能力整库同步通过tables: app_db.\.*正则一次捕获整个app_db库下的所有表表结构变更同步Schema Evolution源端执行ALTER TABLE ADD COLUMN等 DDL 后目标端 StarRocks 表结构随之自动演进分库分表同步通过route路由规则将多张分表汇总到一张目标表。整套流程的骨架是 Flink CDC CLI。仓库中 flink-cdc-cli 模块下的CliFrontend负责接收mysql-to-starrocks.yaml配置并完成作业提交而 YAML 中 source 与 sink 的类型声明则分别由 flink-cdc-pipeline-connector-mysql 与 flink-cdc-pipeline-connector-starrocks 两个 connector 通过工厂模式解析并实例化。二、准备阶段2.1 准备 Flink Standalone 集群下载 Flink 2.2.0 二进制包flink-2.2.0-bin-scala_2.12.tgz并解压进入flink-2.2.0目录cd flink-2.2.0在conf/config.yaml中追加以下参数开启每 3 秒一次的 checkpoint。checkpoint 是 CDC 作业实现故障恢复与 at-least-once 语义的基础execution: checkpointing: interval: 3s启动 Flink 集群./bin/start-cluster.sh启动成功后可访问 http://localhost:8081/ 查看 Flink Web UI。如需更多计算资源可多次执行start-cluster.sh拉起多个 TaskManager。本教程中 Pipeline 的parallelism: 2需要至少 2 个可用 Task Slot。2.2 准备 Docker 环境创建docker-compose.yml内容如下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 包含两个容器MySQL内置商品信息数据库app_db作为数据源StarRocks存储从 MySQL 根据规则映射过来的结果表。在docker-compose.yml所在目录执行docker-compose up -d该命令以 detached 模式启动所有容器。可用docker ps观察容器是否正常启动也可访问 http://localhost:8030/StarRocks FE 的 HTTP 端口9030 为 MySQL 协议查询端口确认 StarRocks 运行状态。注意端口语义StarRocks 的9030是 MySQL 协议端口供 JDBC 查询使用8080/8030是 FE 的 HTTP 端口供 Stream Load 写入使用。后文 YAML 中jdbc-url与load-url正是对应这两类端口。在 MySQL 数据库中准备数据进入 MySQL 容器docker-compose exec MySQL mysql -uroot -p123456创建数据库app_db和表orders、products、shipments并插入数据-- 创建数据库 CREATE DATABASE app_db; USE app_db; -- 创建 orders 表 CREATE TABLE orders ( id INT NOT NULL, price DECIMAL(10,2) NOT NULL, PRIMARY KEY (id) ); -- 插入数据 INSERT INTO orders (id, price) VALUES (1, 4.00); INSERT INTO orders (id, price) VALUES (2, 100.00); -- 创建 shipments 表 CREATE TABLE shipments ( id INT NOT NULL, city VARCHAR(255) NOT NULL, PRIMARY KEY (id) ); -- 插入数据 INSERT INTO shipments (id, city) VALUES (1, beijing); INSERT INTO shipments (id, city) VALUES (2, xian); -- 创建 products 表 CREATE TABLE products ( id INT NOT NULL, product VARCHAR(255) NOT NULL, PRIMARY KEY (id) ); -- 插入数据 INSERT INTO products (id, product) VALUES (1, Beer); INSERT INTO products (id, product) VALUES (2, Cap); INSERT INTO products (id, product) VALUES (3, Peanut);三张表均定义了主键——StarRocks 侧将据此自动创建主键模型Primary Key表这也是 CDC 数据实时 Upsert 更新的前提。三、通过 Flink CDC CLI 提交任务3.1 下载 Flink CDC 发行包与 Connector下载 Flink CDC 发行版二进制包flink-cdc-Version-bin.tar.gz并解压得到目录flink-cdc-Version其中包含bin、lib、log、conf四个目录。下载两个 pipeline connector 包并移动到 Flink CDC Home 的lib目录下注意是Flink CDC Home 的 lib不是 Flink Home 的 libflink-cdc-pipeline-connector-mysqlMySQL 数据源 connectorflink-cdc-pipeline-connector-starrocksStarRocks 数据汇 connector。下载链接只对已发布的稳定版本有效SNAPSHOT 版本需要基于 master 或 release- 分支自行编译。由于 CDC Connectors 不再打包数据库驱动还需要将 MySQL Connector Javamysql-connector-java-8.0.27.jar放入 Flink 的lib目录或通过--jar参数传给 Flink CDC CLI。该驱动同时服务于两部分source 端读取 MySQL binlog/快照sink 端通过jdbc-url连接 StarRocks 执行建表与 Schema 变更StarRocks 的 JDBC 连接兼容 MySQL 协议。3.2 编写任务配置 YAML创建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关键点说明source.tables: app_db.\.*通过正则表达式匹配并同步app_db下的所有表。注意这里的.被用作库名与表名的分隔符因此在正则中要用\.转义。sink.table.create.properties.replication_num: 1该参数是为本教程特设的——Docker 镜像中只有一个 StarRocks BE 节点将副本数设为 1 以避免建表失败。它是table.create.properties.前缀配置的一个实例用于在建表 DDL 中注入自定义属性。source 配置项解析结合源码从 MySqlDataSourceOptions.java 可以看到 source 侧各配置项的含义hostnameMySQL 服务器 IP 或主机名必填portMySQL 端口默认3306username/password连接 MySQL 的账号密码tables待监控的 MySQL 表名支持正则表达式。.,库表分隔符需用反斜杠转义例如db0.\.*、db1.user_table_[0-9]、db[1-2].[app|web]_order_\.*server-id本连接器在 MySQL 集群中作为 binlog 客户端时的唯一 ID支持单值如5400或区间如5400-5404两种写法启用增量快照默认开启时推荐使用区间写法区间长度即快照并行度。不设置时默认在 5400~6400 之间随机生成官方建议显式指定server-time-zone数据库服务器的会话时区不设置时使用 JVM 默认时区。此外该文件还定义了scan.startup.mode默认initial可选earliest-offset、latest-offset、timestamp、specific-offset、snapshot、scan.incremental.snapshot.chunk.size默认 8096 行快照分片大小等进阶参数。其中schema-change.enabled默认值为true意味着 schema 变更事件默认会被发出并传递给下游 sink——这正是本教程表结构变更同步能力的开关。sink 配置项解析结合源码从 StarRocksDataSinkOptions.java 可以看到 sink 侧的关键配置jdbc-urlJDBC 连接地址形如jdbc:mysql://fe_ip1:query_port,fe_ip2:query_port...用于 DDL/DML 查询操作load-urlStream Load 地址形如fe_ip1:http_port;http://fe_ip2:http_port;https://fe_nlb不指定协议前缀时默认使用 http用于数据批量写入username/passwordStarRocks 账号密码sink.buffer-flush.max-bytes单次 flush 的最大字节数默认150MBsink.buffer-flush.interval-ms行批次 flush 间隔默认300000ms5 分钟sink.connect.timeout-ms连接load-url的超时时间默认30000mssink.io.thread-count不同表之间并发 Stream Load 的线程数默认2table.create.num-buckets建表分桶数不设置时由 StarRocks 自动选择table.schema-change.timeoutStarRocks 侧 schema 变更超时时间默认1800s超时后 StarRocks 会取消变更并导致 sink 失败unicode-char.max-bytesCHAR/VARCHAR 类型映射时每个字符分配的字节数1~4默认 3若上游使用 utf8mb4建议设为 4 以避免低估列长度。值得注意的是 StarRocksDataSinkFactory.java 中的两处硬编码Stream Load 强制使用 JSON 格式sink.properties.formatjson、strip_outer_arraytrue、ignore_json_sizetrue这样在 schema 变更后无需重新配置columns属性sink 语义固定为at-least-once这也是当前 CDC 框架统一支持的投递语义。3.3 提交作业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连接jdbc:mysql://127.0.0.1:9030即可在 StarRocks 中看到自动创建并写入数据的三张表。四、同步变更验证 Schema Evolution 与数据实时更新进入 MySQL 容器docker-compose exec mysql mysql -uroot -p123456随后依次执行以下操作StarRocks 中的订单数据将实时联动更新在 MySQL 的orders表中插入一条数据INSERT INTO app_db.orders (id, price) VALUES (3, 100.00);在 MySQL 的orders表中增加一个字段ALTER TABLE app_db.orders ADD amount varchar(100) NULL;在 MySQL 的orders表中更新一条数据UPDATE app_db.orders SET price100.00, amount100.00 WHERE id1;在 MySQL 的orders表中删除一条数据DELETE FROM app_db.orders WHERE id2;每执行一步后在 Dbeaver 中刷新即可看到 StarRocks 的orders表在实时变化新增的行出现、amount新列自动补上、id1的行被更新、id2的行被删除。同样的修改shipments、products表也能在 StarRocks 中实时看到同步结果。底层原理MetadataApplier 如何执行 DDL表结构自动变更背后是 StarRocks connector 的 StarRocksMetadataApplier.java。它实现了 CDC 框架的MetadataApplier接口支持如下 schema 变更事件类型CREATE_TABLE建表时若目标库不存在会自动createDatabase再创建主键模型表ADD_COLUMN向表中追加列。源码中明确指出会忽略列位置信息StarRocks 主键列必须在表头与 MySQL 列顺序可能不一致且不允许 FIRST 位置总是将新列追加到末尾同时通过catalog.getTable二次校验列是否真正添加成功用于容错 failover 后重复执行 DDL 的场景DROP_COLUMN删除列同样带删除后校验的幂等逻辑RENAME_COLUMN、ALTER_COLUMN_TYPE重命名列与修改列类型DROP_TABLE、TRUNCATE_TABLE删表与清空表。其中ADD_COLUMN的幂等校验逻辑值得关注如果 MySQL 执行ALTER TABLE ... ADD amount后作业发生故障重启StarRocks 侧可能已经加过该列此时重复执行会失败。实现中先捕获异常再读取目标表结构确认列是否已存在若已存在则视为成功并忽略异常从而保证 CDC 侧语义与 StarRocks 实际状态一致。五、路由变更表名映射与分库分表同步Flink CDC 提供将源表的表结构/数据路由到其他表名的能力借助它可实现表名、库名替换以及整库同步等场景。完整配置示例如下################################################################################ # 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:8030 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 支持正则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一张表中。注意事项目前尚不支持多张表中存在相同主键数据的场景多表主键冲突会导致数据互相覆盖该场景将在后续版本支持。底层原理RouteRule 的匹配与替换从仓库源码看路由配置在 composer 层被解析为 RouteDef.java其中包含三个要素sourceTable匹配输入表 ID 的正则模式必填sinkTable替换匹配到的表 ID 作为输出的目标表名必填replaceSymbol可选的替换符号用于将源表名中正则捕获的部分按规则映射到目标表名。随后在 SchemaOperatorTranslator.java 中这些RouteDef会被传入SchemaOperatorschema 事件处理算子由算子内部的RouteRule对每条流经的表 ID 进行正则匹配与替换。这意味着路由不仅作用于数据事件也作用于 schema 变更事件——新表、新列、类型变更都会按同一套路由规则落到正确的目标表上。整库同步正是tables正则捕获所有源表 route规则重定向目标表的组合拳。六、环境清理教程结束后在docker-compose.yml所在目录执行docker-compose down停止所有容器再回到 Flink 目录flink-2.2.0下执行./bin/stop-cluster.sh停止 Flink 集群。七、延伸阅读与可验证源码本文涉及的实操能力在仓库中均有对应实现与测试可进一步深入Flink CDC CLI入口位于 flink-cdc-cliCliFrontend负责解析命令行参数与 YAML 配置并提交作业MySQL source connectorMySqlDataSourceFactory.java 负责实例化 MySQL 数据源MySqlDataSourceOptions.java 定义了全部配置项StarRocks sink connectorStarRocksDataSinkFactory.java 与 StarRocksDataSinkOptions.java 定义了 sink 的构建与参数StarRocksMetadataApplier.java 实现 schema 变更落地TableCreateConfig.java 解析table.create.properties.*前缀的建表属性路由能力RouteDef.java 与 SchemaOperatorTranslator.java 是路由规则的解析与注入位置更多玩法本仓库docs/content.zh/docs/get-started/quickstart-for-2.2/下还提供了 MySQL 到 Doris、MySQL 到 Kafka、PostgreSQL 到 Fluss 等 Quickstart 教程读者可在掌握本流程后举一反三。至此你已经完整走通了一条基于 Flink CDC 的 MySQL → StarRocks 实时同步链路并理解了整库同步、Schema Evolution 与分库分表路由背后的配置与源码实现。【免费下载链接】flink-cdcFlink CDC is a streaming data integration tool项目地址: https://gitcode.com/GitHub_Trending/flin/flink-cdc创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考