小易说IT 小易说IT

Flink SQL 实时同步代码范例

本项目演示如何使用 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 目录:

  1. flink-connector-jdbc-1.20.3.jar

  2. mysql-connector-java-8.0.30.jar

  3. 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
  • 编辑 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

方法二:手动执行基础版

  1. 启动 Flink 集群

  2. 使用 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)

  1. 配置管理:集中管理数据库连接信息和性能参数

  2. 性能优化

    • 调整 CDC 读取参数,提高快照和变更捕获效率

    • 优化 JDBC 写入参数,提高写入性能

    • 启用迷你批处理,减少网络开销

  3. 数据转换

    • 确保数据完整性,处理空值

    • 统一数据格式,如邮箱小写化

    • 确保数据合法性,如金额为正数

  4. 错误处理

    • 增加重试次数和延迟

    • 配置并行写入,提高可靠性

优化版 (sync_job_optimized.sql)

在基础版基础上增加:

  1. 高级配置优化

    • 状态管理优化(RocksDB 配置)

    • 检查点策略优化

    • 重启策略优化(指数退避)

    • 批处理参数调优

  2. SQL 优化

    • 谓词下推(在源数据库侧过滤)

    • 数据分区优化(DISTRIBUTED BY)

    • 实时聚合计算

    • 增量更新(ON DUPLICATE KEY UPDATE)

  3. 高级功能

    • 实时用户统计分析

    • 数据质量保证

    • 自动错误恢复

    • 详细的监控指标

在优化版基础上增加:

  1. 企业级特性

    • 数据质量监控(自动计算质量分数)

    • 死信队列(DLQ)处理异常数据

    • 实时聚合统计

    • 批量写入优化(rewriteBatchedStatements)

  2. 数据清洗与验证

    • 姓名格式标准化(TRIM + 默认值)

    • 邮箱格式验证(正则表达式)

    • 金额数据验证(负数处理)

    • 时间戳完整性保证

  3. 容错与恢复

    • 检查点外部化保留(RETAIN_ON_CANCELLATION)

    • 增量检查点(RocksDB 增量)

    • 固定延迟重启策略

    • 多层级错误处理

  4. 监控与可观测性

    • 数据质量日志表

    • 异常数据捕获

    • 统计信息实时更新

    • 详细的注释和监控点

  5. 性能调优

    • CDC 分块大小优化(8096)

    • 数据分布因子调优

    • 内存管理优化

    • 网络缓冲区优化

所有 SQL 脚本均经过验证,确保与 Flink 1.20.3 完全兼容:

  • 使用标准的 Flink SQL 语法

  • 使用 Flink 1.20.3 支持的 CDC 连接器参数

  • 使用 Flink 1.20.3 支持的 JDBC 连接器参数

  • 所有配置参数均为 Flink 1.20.3 支持的标准参数

生产环境建议

  1. 使用高级版脚本sync_job_advanced.sql

  2. 配置环境变量:使用 ${VAR_NAME} 方式配置敏感信息

  3. 启用 SSL:配置 SSL 连接保证数据安全

  4. 监控告警:配置 Prometheus + Grafana 监控

  5. 定期检查:定期验证数据一致性

注意事项

  • 确保 MySQL 5.7 已启用 binlog 且格式为 ROW

  • 确保源数据库和目标数据库的网络连通性

  • 根据实际数据量调整性能参数

  • 定期监控同步作业状态,确保数据一致性

  • 生产环境建议使用环境变量管理敏感配置