小易说IT 小易说IT

Flink SQL MySQL → PostgreSQL 实时同步完整范例

场景说明:CDC 捕获 MySQL binlog(增量 + 全量快照),实时写入 PostgreSQL;使用 Flink CDC MySQL Connector + Flink JDBC PostgreSQL Sink。 环境版本参考:Flink 1.18 /flink-cdc-connector 3.0.1;MySQL 8.0;PostgreSQL 14+ 依赖包(提交任务时需要放到 lib 或通过 -j 参数)

  • flink-sql-connector-mysql-cdc-3.0.1.jar

  • flink-connector-jdbc-3.1.2-1.18.jar

  • postgresql-42.6.0.jar

  • flink-table-planner-loader-1.18.0.jar

  • flink-runtime-1.18.0.jar

⚠️ 前置准备

  1. MySQL 开启 binlog,binlog_format=ROW;创建 CDC 账号,授予 REPLICATION SLAVE, REPLICATION CLIENT, SELECT 权限

  2. PostgreSQL 提前建好目标表,表结构必须和源表对齐(字段名、类型、主键)

  3. 源表必须有主键(CDC 需要主键,无主键只能做快照,无法消费增量变更)

-- 开启 checkpoint,保证端到端 Exactly-Once(PostgreSQL JDBC Sink 依赖 checkpoint 提交事务)
SET execution.checkpointing.interval = 30s;
SET execution.checkpointing.mode = EXACTLY_ONCE;
SET table.exec.sink.upsert-materialize = AUTO;

-- 【1】定义 MySQL CDC 源表:捕获源库 binlog,支持全量快照 + 增量变更
CREATE TABLE mysql_source (
    id INT PRIMARY KEY,
    name STRING,
    create_time TIMESTAMP(3),
    update_time TIMESTAMP(3),
    remark STRING
) WITH (
    'connector' = 'mysql-cdc',
    'hostname' = '127.0.0.1',
    'port' = '3306',
    'username' = 'cdc_user',
    'password' = 'cdc_passwd',
    'database-name' = 'source_db',
    'table-name' = 'user_info',
    'server-id' = '5400', -- 每个Flink任务server-id必须唯一,不能和mysql其他slave冲突
    'server-time-zone' = 'Asia/Shanghai',
    'scan.startup.mode' = 'initial', -- initial:先全量快照,再消费binlog;latest-offset:只消费新binlog
    'debezium.snapshot.locking.mode' = 'none' -- 生产建议:避免锁表,MySQL8.0+支持
);

-- 【2】定义 PostgreSQL JDBC 目标表,UPSERT SINK(主键存在更新,不存在插入)
CREATE TABLE pg_sink (
    id INT PRIMARY KEY,
    name STRING,
    create_time TIMESTAMP(3),
    update_time TIMESTAMP(3),
    remark STRING
) WITH (
    'connector' = 'jdbc',
    'url' = 'jdbc:postgresql://127.0.0.1:5432/target_db',
    'table-name' = 'user_info',
    'username' = 'pg_user',
    'password' = 'pg_passwd',
    'driver' = 'org.postgresql.Driver',
    -- upsert 模式,根据主键更新
    'sink.max-retries' = '3',
    'sink.buffer-flush.max-rows' = '1000', -- 批量攒多少条提交PG
    'sink.buffer-flush.interval' = '2s', -- 最大攒多久刷一次
    'sink.parallelism' = '2'
);

-- 【3】插入语句,流同步:mysql_source 的变更实时写入pg_sink
INSERT INTO pg_sink
SELECT id, name, create_time, update_time, remark FROM mysql_source;

多表同步扩展(整库同步,多表路由)

如果需要一个任务同步 MySQL 多个表到 PG 对应表,使用 CDC 整库读取 + 动态表路由(需要使用 mysql-cdc 的整库模式,搭配 Flink CDC pipeline 或者自定义表名元数据路由)

