目录
- 一、这段源码处在整个同步流程的哪一步?
- 二、从命令行到 Action:先看完整调用链
- 三、六个核心类分别负责什么?
- 四、FlinkActions:统一入口,只负责启动
- 五、ActionFactory:根据命令名称找到具体工厂
- 1. 统一 action 名称格式
- 2. 通过 FactoryUtil 发现工厂
- 3. 解析重复出现的参数
- 六、为什么要把 Factory 分成三层?
- 七、SynchronizationActionFactoryBase:解析三类通用配置
- 1. 检查 CDC Source 配置
- 2. 分别解析 Catalog 和 Source 配置
- 3. 创建 Action,再补充目标表配置
- 八、SyncDatabaseActionFactoryBase:补齐整库同步规则
- 九、MySqlSyncDatabaseActionFactory:处理 MySQL 特有差异
- 1. 声明 Action 标识
- 2. 指定源端配置参数
- 3. 创建 MySqlSyncDatabaseAction
- 4. 解析 MySQL 专属参数
- 十、DIVIDED 与 COMBINED 有什么区别?
- 十一、Factory 最终组装出了什么?
- 十二、为什么不在 main 方法中一次写完?
- 1. 让 Action 可以独立扩展
- 2. 分离参数解析与任务执行
- 3. 复用公共逻辑
- 4. 避免不同配置相互污染
- 十三、小结
- 十四、下一篇读什么?
本文基于 Apache Paimon 1.4.2 源码,面向刚接触 Flink 和 Paimon 的读者。我们从一个具体问题出发:执行
mysql_sync_database命令后,Paimon 如何识别这条命令、解析参数,并创建出对应的同步任务?
一、这段源码处在整个同步流程的哪一步?
先看完整链路。一个 MySQL CDC 整库同步任务,从启动到数据最终可查询,大致会经历两个阶段。
第一个阶段是任务构建与提交:
命令行参数 ↓ 识别 Action 类型 ↓ 解析 MySQL、Catalog 和目标表配置 ↓ 创建 MySqlSyncDatabaseAction ↓ 构建 Flink Source、数据处理逻辑和 Paimon Sink ↓ 提交 Flink 作业第二个阶段是Flink 作业运行:
MySQL 全量快照 / Binlog ↓ Flink CDC Source ↓ EventParser 解析数据变更和表结构变更 ↓ Paimon Sink ↓ Writer 写入数据文件 ↓ Committer 提交文件和元数据 ↓ 生成新的 Paimon Snapshot本文聚焦第一阶段的前半段,也就是:
命令行 → ActionFactory → MySqlSyncDatabaseAction这一段代码还没有开始读取 Binlog,也没有写入 Paimon。它的任务是把命令行中的字符串参数整理成一个结构化的 Action 对象,为后续构建 Flink 作业做好准备。
对源码初学者来说,先明确自己处在主流程的哪一段非常重要。否则直接进入 Writer、Committer 或 ConflictDetection,很容易看到很多类,却不知道它们为什么会被调用。
二、从命令行到 Action:先看完整调用链
mysql_sync_database的入口调用链如下:
如果不看类名,可以把这条链路理解为:
接收命令 → 找到处理这条命令的工厂 → 分层解析参数 → 组装同步任务对象 → 运行任务其中最容易混淆的是“工厂”和“Action”:
- Factory负责解析参数并创建对象;
- Action保存任务配置,并在后续构建、提交 Flink 作业。
因此,本篇看到的大部分 Factory 代码都属于任务准备阶段。
三、六个核心类分别负责什么?
| 类 | 主要职责 | 所处阶段 |
|---|---|---|
| FlinkActions | Action 命令的统一 Java 入口 | 接收参数并启动 Action |
| ActionFactory | 根据 action 名称寻找具体工厂 | 完成命令路由 |
| SynchronizationActionFactoryBase | 解析 CDC 同步任务共有的配置 | 处理 Source、Catalog 和目标表配置 |
| SyncDatabaseActionFactoryBase | 解析整库同步共有的参数 | 处理目标库、选表、表名和 Schema 规则 |
| MySqlSyncDatabaseActionFactory | 处理 MySQL 整库同步特有参数 | 创建 MySqlSyncDatabaseAction |
| MySqlSyncDatabaseAction | 表示本次 MySQL 整库同步任务 | 后续负责构建 Flink 数据链路 |
这几个类形成了一个清晰的分工:
FlinkActions 负责启动 ActionFactory 负责找到工厂 三层具体 Factory 负责解析和组装 MySqlSyncDatabaseAction 负责构建并运行任务四、FlinkActions:统一入口,只负责启动
FlinkActions 的 main 方法很短,核心逻辑可以概括为:
Optional<Action>action=ActionFactory.createAction(args);if(action.isPresent()){action.get().run();}它只完成三件事:
- 检查命令行中是否包含 action 名称;
- 调用 ActionFactory 创建具体的 Action;
- Action 创建成功后调用
run()。
为什么入口类要写得这么简单?
因为 Paimon 不只有mysql_sync_database。它还包含单表同步、其他数据库 CDC 同步以及表维护等多种 Action。如果所有分支都写进 main 方法,入口很快就会变成一个庞大的 if-else。
Paimon 的做法是:FlinkActions 只保留统一启动流程,具体差异交给各自的 ActionFactory。
五、ActionFactory:根据命令名称找到具体工厂
假设用户执行下面的命令:
mysql_sync_database\--warehousehdfs:///paimon/warehouse\--databaseods\--mysql_confhostname=127.0.0.1\--mysql_confusername=root\--mysql_confpassword=******\--mysql_confdatabase-name=source_db\--table_confbucket=4数组args中,第一个元素是 action 名称,后面的内容才是该 Action 的参数。
1. 统一 action 名称格式
ActionFactory 首先进行名称标准化:
Stringaction=args[0].toLowerCase().replaceAll("-","_");这段代码带来两个效果:
- action 名称不区分大小写;
- 连字符会被替换为下划线。
因此,mysql-sync-database和mysql_sync_database最终都会匹配到同一个标识。
这一步属于命令行兼容处理,可以减少用户因为书写形式不同而遇到的“找不到 Action”问题。
2. 通过 FactoryUtil 发现工厂
名称处理完成后,ActionFactory 调用FactoryUtil.discoverFactory(),寻找 identifier 为mysql_sync_database的工厂。
Paimon 在这里采用了 Java SPI 的扩展方式。模块中的服务注册文件为:
META-INF/services/org.apache.paimon.factories.Factory文件中注册了:
org.apache.paimon.flink.action.cdc.mysql.MySqlSyncDatabaseActionFactoryFactoryUtil 会扫描这些注册信息,再通过每个 Factory 的identifier()判断谁负责当前命令。
这样设计的好处是:增加新的 Action 时,通常只需要新增实现类并完成 SPI 注册,不必修改 FlinkActions 入口。
3. 解析重复出现的参数
去掉 action 名称后,其余参数会被包装成 MultipleParameterToolAdapter。
之所以不能简单地转成一个普通 Map,是因为某些参数允许重复出现。例如:
--mysql_confhostname=127.0.0.1--mysql_confusername=root--mysql_confdatabase-name=source_db多个--mysql_conf最终会被合并为一组 MySQL CDC Source 配置。
六、为什么要把 Factory 分成三层?
MySqlSyncDatabaseActionFactory 的继承关系如下:
ActionFactory └── SynchronizationActionFactoryBase └── SyncDatabaseActionFactoryBase └── MySqlSyncDatabaseActionFactory这三层并不是为了增加代码复杂度,而是在按“通用程度”拆分职责:
- SynchronizationActionFactoryBase 处理所有 CDC 同步任务都需要的配置;
- SyncDatabaseActionFactoryBase 处理所有整库同步都需要的配置;
- MySqlSyncDatabaseActionFactory 只处理 MySQL 特有的配置。
越靠上的父类,逻辑越通用;越靠下的子类,逻辑越具体。
例如,表过滤和目标表前缀并不是 MySQL 独有能力,所以放在整库同步父类中。merge_shards与 MySQL 分库分表场景直接相关,因此保留在 MySQL 工厂中。
这种“父类规定流程,子类补充差异”的写法,就是常见的模板方法思想。
七、SynchronizationActionFactoryBase:解析三类通用配置
SynchronizationActionFactoryBase 是 CDC 同步工厂的公共骨架,它的create()方法完成了三个关键步骤。
1. 检查 CDC Source 配置
checkArgument(params.has(cdcConfigIdentifier()),...);cdcConfigIdentifier()由具体子类实现。MySQL 工厂返回的是mysql_conf,因此 MySQL 同步任务必须提供--mysql_conf。
父类不需要知道具体数据源是 MySQL、Kafka 还是其他系统,只需要让子类告诉它“源端配置使用什么参数名”。
2. 分别解析 Catalog 和 Source 配置
this.catalogConfig=catalogConfigMap(params);this.cdcSourceConfig=optionalConfigMap(params,cdcConfigIdentifier());这两组配置有完全不同的用途:
| 配置 | 表示什么 | 后续交给谁使用 |
|---|---|---|
--catalog_conf和--warehouse | Paimon Catalog 配置 | 决定目标库表存放位置和元数据管理方式 |
--mysql_conf | MySQL CDC Source 配置 | 决定从哪个 MySQL 实例和数据库读取数据 |
3. 创建 Action,再补充目标表配置
Taction=createAction();action.withTableConfig(optionalConfigMap(params,TABLE_CONF));withParams(params,action);--table_conf表示目标 Paimon 表的公共配置,例如 bucket 数量、changelog producer 或 Sink 并行度。同一次整库同步创建或使用的目标表,会共享这组配置。
这里有一个非常容易混淆的参数:
--database ods:Paimon 的目标数据库;--mysql_conf database-name=source_db:MySQL 的源数据库。
虽然两者都出现了 database,但它们分别属于目标端和源端,不能混为一谈。
八、SyncDatabaseActionFactoryBase:补齐整库同步规则
SyncDatabaseActionFactoryBase 在公共同步工厂的基础上,增加了整库同步所需的参数。
首先,它读取 Paimon 目标数据库:
this.database=params.getRequired(DATABASE);随后,withParams()将整库同步规则写入 Action,主要包括:
table_prefix、table_suffix:为目标表统一添加前缀或后缀;table_mapping:显式指定源表与目标表的映射关系;including_tables、excluding_tables:选择或排除源表;including_dbs、excluding_dbs:选择或排除源数据库;partition_keys、primary_keys:指定目标表分区键和主键;type_mapping:控制 MySQL 类型到 Paimon 类型的映射;computed_column:定义计算列;eager_init:控制是否提前初始化目标表;sync_pkeys_from_source_schema:控制是否从源端 Schema 同步主键信息。
这些能力并不局限于 MySQL,因此统一放在整库同步父类中,供其他数据库类型复用。
九、MySqlSyncDatabaseActionFactory:处理 MySQL 特有差异
经过前两层父类后,通用参数和整库参数已经解析完毕。MySqlSyncDatabaseActionFactory 只需要处理 MySQL 相关的差异。
1. 声明 Action 标识
publicstaticfinalStringIDENTIFIER="mysql_sync_database";FactoryUtil 正是通过这个 identifier,把命令行中的 action 名称与当前工厂对应起来。
2. 指定源端配置参数
protectedStringcdcConfigIdentifier(){returnMYSQL_CONF;}这相当于告诉父类:当前 Action 的 CDC Source 配置来自--mysql_conf。
3. 创建 MySqlSyncDatabaseAction
returnnewMySqlSyncDatabaseAction(database,catalogConfig,cdcSourceConfig);创建 Action 时传入了三项核心信息:
database → Paimon 目标数据库 catalogConfig → Paimon Catalog 配置 cdcSourceConfig → MySQL CDC Source 配置目标表配置、过滤规则和其他可选参数,会在 Action 创建后通过一系列withXxx()方法继续补充。
4. 解析 MySQL 专属参数
MySQL 工厂还会处理:
ignore_incompatible:源表与已有 Paimon 表的 Schema 不兼容时,是抛出异常还是忽略该表;merge_shards:不同数据库中的同名分表是否合并到一张 Paimon 表;mode:多表 Sink 使用 DIVIDED 还是 COMBINED 模式;metadata_column:是否把指定的 CDC 元数据列写入目标表。
到这里,命令行参数已经被完整地转换为 MySqlSyncDatabaseAction 的字段。
十、DIVIDED 与 COMBINED 有什么区别?
MySQL 整库同步支持两种多表 Sink 模式:
| 模式 | Sink 组织方式 | 任务启动后出现新表时 |
|---|---|---|
| DIVIDED | 每张表建立独立 Sink | 需要重启任务才能同步新表 |
| COMBINED | 所有表共用一个组合 Sink | 可以自动发现并同步新表 |
MySqlSyncDatabaseAction 默认使用 DIVIDED。
源码帮助信息中有一个值得注意的细节:开头的概述写着“任务启动后新建的 MySQL 表不会被包含”,后面的 mode 说明却指出 COMBINED 可以自动同步新表。
结合实际分支逻辑,更准确的理解是:
- DIVIDED 模式下,新增表需要重启任务;
- COMBINED 模式下,任务运行期间可以接入新增表。
这也是阅读源码时常见的情况:帮助文案、默认值和真正的分支逻辑需要相互核对,不能只依据其中一句话下结论。
十一、Factory 最终组装出了什么?
Factory 链执行结束后,得到的 MySqlSyncDatabaseAction 大致包含以下信息:
MySqlSyncDatabaseAction ├── Paimon 目标端 │ ├── database │ └── catalogConfig ├── MySQL 源端 │ └── cdcSourceConfig ├── 目标表公共配置 │ └── tableConfig ├── 选表与表名规则 │ ├── including / excluding │ ├── prefix / suffix │ └── tableMapping ├── Schema 规则 │ ├── primaryKeys │ ├── partitionKeys │ ├── typeMapping │ └── computedColumns └── MySQL 多表同步参数 ├── mergeShards ├── ignoreIncompatible ├── mode └── metadataColumns此时还没有创建 Flink CDC Source,也没有 Paimon Writer。这个 Action 更像一份已经完成结构化整理的“任务配置说明”。
下一步调用action.run()时,Paimon 才会根据这些配置创建执行环境、Source、事件解析器和 Sink,并最终构建出完整的 Flink 作业。
十二、为什么不在 main 方法中一次写完?
回过头看,这套设计主要解决了四个问题。
1. 让 Action 可以独立扩展
SPI 和 identifier 将命令入口与具体实现解耦。增加新的 Action 时,不必不断修改 FlinkActions。
2. 分离参数解析与任务执行
Factory 负责把字符串参数转换成结构化对象,Action 负责构建任务。参数问题可以在作业运行前尽早暴露,运行逻辑也更容易阅读。
3. 复用公共逻辑
CDC 通用参数、整库同步参数和 MySQL 专属参数被放在不同层级。既避免重复代码,也不会把所有数据源差异堆进同一个类。
4. 避免不同配置相互污染
Source、Catalog 和目标表配置分别保存。后续构建组件时,各组件只读取自己需要的配置,源端参数不会误传到目标端。
十三、小结
读完这一段源码,需要记住以下五点:
- FlinkActions 是统一入口,具体命令路由由 ActionFactory 完成。
- FactoryUtil 通过 SPI 和 identifier 找到 MySqlSyncDatabaseActionFactory。
- 三层 Factory 分别处理 CDC 通用参数、整库通用参数和 MySQL 专属参数。
--database表示 Paimon 目标库,--mysql_conf database-name表示 MySQL 源库。- Factory 阶段只是在组装任务对象,真正的数据读取与写入要从
action.run()之后开始。
十四、下一篇读什么?
下一篇将进入 MySqlSyncDatabaseAction 和 SyncDatabaseActionBase,沿着任务构建主线继续阅读:
- Action 如何创建 Flink CDC Source;
recordParse()如何把原始 CDC 消息转换成 Paimon 能处理的事件;buildEventParserFactory()为什么返回解析器工厂,而不是共享一个解析器实例;buildSink()如何把 mode、tables 和全局 tableConfig 传给 Sink Builder;- DataStream 最终如何连接到 Paimon Sink。
等“Source → 事件解析 → Sink”这段主干打通后,再继续阅读 Writer、PrepareCommit、Committer、ConflictDetection 和 Snapshot,就能始终知道每个类在整条读写链路中的位置。
源码版本:Apache Paimon 1.4.2
涉及模块:paimon-flink-action、paimon-flink-common、paimon-flink-cdc