FlinkCDC与Dinky实现StarRocks整库同步
概述
Dinky 是一个基于 Apache Flink 的实时计算平台,它提供了一站式的 Flink 任务开发、运维、监控等功能。本教程将一步一步教你如何使用 Dinky 运行 CDC Pipeline 任务实现整库同步到 StarRocks 并自动建表功能。Doris 同理。
前置条件
Dinky 目前只支持 JDK8 和 JDK11,MySQL 5.7 以上。
- 已部署好的 Dinky
- 准备好 Flink 集群、StarRocks 集群
环境准备
- FLINK_VERSION: 1.19
- dinky-1.19-1.2.0
- StarRocks 3.3.7 存算分离
- MySQL 8.0.30
Dinky 上使用 FlinkCDC 3.1 Pipeline 同步 MySQL 到 StarRocks 需要的依赖。
Jar 包下载目录
在 Flink/lib 和 dinky/extends 目录下放置以下 Jar 包:
- mysql-connector-java-8.0.30.jar
- flink-sql-connector-mysql-cdc-3.2.1.jar
- flink-connector-starrocks-1.2.10_flink-1.19.jar
- flink-cdc-pipeline-connector-starrocks-3.2.1.jar
- flink-cdc-pipeline-connector-mysql-3.2.1.jar
- flink-cdc-dist-new.3.2.1.jar
- flink-cdc-pipeline-connector-paimon-3.2.1.jar
如果直接在 Dinky 使用 flink-cdc-dist-3.2.1.jar 会有
java.lang.NoSuchMethodError: org.apache.calcite.tools.FrameworkConfig.getTraitDefs()Lorg/apache/flink/calcite/shaded/com/google/common/collect/ImmutableList;错误,所以我们需要先处理一下。
bash
# 解压 flink-cdc-3.2.1-bin.tar.gz
tar -zxvf flink-cdc-3.2.1-bin.tar.gz
cd flink-cdc-3.1.0/lib/
# 解压 jar 文件
jar -xvf flink-cdc-dist-3.1.0.jar
# 删除冲突包
rm -rf org/apache/calcite
# 重新打包
jar -cvf flink-cdc-dist-3.1.0-new.jar *
# 重命名
mv flink-cdc-dist-3.1.0-new.jar flink-cdc-dist-3.1.0.jar
把最新的 flink-cdc-dist-3.2.1.jar 文件放到 Dinky 依赖目录下,重启 Dinky。
所有 Jar 包:


开始运行
打开 Dinky 页面,新建 Flink SQL 任务,输入以下代码,注意把相关 IP 替换成你自己的。Flink 集群需要自己提前注册好,选择对应集群。
sql
SET 'execution.checkpointing.interval' = '10s';
EXECUTE PIPELINE WITH YAML (
source:
type: mysql
hostname: xxxx
port: 3306
username: root
password: 'xxx!'
tables: test.t_user, test.t_teacher
server-id: 5400-5404
sink:
type: starrocks
name: Starrocks Sink
jdbc-url: jdbc:mysql://ip:9030
load-url: ip:8029
username: root
password: 'xxxx'
table.create.properties.replication_num: 1
table.create.properties.fast_schema_evolution: true
route:
- source-table: test.t_user
sink-table: starrocks_test.t_user
- source-table: test.t_teacher
sink-table: starrocks_test.teacher
pipeline:
name: Sync MySQL Database to Starrocks Route
parallelism: 1
)
运行效果和常见错误

错误 1
java
java.lang.RuntimeException: com.starrocks.data.load.stream.exception.StreamLoadFailException: Transaction start failed, db: wsmg, table: bus_analog, label: flink-117dd791-84b4-4361-96b4-061c56d6c112, responseBody: {
"Status": "THRIFT_RPC_ERROR",
"Message": "call frontend service failed, address=TNetworkAddress(hostname=172.17.0.1, port=9020), reason=THRIFT_EAGAIN (timed out)"
}
解决方案:目前调大 CN 配置参数 thrift_rpc_timeout_ms 和 txn_commit_rpc_timeout_ms 超时时间。


大数据相关学习资料、大数据项目、湖仓一体、架构师必知必会、数据中台建设方法论...
共有 1400 多份文档资料,另专为星球成员整理了一份比较详细的语雀知识库合计 190 万字和飞书文档资料。(内容太多,仅展示部分内容...),欢迎大家踊跃加入星球,您将获得:
- 提供最全的大数据知识库,不限设备,随时随地打开看的在线文档。
- 免费答疑解惑、交流技术。
- 面试指导、模拟面试。
- 各类 PDF 文档下载、星球代码下载。
- 提供简历模板,简历修改指导服务,星球成员免费提供简历修改指导。
另外说明加入星球后支持三天无理由退款,不满意无条件随时退。
需要资料请加微信:D1435221412

end