-- 整库CDC源表,读取source_db下所有表,获取元数据 table_name
CREATE TABLE mysql_all_table_source (
    id INT,
    name STRING,
    create_time TIMESTAMP(3),
    update_time TIMESTAMP(3),
    remark STRING,
    -- 内置元数据字段:当前变更所属表名
    table_name STRING METADATA FROM 'table_name' VIRTUAL,
    PRIMARY KEY (id, table_name) NOT ENFORCED
) WITH (
    'connector' = 'mysql-cdc',
    'hostname' = '127.0.0.1',
    'port' = '3306',
    'username' = 'cdc_user',
    'password' = 'cdc_passwd',
    'database-name' = 'source_db',
    'table-name' = '.*', -- 正则匹配所有表
    'server-id' = '5401',
    'server-time-zone' = 'Asia/Shanghai',
    'scan.startup.mode' = 'initial',
    'debezium.snapshot.locking.mode' = 'none'
);

多表场景下,Flink SQL 原生 JDBC Sink不支持动态目标表名。两种方案:

  1. 多个独立 Flink 任务,单表同步(稳定,生产首选)

  2. 使用 Flink CDC Pipeline(flink-cdc-pipeline-connector),声明式配置整库同步,自动映射表名,不需要写 SQL

关键参数说明

MySQL CDC 参数

参数

说明

scan.startup.mode

initial:全量快照 + 增量;latest-offset:只消费后续 binlog;specific-offset:指定 binlog 位点

debezium.snapshot.locking.mode

none:无锁快照,生产推荐;minimal:轻量锁

server-id

Flink CDC 消费 binlog,相当于 MySQL 从库,id 全局唯一

PostgreSQL JDBC Sink 参数

参数

说明

sink.buffer-flush.max-rows

批量写入大小,调大提升吞吐,增大事务压力

sink.buffer-flush.interval

超时强制刷盘,保证延迟上限

sink.max-retries

写入失败重试次数

JDBC Sink 是 UPSERT 模式,要求目标表定义主键,收到 DELETE 消息会执行 DELETE FROM table WHERE pk=?

常见坑点

  1. 时区问题:MySQL CDC 必须指定 server-time-zone='Asia/Shanghai',否则时间字段偏移 8 小时

  2. 主键缺失:源表无主键,CDC 无法识别 UPDATE/DELETE,只能 INSERT,不能做 UPSERT 同步

  3. Checkpoint 必须开启:Exactly-Once 依赖 checkpoint,checkpoint 失败会导致重复写入

  4. 类型映射

    • MySQL DATETIME → Flink TIMESTAMP(3) → PostgreSQL TIMESTAMPTZ / TIMESTAMP

    • MySQL BIGINT → Flink BIGINT → PostgreSQL BIGINT

    • MySQL VARCHAR → Flink STRING → PostgreSQL VARCHAR/TEXT

  5. PostgreSQL 连接参数:生产建议增加 ?tcpKeepAlive=true 防止长连接断连

    'url' = 'jdbc:postgresql://127.0.0.1:5432/target_db?tcpKeepAlive=true'
    

提交方式(sql-client.sh)

./sql-client.sh -f mysql2pg-cdc.sql
pipeline:
  name: MySQL to PG sync
  parallelism: 2
source:
  type: mysql-cdc
  hostname: 127.0.0.1
  port: 3306
  username: cdc_user
  password: cdc_passwd
  database-name: source_db
  table-name: "source_db.user_info"
  server-id: 5400
  server-time-zone: Asia/Shanghai
sink:
  type: jdbc
  url: jdbc:postgresql://127.0.0.1:5432/target_db
  table-name: user_info
  username: pg_user
  password: pg_passwd
  driver: org.postgresql.Driver

本文原创作者:易君召,详见:https://www.yijunzhao.cc/about,转载请注明出处。

原文链接 https://www.yijunzhao.cc/archives/flink-sql-mysql-to-postgresql-real-time-sync-example

欢迎访问 https://www.yijunzhao.cc/

https://www.yijunzhao.cc/