场景说明: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
⚠️ 前置准备
MySQL 开启 binlog,binlog_format=ROW;创建 CDC 账号,授予 REPLICATION SLAVE, REPLICATION CLIENT, SELECT 权限
PostgreSQL 提前建好目标表,表结构必须和源表对齐(字段名、类型、主键)
源表必须有主键(CDC 需要主键,无主键只能做快照,无法消费增量变更)
Flink SQL 完整脚本
-- 开启 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不支持动态目标表名。两种方案:
多个独立 Flink 任务,单表同步(稳定,生产首选)
使用 Flink CDC Pipeline(
flink-cdc-pipeline-connector),声明式配置整库同步,自动映射表名,不需要写 SQL
关键参数说明
MySQL CDC 参数
PostgreSQL JDBC Sink 参数
JDBC Sink 是 UPSERT 模式,要求目标表定义主键,收到 DELETE 消息会执行
DELETE FROM table WHERE pk=?
常见坑点
时区问题:MySQL CDC 必须指定
server-time-zone='Asia/Shanghai',否则时间字段偏移 8 小时主键缺失:源表无主键,CDC 无法识别 UPDATE/DELETE,只能 INSERT,不能做 UPSERT 同步
Checkpoint 必须开启:Exactly-Once 依赖 checkpoint,checkpoint 失败会导致重复写入
类型映射
MySQL DATETIME → Flink TIMESTAMP(3) → PostgreSQL TIMESTAMPTZ / TIMESTAMP
MySQL BIGINT → Flink BIGINT → PostgreSQL BIGINT
MySQL VARCHAR → Flink STRING → PostgreSQL VARCHAR/TEXT
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
备选方案:Flink CDC Pipeline(声明式 YAML,无需写 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,转载请注明出处。
原文链接
欢迎访问