Flink CDC零基础教程:MySQL到PostgreSQL实时同步全流程
·
Flink CDC零基础教程:MySQL到PostgreSQL实时同步全流程
本教程将逐步指导你使用Flink CDC实现MySQL到PostgreSQL的实时数据同步。整个过程分为5个阶段,无需编程基础,只需基础命令行操作能力。
阶段1:环境准备
-
安装组件(需提前安装):
-
配置Flink:
# 将下载的flink-sql-connector-mysql-cdc.jar和 # flink-sql-connector-postgres-cdc.jar放入Flink的lib目录 cp *.jar /path/to/flink/lib/
阶段2:MySQL数据源配置
-
启用MySQL Binlog:
# 在MySQL执行 SET GLOBAL log_bin = ON; SET GLOBAL binlog_format = 'ROW'; CREATE USER 'flink'@'%' IDENTIFIED BY 'password'; GRANT SELECT, RELOAD, SHOW DATABASES ON *.* TO 'flink'@'%'; -
创建测试表:
CREATE DATABASE test_db; USE test_db; CREATE TABLE users ( id INT PRIMARY KEY, name VARCHAR(50) );
阶段3:PostgreSQL目标库配置
-
创建接收表:
CREATE DATABASE sync_db; \c sync_db CREATE TABLE users ( id INT PRIMARY KEY, name VARCHAR(50) ); -
开放权限:
CREATE USER flink WITH PASSWORD 'password'; GRANT ALL PRIVILEGES ON TABLE users TO flink;
阶段4:Flink SQL作业开发
-
启动Flink SQL客户端:
./bin/sql-client.sh -
执行同步SQL:
-- 创建MySQL数据源表 CREATE TABLE mysql_users ( id INT, name STRING, PRIMARY KEY (id) NOT ENFORCED ) WITH ( 'connector' = 'mysql-cdc', 'hostname' = 'localhost', 'port' = '3306', 'username' = 'flink', 'password' = 'password', 'database-name' = 'test_db', 'table-name' = 'users' ); -- 创建PostgreSQL目标表 CREATE TABLE pg_users ( id INT, name STRING, PRIMARY KEY (id) NOT ENFORCED ) WITH ( 'connector' = 'jdbc', 'url' = 'jdbc:postgresql://localhost:5432/sync_db', 'username' = 'flink', 'password' = 'password', 'table-name' = 'users' ); -- 启动同步任务 INSERT INTO pg_users SELECT * FROM mysql_users;
阶段5:验证与监控
-
数据验证:
- 在MySQL插入数据:
INSERT INTO users VALUES (1, 'Alice'), (2, 'Bob'); - 在PostgreSQL检查:
SELECT * FROM users; -- 应实时显示Alice和Bob
- 在MySQL插入数据:
-
任务监控:
- 访问Flink Web UI(默认8081端口)
- 查看
Running Jobs状态和numRecordsOut指标
常见问题处理
-
Binlog权限错误:
Access denied; you need (at least one of) the SUPER privilege(s)解决方案:确认MySQL用户拥有
RELOAD权限 -
数据延迟高:
- 增加Flink任务并行度
- 检查网络带宽
-
字段类型映射错误:
- 在DDL中显式指定类型,例如:
name STRING COMMENT '映射VARCHAR(50)'
- 在DDL中显式指定类型,例如:
提示:生产环境建议开启Checkpoint(在Flink配置中设置
execution.checkpointing.interval: 5000)
通过以上步骤,你已实现完整的MySQL到PostgreSQL实时同步管道。任何数据变更(增删改)将在秒级内完成同步。
更多推荐
所有评论(0)