Flink 1.20 实战:MySQL 实时同步到 ClickHouse
📝 摘要:基于 Flink 1.20 + Flink CDC 3.5,监听 MySQL Binlog 实现到 ClickHouse 的实时同步,取代 T+1 批量同步。文中给出 Binlog 开启与最小权限账号配置、Maven 依赖、Flink SQL 建源表与 JDBC 批量写入代码,用 ReplacingMergeTree 按 _version 去重保幂等、Checkpoint 保证 Exactly-Once,并附时区、打包 SPI、去重等常见问题解决。

用 Flink CDC 实现 MySQL 到 ClickHouse 的实时数据同步,增量捕获、批量写入、断点续传,一套代码搞定。本文提供完整可运行的代码,拿走即用。
前言
数据同步是大数据领域的老生常谈。传统的 T+1 批量同步已经满足不了业务需求,实时同步才是王道。
Flink CDC 是目前最主流的实时数据捕获方案:
- 基于 Binlog 增量捕获,对源库压力小
- 支持全量 + 增量无缝衔接
- Checkpoint 机制保证 Exactly-Once
本文基于 Flink 1.20 + Flink CDC 3.5,实现 MySQL 到 ClickHouse 的实时同步。
前置依赖:
- Flink 环境:Flink 1.20 基于 Docker 单机部署实战指南
- ClickHouse 环境:ClickHouse 25.4 基于 Docker 单机部署实战指南
一、环境准备
1.1 MySQL 配置
Flink CDC 依赖 MySQL 的 Binlog,需要确保以下配置:
-- 检查 Binlog 是否开启
SHOW VARIABLES LIKE 'log_bin'; -- 必须为 ON
SHOW VARIABLES LIKE 'binlog_format'; -- 必须为 ROW
如果未开启,修改 my.cnf:
[mysqld]
server-id = 1
log-bin = mysql-bin
binlog-format = ROW
创建同步专用账号(最小权限原则):
CREATE USER 'flink'@'%' IDENTIFIED BY 'flinkpw';
GRANT SELECT, RELOAD, SHOW DATABASES, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO 'flink'@'%';
FLUSH PRIVILEGES;
权限说明:
REPLICATION SLAVE和REPLICATION CLIENT是读取 Binlog 的必要权限,SELECT用于全量快照阶段。
1.2 MySQL 源表
准备一张测试表:
CREATE TABLE `test` (
`id` int(11) NOT NULL AUTO_INCREMENT COMMENT 'id',
`time` datetime DEFAULT NULL COMMENT '时间',
`error_message` text COMMENT '错误信息',
`send_error_message` text COMMENT '发送的错误信息',
PRIMARY KEY (`id`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
-- 插入一些测试数据
INSERT INTO test (time, error_message, send_error_message) VALUES
(NOW(), '测试错误1', '发送内容1'),
(NOW(), '测试错误2', '发送内容2');
1.3 ClickHouse 目标表
在 ClickHouse 创建对应的 ODS 层表(ODS = Operational Data Store,贴源层,存放与源库结构基本一致的原始数据),推荐使用 ReplacingMergeTree 引擎:
CREATE TABLE ods.ods_test (
id UInt64 COMMENT '自增主键',
time Nullable(DateTime('Asia/Shanghai')) COMMENT '时间',
error_message Nullable(String) COMMENT '错误信息',
send_error_message Nullable(String) COMMENT '发送的错误信息',
_version DateTime('Asia/Shanghai') DEFAULT now() COMMENT '版本时间,用于去重'
)
ENGINE = ReplacingMergeTree(_version)
PARTITION BY toYYYYMM(time)
ORDER BY (id)
SETTINGS
deduplicate_merge_projection_mode = 'rebuild',
allow_nullable_key = 1
COMMENT 'ODS层 - test表同步';
为什么用 ReplacingMergeTree?
| 引擎 | 特点 | 适用场景 |
|---|---|---|
| MergeTree | 不去重,追加写入 | 日志类数据 |
| ReplacingMergeTree | 按 ORDER BY 去重,保留最新版本 | CDC 同步(有更新操作) |
| CollapsingMergeTree | 支持删除标记 | 需要物理删除的场景 |
CDC 同步会产生多次 INSERT(全量 + 增量更新),ReplacingMergeTree 通过 _version 字段自动保留最新数据。
📖 注意:ReplacingMergeTree 的去重只在后台 merge 时发生、时机不确定,查询要加
FINAL才能立刻看到去重结果(见 §四 Q5)。机制细节以 ClickHouse 官方 · ReplacingMergeTree 为准。
二、Flink 应用开发
2.1 Maven 依赖
创建 Maven 项目,核心依赖如下:
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>
<groupId>com.example</groupId>
<artifactId>flink-mysql-to-ck</artifactId>
<version>1.0-SNAPSHOT</version>
<properties>
<maven.compiler.source>17</maven.compiler.source>
<maven.compiler.target>17</maven.compiler.target>
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
<flink.version>1.20.1</flink.version>
</properties>
<dependencies>
<!-- Flink 核心 -->
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-java</artifactId>
<version>${flink.version}</version>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-table-api-java-bridge</artifactId>
<version>${flink.version}</version>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-connector-base</artifactId>
<version>${flink.version}</version>
</dependency>
<!-- Flink CDC:MySQL 数据捕获 -->
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-connector-mysql-cdc</artifactId>
<version>3.5.0</version>
</dependency>
<!-- MySQL 驱动 -->
<dependency>
<groupId>com.mysql</groupId>
<artifactId>mysql-connector-j</artifactId>
<version>8.4.0</version>
</dependency>
<!-- ClickHouse JDBC 驱动 -->
<dependency>
<groupId>com.clickhouse</groupId>
<artifactId>clickhouse-jdbc</artifactId>
<version>0.8.5</version>
<classifier>all</classifier>
</dependency>
<!-- Flink JDBC Connector:写入 ClickHouse -->
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-connector-jdbc</artifactId>
<version>3.3.0-1.20</version>
</dependency>
</dependencies>
<build>
<plugins>
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-compiler-plugin</artifactId>
<version>3.8.1</version>
<configuration>
<source>17</source>
<target>17</target>
</configuration>
</plugin>
<!-- Shade 插件:打 Fat Jar -->
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-shade-plugin</artifactId>
<version>3.5.0</version>
<executions>
<execution>
<phase>package</phase>
<goals>
<goal>shade</goal>
</goals>
<configuration>
<transformers>
<transformer implementation="org.apache.maven.plugins.shade.resource.ServicesResourceTransformer"/>
<transformer implementation="org.apache.maven.plugins.shade.resource.ManifestResourceTransformer">
<mainClass>com.example.MySQLToClickHouseSync</mainClass>
</transformer>
</transformers>
<filters>
<filter>
<artifact>*:*</artifact>
<excludes>
<exclude>META-INF/*.SF</exclude>
<exclude>META-INF/*.DSA</exclude>
<exclude>META-INF/*.RSA</exclude>
</excludes>
</filter>
</filters>
</configuration>
</execution>
</executions>
</plugin>
</plugins>
</build>
<profiles>
<!-- 本地开发:包含 Flink 运行时 -->
<profile>
<id>dev</id>
<activation>
<activeByDefault>true</activeByDefault>
</activation>
<dependencies>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-streaming-java</artifactId>
<version>${flink.version}</version>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-clients</artifactId>
<version>${flink.version}</version>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-table-planner_2.12</artifactId>
<version>${flink.version}</version>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-table-runtime</artifactId>
<version>${flink.version}</version>
</dependency>
<!-- 日志 -->
<dependency>
<groupId>org.apache.logging.log4j</groupId>
<artifactId>log4j-slf4j-impl</artifactId>
<version>2.17.2</version>
</dependency>
<dependency>
<groupId>org.apache.logging.log4j</groupId>
<artifactId>log4j-core</artifactId>
<version>2.17.2</version>
</dependency>
</dependencies>
</profile>
<!-- 生产打包:Flink 运行时由集群提供 -->
<profile>
<id>prod</id>
<dependencies>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-streaming-java</artifactId>
<version>${flink.version}</version>
<scope>provided</scope>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-clients</artifactId>
<version>${flink.version}</version>
<scope>provided</scope>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-table-planner_2.12</artifactId>
<version>${flink.version}</version>
<scope>provided</scope>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-table-runtime</artifactId>
<version>${flink.version}</version>
<scope>provided</scope>
</dependency>
</dependencies>
</profile>
</profiles>
</project>
依赖说明:
| 依赖 | 版本 | 作用 |
|---|---|---|
flink-connector-mysql-cdc | 3.5.0 | 监听 MySQL Binlog |
mysql-connector-j | 8.4.0 | MySQL JDBC 驱动 |
clickhouse-jdbc | 0.8.5 | ClickHouse JDBC 驱动 |
flink-connector-jdbc | 3.3.0-1.20 | Flink JDBC Sink |
Profile 说明:
dev用于本地 IDE 调试,prod用于提交到集群(Flink 运行时由集群提供,设为 provided 减小包体积)。
2.2 同步代码
import org.apache.flink.connector.jdbc.JdbcConnectionOptions;
import org.apache.flink.connector.jdbc.JdbcExecutionOptions;
import org.apache.flink.connector.jdbc.JdbcSink;
import org.apache.flink.connector.jdbc.JdbcStatementBuilder;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.table.api.Table;
import org.apache.flink.table.api.bridge.java.StreamTableEnvironment;
import org.apache.flink.types.Row;
import java.sql.Types;
/**
* MySQL -> ClickHouse 实时同步
* 使用 Flink CDC + JDBC Sink,支持批量写入
*/
public class MySQLToClickHouseSync {
public static void main(String[] args) throws Exception {
// 1. 创建执行环境
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.setParallelism(1); // CDC 源建议并行度为 1
env.enableCheckpointing(5000L); // 5 秒一次 Checkpoint
StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);
// 2. 创建 MySQL CDC 源表
tableEnv.executeSql("""
CREATE TABLE mysql_source (
id INT,
`time` TIMESTAMP(0),
error_message STRING,
send_error_message STRING,
PRIMARY KEY (id) NOT ENFORCED
) WITH (
'connector' = 'mysql-cdc',
'hostname' = '192.168.1.100',
'port' = '3306',
'username' = 'flink',
'password' = 'flinkpw',
'database-name' = 'test_db',
'table-name' = 'test',
'scan.incremental.snapshot.enabled' = 'true',
'server-time-zone' = 'Asia/Shanghai',
'scan.startup.mode' = 'initial'
)
""");
// 3. 转换为 DataStream
Table sourceTable = tableEnv.from("mysql_source");
DataStream<Row> dataStream = tableEnv.toChangelogStream(sourceTable);
// 4. 写入 ClickHouse(批量 Sink)
dataStream.addSink(JdbcSink.sink(
// INSERT 语句
"INSERT INTO ods_test (id, `time`, error_message, send_error_message, _version) " +
"VALUES (?, ?, ?, ?, now())",
// 参数绑定
(JdbcStatementBuilder<Row>) (ps, row) -> {
ps.setInt(1, (Integer) row.getField(0));
Object timeValue = row.getField(1);
if (timeValue != null) {
ps.setObject(2, timeValue);
} else {
ps.setNull(2, Types.TIMESTAMP);
}
ps.setString(3, (String) row.getField(2));
ps.setString(4, (String) row.getField(3));
},
// 执行选项:批量写入配置
JdbcExecutionOptions.builder()
.withBatchSize(2000) // 每批 2000 条
.withBatchIntervalMs(5000) // 最长 5 秒刷一次
.withMaxRetries(3) // 失败重试 3 次
.build(),
// 连接选项
new JdbcConnectionOptions.JdbcConnectionOptionsBuilder()
.withUrl("jdbc:clickhouse://192.168.1.101:8123/ods")
.withDriverName("com.clickhouse.jdbc.ClickHouseDriver")
.withUsername("admin")
.withPassword("admin123")
.build()
)).name("ClickHouse Sink");
// 5. 执行
env.execute("MySQL to ClickHouse Sync");
}
}
核心参数说明:
| 参数 | 值 | 说明 |
|---|---|---|
scan.startup.mode | initial | 先全量快照,再增量同步 |
scan.incremental.snapshot.enabled | true | 无锁快照,不阻塞源库 |
server-time-zone | Asia/Shanghai | 时区,必须和 MySQL 一致 |
withBatchSize | 2000 | 批量写入大小,提升 CK 写入性能 |
withBatchIntervalMs | 5000 | 批次间隔,避免数据延迟太久 |
三、打包部署
3.1 打包
# 生产环境打包(排除 Flink 运行时)
mvn clean package -P prod -DskipTests
# 生成的 jar 包在 target 目录下
ls target/*.jar
3.2 提交任务
通过 Flink Web UI 提交:
- 打开
http://<Flink-IP>:8081 - 点击 Submit New Job
- 上传打好的 jar 包
- 填写 Entry Class(如果 pom 中配置了 mainClass 可省略)
- 点击 Submit

3.3 验证运行
进入 Jobs → Running Jobs,可以看到任务正常运行:

在 MySQL 插入或更新数据,几秒后在 ClickHouse 查询验证:
SELECT * FROM ods.ods_test ORDER BY _version DESC LIMIT 10;
四、常见问题
Q1: 报错 “Access denied for user ‘flink’@‘%’”?
检查 MySQL 用户权限,确保授予了 REPLICATION SLAVE 和 REPLICATION CLIENT。
Q2: 时间字段差 8 小时?
检查 server-time-zone 配置是否为 Asia/Shanghai,同时确保 ClickHouse 表字段也指定了时区。
Q3: ClickHouse 写入慢?
- 增大
withBatchSize(如 5000) - 检查 CK 的
max_insert_block_size配置 - 避免过高的 Checkpoint 频率
Q4: 全量阶段太慢?
scan.incremental.snapshot.enabled = true 使用无锁快照,但大表仍需要时间。可以考虑:
- 分库分表并行同步
- 先用 DataX 做全量,再用 Flink CDC 做增量
Q5: 任务重启后数据重复?
这是正常现象,ReplacingMergeTree 会自动去重。查询时加 FINAL 可立即看到去重结果:
SELECT * FROM ods_test FINAL WHERE id = 1;
Q6: 报错 “Unable to create a source for reading table”?
报错信息类似:
ValidationException: Unable to create a source for reading table
'default_catalog.default_database.xxx_cdc'.
原因:使用了 maven-assembly-plugin 打包,该插件无法正确处理 Flink CDC 的 SPI(Service Provider Interface)文件,导致连接器无法被识别。
解决:改用 maven-shade-plugin 打包(本文 2.1 节的 pom.xml 已使用正确插件)。
⚠️ 注意是「shade +
ServicesResourceTransformer」,不是换成 shade 就完事。真正解决问题的是这个 transformer——它把各个 jar 里META-INF/services/下的同名 SPI 文件合并成一个;不配它的话,shade 打包时后一个 jar 的 SPI 文件会直接覆盖前一个,连接器照样注册不上,报的还是同一个错。所以如果你已经在用 shade 却仍然报
Unable to create a source for reading table,先去 pom 里确认这段在不在:<transformers> <transformer implementation="org.apache.maven.plugins.shade.resource.ServicesResourceTransformer"/> </transformers>
Q7: 源表有频繁 UPDATE / DELETE,这套还适用吗?
要留意适用边界。本文这套「toChangelogStream + JdbcSink 无脑 INSERT + ReplacingMergeTree」组合,最适合以 INSERT 为主的 ODS 贴源场景(日志、流水、埋点这类只追加的表)——这也是它日同步千万级还稳的前提。一旦源表有频繁 UPDATE / DELETE,有两点要处理:
- UPDATE:注意
-U(旧值)也会被插进 CK——一次更新下发-U(旧值)和+U(新值)两条,本文的 Sink 不看 RowKind、两条都INSERT,_version都是同一秒的now()。ReplacingMergeTree 在_version相等时按「插入顺序保留最后一条」,而+U后插,所以同一 data part 内多半留到的是新值、看着没问题;但这属未定义行为,跨 part / 多次 merge 后不保证。→ 想让结果确定,把_version升到毫秒DateTime64(3),+U就严格大于-U。 - DELETE:源库删一行会下发
-D,本文代码把它当普通 INSERT 又写回 CK,删除不会传播、数据会「复活」。→ 在 Sink 前按RowKind过滤(只留+I/+U),或改用CollapsingMergeTree/VersionedCollapsingMergeTree,用sign列表达删除。
一句话:只增不改的表直接用本文方案;有增删改的表,按上面两点加固。
五、总结
本文实现了 MySQL → ClickHouse 的实时同步,核心技术栈:
| 组件 | 版本 | 作用 |
|---|---|---|
| Flink | 1.20.1 | 流处理引擎 |
| Flink CDC | 3.5.0 | Binlog 捕获 |
| ClickHouse | 25.4 | OLAP 存储 |
| ReplacingMergeTree | - | 去重引擎 |
核心要点:
- MySQL 必须开启 ROW 格式的 Binlog
- 使用 JDBC Sink 批量写入,提升 CK 性能
- ReplacingMergeTree 配合
_version字段实现幂等更新 - Checkpoint 保证 Exactly-Once 语义
这套方案经过生产验证,日同步千万级数据稳定运行。
如果这篇文章对你有帮助,欢迎点赞收藏。有问题欢迎评论区交流。
延伸阅读
- Flink 1.20 实战:零代码配置实现 MySQL 百表到 ClickHouse 实时同步 —— 进阶版:从单表到百表,新增表只改配置不写代码
- Flink 1.20 + CDC 3.5 实战:MongoDB 实时同步到 ClickHouse,从踩坑到上线 —— 换个数据源:MongoDB Change Stream 同步到 CK
- Flink Jenkinsfile 怎么写不出 bug:10 条设计要点 + 完整 Demo —— 把同步作业做成 savepoint 优雅停启的部署流水线
- Flink 重启变双开:一次部署引发的两个 CDC 任务并发消费 —— 同步作业重启不当会双写下游,这篇复盘讲清楚原因
🏷️ 标签:Flink CDC MySQL ClickHouse Binlog ReplacingMergeTree 实时同步
更多推荐
所有评论(0)