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. 进阶优化
  1. Exactly-Once语义:启用Flink Checkpointing
    env.enableCheckpointing(5000); // 每5秒一次Checkpoint
    

  2. 批量写入:调整JDBC Sink参数
    JdbcExecutionOptions.builder()
        .withBatchSize(100)  // 每批100条
        .withBatchIntervalMs(2000) // 2秒刷新
        .build()
    

  3. Schema自动同步:结合Flink SQL动态建表

完整代码示例参考:Flink CDC官方示例库

Logo

码道开发者社区,聚焦华为云码道 CodeArts 代码智能体,沉淀 Agent、Skill、鸿蒙开发实战内容,供开发者查阅资料、交流技术、分享工程实践

更多推荐