☰
Apache Paimon 源码导读(一):MySQL CDC Action 从命令行到任务创建
2026/9/25 6:32:04 网站建设 项目流程

目录

    • 一、这段源码处在整个同步流程的哪一步?
    • 二、从命令行到 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的入口调用链如下:

mysql_sync_database 命令

FlinkActions.main

ActionFactory.createAction

FactoryUtil.discoverFactory

MySqlSyncDatabaseActionFactory

SynchronizationActionFactoryBase.create

SyncDatabaseActionFactoryBase 解析整库参数

MySQL 工厂补充专属参数

创建 MySqlSyncDatabaseAction

action.run 进入下一阶段

如果不看类名,可以把这条链路理解为:

接收命令 → 找到处理这条命令的工厂 → 分层解析参数 → 组装同步任务对象 → 运行任务

其中最容易混淆的是“工厂”和“Action”:

  • Factory负责解析参数并创建对象;
  • Action保存任务配置,并在后续构建、提交 Flink 作业。

因此,本篇看到的大部分 Factory 代码都属于任务准备阶段。

三、六个核心类分别负责什么?

类主要职责所处阶段
FlinkActionsAction 命令的统一 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();}

它只完成三件事:

  1. 检查命令行中是否包含 action 名称;
  2. 调用 ActionFactory 创建具体的 Action;
  3. 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.MySqlSyncDatabaseActionFactory

FactoryUtil 会扫描这些注册信息,再通过每个 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和--warehousePaimon Catalog 配置决定目标库表存放位置和元数据管理方式
--mysql_confMySQL 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 和目标表配置分别保存。后续构建组件时,各组件只读取自己需要的配置,源端参数不会误传到目标端。

十三、小结

读完这一段源码,需要记住以下五点:

  1. FlinkActions 是统一入口,具体命令路由由 ActionFactory 完成。
  2. FactoryUtil 通过 SPI 和 identifier 找到 MySqlSyncDatabaseActionFactory。
  3. 三层 Factory 分别处理 CDC 通用参数、整库通用参数和 MySQL 专属参数。
  4. --database表示 Paimon 目标库,--mysql_conf database-name表示 MySQL 源库。
  5. 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

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询