Flink CDC零基础教程:MySQL到PostgreSQL实时同步全流程
·
Flink CDC零基础教程:MySQL到PostgreSQL实时同步全流程
本教程将逐步指导您使用Flink CDC实现MySQL到PostgreSQL的实时数据同步。无需前置经验,只需按步骤操作即可完成全流程部署。
1. 核心概念
- Flink CDC:基于Apache Flink的Change Data Capture工具,实时捕获数据库变更。
- 同步原理:
- 通过MySQL的binlog捕获数据变更(插入/更新/删除)
- 将变更事件转换为Flink数据流
- 写入PostgreSQL目标表
- 优势:毫秒级延迟、端到端一致性、零代码侵入。
2. 环境准备
| 组件 | 版本要求 | 配置说明 |
|---|---|---|
| Flink | ≥1.13 | 下载地址:Flink官网 |
| MySQL | ≥5.7 | 需开启binlog:log_bin=ON |
| PostgreSQL | ≥10 | 创建目标表结构 |
| Flink CDC Connector | 2.3+ | 下载JAR包:Maven仓库 |
关键配置:在MySQL的
my.cnf中添加:[mysqld] server_id=1 log_bin=mysql-bin binlog_format=ROW
3. 同步流程实现
步骤1:添加依赖
在Flink项目的pom.xml中引入:
<dependencies>
<!-- Flink CDC MySQL -->
<dependency>
<groupId>com.ververica</groupId>
<artifactId>flink-connector-mysql-cdc</artifactId>
<version>2.3.0</version>
</dependency>
<!-- Flink JDBC for PostgreSQL -->
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-connector-jdbc_2.12</artifactId>
<version>1.14.4</version>
</dependency>
</dependencies>
步骤2:编写同步程序(Java示例)
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import com.ververica.cdc.debezium.JsonDebeziumDeserializationSchema;
import com.ververica.cdc.connectors.mysql.source.MySqlSource;
public class MySQL2PostgreSQLSync {
public static void main(String[] args) throws Exception {
// 1. 创建Flink环境
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// 2. 配置MySQL源(需替换实际参数)
MySqlSource<String> mysqlSource = MySqlSource.<String>builder()
.hostname("localhost")
.port(3306)
.databaseList("your_db") // 数据库名
.tableList("your_db.your_table") // 表名
.username("mysql_user")
.password("mysql_pwd")
.deserializer(new JsonDebeziumDeserializationSchema())
.build();
// 3. 定义PostgreSQL写入器
String pgSinkSQL = "INSERT INTO target_table (id, name) VALUES (?, ?) " +
"ON CONFLICT (id) DO UPDATE SET name = EXCLUDED.name";
// 4. 构建同步管道
env.fromSource(mysqlSource, WatermarkStrategy.noWatermarks(), "MySQL Source")
.map(record -> {
// 解析JSON变更数据(实际需根据业务处理)
return new Tuple2<>(record.getInt("id"), record.getString("name"));
})
.addSink(JdbcSink.sink(
pgSinkSQL,
(stmt, data) -> {
stmt.setInt(1, data.f0); // id
stmt.setString(2, data.f1); // name
},
JdbcExecutionOptions.builder().build(),
new JdbcConnectionOptions.JdbcConnectionOptionsBuilder()
.withUrl("jdbc:postgresql://localhost:5432/pg_db")
.withDriverName("org.postgresql.Driver")
.withUsername("pg_user")
.withPassword("pg_pwd")
.build()
));
// 5. 启动作业
env.execute("MySQL-to-PostgreSQL-CDC-Sync");
}
}
步骤3:启动任务
# 提交Flink作业(假设JAR包为sync-job.jar)
./bin/flink run -c com.yourapp.MySQL2PostgreSQLSync sync-job.jar
4. 验证同步效果
| 操作 | 验证方法 |
|---|---|
| 在MySQL插入新数据 | 查询PostgreSQL目标表是否同步 |
| 更新MySQL记录 | 检查PostgreSQL对应字段更新 |
| 删除MySQL数据 | 确认PostgreSQL数据被删除 |
监控建议:通过Flink Web UI(默认端口8081)观察:
- Source的
numRecordsIn(输入记录数)- Sink的
numRecordsOut(输出记录数)
5. 常见问题排查
| 问题现象 | 解决方案 |
|---|---|
| MySQL连接失败 | 检查my.cnf中binlog配置 |
| PostgreSQL写入冲突 | 优化Sink SQL的ON CONFLICT子句 |
| 数据延迟高 | 增加Flink并行度(env.setParallelism(4)) |
| 字段类型不匹配 | 在map()函数中转换数据类型 |
6. 进阶优化
- Exactly-Once语义:启用Flink Checkpointing
env.enableCheckpointing(5000); // 每5秒一次Checkpoint - 批量写入:调整JDBC Sink参数
JdbcExecutionOptions.builder() .withBatchSize(100) // 每批100条 .withBatchIntervalMs(2000) // 2秒刷新 .build() - Schema自动同步:结合Flink SQL动态建表
完整代码示例参考:Flink CDC官方示例库
更多推荐


所有评论(0)