Flink CDC与Dinky实时同步数据至Paimon HDFS
概述
本文主要讲述通过 Flink CDC + Dinky 同步 MySQL 数据到 Paimon 中自动建表的实践。Dinky 通过定义 CDCSOURCE 语法,可以直接自动构建一个整库入仓入湖的实时任务,避免了大量的数据库连接和 DDL 编写。同时,采用多 source 合并的优化策略,减少了同一作业中的 source 数量,避免了 Binlog 的重复读取,从而减轻了源库的压力。CDCSOURCE 语句用于将上游指定数据库的所有表的数据采用一个任务同步到下游系统。

前置条件
- 安装好 Dinky
- 安装好 Flink HA 集群
环境准备
相关 jar 包
flink-sql-connector-mysql-cdc-3.2.1.jarmysql-connector-java-8.0.30.jarpaimon-flink-1.19-0.9.0.jarflink-shaded-hadoop-2-uber-2.8.3-10.0.jarflink-cdc-dist-3.2.1.jar
同步任务
sql
EXECUTE CDCSOURCE dinky_paimon_test WITH (
'connector' = 'mysql-cdc',
'hostname' = 'ip',
'port' = '3306',
'username' = 'root',
'password' = 'xxx',
'checkpoint' = '10000',
'parallelism' = '2',
'scan.startup.mode' = 'initial',
'database-name' = 'lzs-ad',
'table-name' = 'lzs-ad.t_report',
'server-id' = '5400-5499',
'sink.connector' = 'paimon',
'sink.path' = 'hdfs://mycluster/user/data/paimon/#{schemaName}.db/#{tableName}',
'sink.auto-create' = 'true',
'sink.bucket' = '4'
);
错误1
plaintext
Caused by: java.util.ServiceConfigurationError: org.dinky.cdc.SinkBuilder: org.dinky.cdc.doris.DorisSinkBuilder Unable to get public no-arg constructor

解决:添加 flink-doris-connector jar 包。

错误2

解决:在执行前设置:
sql
set table.exec.sink.upsert-materialize = NONE;
set table.exec.sink.parquet.snapshot.time-retained = 1d;
修改后的完整 SQL 代码:
sql
set table.exec.sink.upsert-materialize = NONE;
set table.exec.sink.parquet.snapshot.time-retained = 1d;
EXECUTE CDCSOURCE dinky_paimon_test WITH (
'connector' = 'mysql-cdc',
'hostname' = 'ip',
'port' = '3306',
'username' = 'root',
'password' = 'xxx!',
'checkpoint' = '10000',
'parallelism' = '2',
'scan.startup.mode' = 'initial',
'database-name' = 'lzs-ad',
'table-name' = 'lzs-ad.t_report',
'server-id' = '5400-5499',
'sink.connector' = 'paimon',
'sink.path' = 'hdfs://mycluster/user/data/paimon/#{schemaName}.db/#{tableName}',
'sink.auto-create' = 'true',
'sink.bucket' = '4'
);


StarRocks 查询 Paimon 表
sql
CREATE EXTERNAL CATALOG paimon_catalog_fs PROPERTIES (
"type" = "paimon",
"paimon.catalog.type" = "filesystem",
"paimon.catalog.warehouse" = "hdfs://bigdata02:9000/user/data/paimon"
);
-- 切换数据 CATALOG
set CATALOG paimon_catalog_fs;
use paimon_catalog_fs;
-- 删除 catalog
drop CATALOG paimon_catalog_fs;
SHOW DATABASES FROM paimon_catalog_fs;
USE paimon_catalog_fs.lzs-ad;
select * from t_report limit 3;



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