基于 Flink 2.2.1 + Flink CDC 3.6.0-2.2 实现的多数据库实时数据同步项目。
功能特性
实时同步: 基于 CDC 实现毫秒级延迟的数据同步
多数据库支持: 支持 MySQL、PostgreSQL 等多种数据库作为 Source 和 Sink
动态字段过滤: 支持白名单(include)和黑名单(exclude)两种字段过滤模式
字段映射: 支持源表和目标表之间的字段名映射
数据过滤: 支持基于字段值的条件过滤
全量+增量: 支持全量历史数据同步和增量实时同步
断点续传: 基于 Checkpoint 机制实现故障恢复
幂等写入: 使用 UPSERT 语义保证数据一致性
批量写入优化: 自适应批量大小调整,提升写入性能
多表路由: 支持单表到多表、多表到单表的数据路由
DDL同步: 支持表结构变更自动同步
死信队列: 失败消息自动进入死信队列,便于问题排查
配置热更新: 支持配置文件修改后自动热加载
健康检查: 内置健康检查机制,支持监控集成
密码加密: 支持 Jasypt 加密敏感配置
高可用: 支持 Flink 集群部署和自动故障转移
多模块架构: 采用 Maven 多模块设计,职责清晰,易于扩展
项目信息
项目名称:
zct-flink-cdc-syncMaven GroupId:
com.rongzer父工程 ArtifactId:
zct-flink-cdc-syncJava 基础包名:
com.rongzer.flink启动类:
com.rongzer.flink.DataSyncJob/com.rongzer.flink.MySqlSyncJob模块命名: 当前子模块仍沿用
flink-mysql-sync-*命名
项目结构
zct-flink-cdc-sync/
├── pom.xml # 父 POM,统一管理版本与模块
├── README.md # 项目说明
├── doc/
│ └── TechArch.md # 技术架构说明
├── flink-mysql-sync-common/ # 公共模块:模型、常量、异常
│ └── src/main/java/com/rongzer/flink/
│ ├── common/ # 公共常量与工具
│ │ ├── Constants.java # 全局常量
│ │ └── TraceIdGenerator.java # TraceId 生成器
│ ├── exception/
│ │ └── SyncException.java # 统一同步异常
│ └── model/ # 事件与映射模型
│ ├── DataChangeEvent.java # CDC 数据变更事件
│ ├── FieldMapping.java # 字段映射模型
│ └── SourcePosition.java # CDC 位点信息
├── flink-mysql-sync-core/ # 核心模块:配置、Pipeline、过滤与监控
│ └── src/main/java/com/rongzer/flink/
│ ├── config/ # 同步配置模型
│ │ ├── SyncJobConfig.java # 作业总配置
│ │ ├── SourceConfig.java # Source 配置
│ │ ├── SinkConfig.java # Sink 配置
│ │ ├── MetricsConfig.java # 指标配置
│ │ └── loader/ # 配置加载与分域解析
│ │ ├── ConfigLoader.java # 配置加载门面
│ │ ├── ConfigParser.java # 配置解析编排门面
│ │ ├── SourceConfigParser.java # Source 配置解析
│ │ ├── SinkConfigParser.java # Sink 配置解析
│ │ ├── FilterConfigParser.java # 过滤配置解析
│ │ ├── MetricsConfigParser.java # 指标配置解析
│ │ ├── ConfigValueReader.java # 配置取值工具
│ │ ├── ConfigWatcher.java # 配置文件监听
│ │ └── WhereConditionParser.java # where 条件解析
│ ├── pipeline/ # Pipeline 编排抽象
│ │ ├── DataSyncPipeline.java # 通用数据同步 Pipeline
│ │ ├── MySqlSyncPipeline.java # 兼容旧入口的包装类
│ │ ├── PipelineFactory.java # Source/Sink Builder 注册工厂
│ │ ├── SourceBuilder.java # Source Builder 接口
│ │ └── SinkBuilder.java # Sink Builder 接口
│ ├── transform/ # 字段与条件过滤
│ ├── router/ # 表路由
│ ├── metrics/ # 延迟等运行指标
│ ├── health/ # 健康检查
│ ├── ddl/ # DDL 事件处理
│ └── dlq/ # 死信队列
├── flink-mysql-sync-connector/ # 连接器模块:CDC Source 与 JDBC Sink
│ └── src/main/java/com/rongzer/flink/
│ ├── pipeline/ # Source/Sink Builder 实现
│ │ ├── MySqlCdcSourceBuilder.java # MySQL CDC Source Builder
│ │ ├── PostgresCdcSourceBuilder.java # PostgreSQL CDC Source Builder
│ │ ├── JdbcSinkBuilder.java # MySQL/JDBC Sink Builder
│ │ └── PostgresSinkBuilder.java # PostgreSQL Sink Builder
│ ├── sink/ # JDBC Sink 实现
│ │ ├── AbstractJdbcSink.java # JDBC Sink 公共基类
│ │ ├── PooledJdbcSink.java # 连接池 JDBC Sink
│ │ ├── PostgresJdbcSink.java # PostgreSQL 专用 Sink
│ │ ├── SyncResultMetricsSink.java # 同步结果统计 Sink
│ │ ├── dialect/ # 数据库 SQL 方言
│ │ │ ├── JdbcDialect.java # SQL 方言接口
│ │ │ ├── MySqlJdbcDialect.java # MySQL 方言
│ │ │ └── PostgresJdbcDialect.java # PostgreSQL 方言
│ │ └── metrics/ # 同步结果统计服务
│ │ ├── SyncResultReporter.java # 结果上报接口
│ │ ├── FileSyncResultReporter.java # 文件结果上报
│ │ └── RecordCountQueryService.java # 源/目标记录数查询
│ ├── schema/ # 表结构管理
│ ├── resource/ # 数据源资源管理
│ └── util/ # 密码加密工具
└── flink-mysql-sync-starter/ # 启动模块:应用入口与示例配置
└── src/main/
├── java/com/rongzer/flink/
│ ├── DataSyncJob.java # 通用数据同步主程序
│ ├── MySqlSyncJob.java # 兼容旧 MySQL 入口
│ ├── bootstrap/ # 启动公共服务
│ │ ├── JobConfigBootstrapService.java # 配置加载、覆盖与解密编排
│ │ ├── CommandLineConfigOverride.java # 命令行参数覆盖
│ │ ├── PasswordDecryptService.java # 密码解密 fail-fast
│ │ ├── PipelineFactoryProvider.java # PipelineFactory 注册
│ │ └── DefaultJobConfigFactory.java # 本地开发默认配置
│ └── util/ # 加密命令行工具与示例
└── resources/
├── application.properties # MySQL 示例配置
├── application-postgres.properties # PostgreSQL 示例配置
└── log4j2.xml # 日志配置模块说明
支持的数据库
Source (CDC)
✅ MySQL CDC (flink-connector-mysql-cdc)
✅ PostgreSQL CDC (flink-connector-postgres-cdc)
🔄 Oracle CDC (待扩展)
🔄 SQLServer CDC (待扩展)
🔄 Kafka (待扩展)
Sink (JDBC)
✅ MySQL (JDBC)
✅ PostgreSQL (JDBC + ON CONFLICT UPSERT)
🔄 Oracle (待扩展)
🔄 SQLServer (待扩展)
🔄 Kafka (待扩展)
🔄 Elasticsearch (待扩展)
项目源码地址
GitHub:
Gitee:https://gitee.com/yijunzhao/flink-mysql-sync
本文原创作者:易君召,详见:https://www.yijunzhao.cc/about,转载请注明出处。
原文链接
欢迎访问