小易说IT 小易说IT

基于Apache Flink CDC的多数据库实时数据同步项目

基于 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-sync

  • Maven GroupId: com.rongzer

  • 父工程 ArtifactId: zct-flink-cdc-sync

  • Java 基础包名: 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                   # 日志配置

模块说明

模块

职责

依赖

flink-mysql-sync-common

数据模型、常量、工具类、异常定义

flink-mysql-sync-core

配置类、转换逻辑、监控指标、健康检查、路由、DLQ

common

flink-mysql-sync-connector

Source/Sink 实现、连接管理、表结构管理

common, core

flink-mysql-sync-starter

应用程序入口、打包

所有模块

支持的数据库

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,转载请注明出处。

原文链接 https://www.yijunzhao.cc/archives/apache-flink-cdc-multi-database-real-time-data-sync-project

欢迎访问 https://www.yijunzhao.cc/

https://www.yijunzhao.cc/