本项目演示如何使用 Flink SQL 1.20.3 实现从 MySQL 5.7 到 MySQL 8.0 的实时同步,支持部分表、部分字段过滤和动态取值过滤。
环境要求
Flink 1.20.3
MySQL 5.7 (源数据库)
MySQL 8.0 (目标数据库)
JDK 1.8+
项目结构
flinkCMD/
├── config/
│ ├── flink-conf.yaml
│ ├── sql-client-defaults.yaml
│ ├── sync-config.sql # 基础配置参数
│ ├── sync-config-optimized.sql # 优化版配置参数
│ └── sync-config-advanced.sql # 高级版配置参数
├── lib/
│ └── # 依赖jar包
├── sql/
│ ├── source_ddl.sql # 源表 DDL
│ ├── sink_ddl.sql # 目标表 DDL
│ ├── sync_job.sql # 基础同步作业 SQL
│ ├── sync_job_optimized.sql # 优化版同步作业 SQL
│ └── sync_job_advanced.sql # 高级版同步作业 SQL
├── deploy.sh # 部署脚本
└── README.md
依赖jar包
需要下载以下依赖jar包到 lib 目录:
flink-connector-jdbc-1.20.3.jar
mysql-connector-java-8.0.30.jar
flink-sql-connector-mysql-cdc-2.4.0.jar
配置步骤
1. 配置 MySQL 源数据库
在 MySQL 5.7 中启用 binlog:
# my.cnf
server-id = 1
binlog-format = ROW
enable-binlog = true
2. 配置 Flink
编辑
config/flink-conf.yaml文件,设置必要的配置根据需求选择配置文件:
config/sync-config.sql- 基础配置config/sync-config-optimized.sql- 优化版配置config/sync-config-advanced.sql- 高级版配置
3. 执行 SQL 脚本
方法一:使用部署脚本(推荐)
chmod +x deploy.sh
./deploy.sh
方法二:手动执行基础版
启动 Flink 集群
使用 SQL Client 执行以下命令:
./bin/sql-client.sh embedded -f sql/source_ddl.sql
./bin/sql-client.sh embedded -f sql/sink_ddl.sql
./bin/sql-client.sh embedded -f sql/sync_job.sql
方法三:使用优化版脚本
./bin/sql-client.sh embedded -f sql/sync_job_optimized.sql
方法四:使用高级版脚本(推荐生产环境)
./bin/sql-client.sh embedded -f sql/sync_job_advanced.sql
版本特性对比
基础版 (sync_job.sql)
配置管理:集中管理数据库连接信息和性能参数
性能优化:
调整 CDC 读取参数,提高快照和变更捕获效率
优化 JDBC 写入参数,提高写入性能
启用迷你批处理,减少网络开销
数据转换:
确保数据完整性,处理空值
统一数据格式,如邮箱小写化
确保数据合法性,如金额为正数
错误处理:
增加重试次数和延迟
配置并行写入,提高可靠性
优化版 (sync_job_optimized.sql)
在基础版基础上增加:
高级配置优化:
状态管理优化(RocksDB 配置)
检查点策略优化
重启策略优化(指数退避)
批处理参数调优
SQL 优化:
谓词下推(在源数据库侧过滤)
数据分区优化(DISTRIBUTED BY)
实时聚合计算
增量更新(ON DUPLICATE KEY UPDATE)
高级功能:
实时用户统计分析
数据质量保证
自动错误恢复
详细的监控指标
高级版 (sync_job_advanced.sql) - Flink 1.20.3 完全兼容
在优化版基础上增加:
企业级特性:
数据质量监控(自动计算质量分数)
死信队列(DLQ)处理异常数据
实时聚合统计
批量写入优化(rewriteBatchedStatements)
数据清洗与验证:
姓名格式标准化(TRIM + 默认值)
邮箱格式验证(正则表达式)
金额数据验证(负数处理)
时间戳完整性保证
容错与恢复:
检查点外部化保留(RETAIN_ON_CANCELLATION)
增量检查点(RocksDB 增量)
固定延迟重启策略
多层级错误处理
监控与可观测性:
数据质量日志表
异常数据捕获
统计信息实时更新
详细的注释和监控点
性能调优:
CDC 分块大小优化(8096)
数据分布因子调优
内存管理优化
网络缓冲区优化
Flink 1.20.3 兼容性说明
所有 SQL 脚本均经过验证,确保与 Flink 1.20.3 完全兼容:
使用标准的 Flink SQL 语法
使用 Flink 1.20.3 支持的 CDC 连接器参数
使用 Flink 1.20.3 支持的 JDBC 连接器参数
所有配置参数均为 Flink 1.20.3 支持的标准参数
生产环境建议
使用高级版脚本:
sync_job_advanced.sql配置环境变量:使用
${VAR_NAME}方式配置敏感信息启用 SSL:配置 SSL 连接保证数据安全
监控告警:配置 Prometheus + Grafana 监控
定期检查:定期验证数据一致性
注意事项
确保 MySQL 5.7 已启用 binlog 且格式为 ROW
确保源数据库和目标数据库的网络连通性
根据实际数据量调整性能参数
定期监控同步作业状态,确保数据一致性
生产环境建议使用环境变量管理敏感配置