Apache Paimon 源码导读(一):MySQL CDC Action 从命令行到任务创建

发布时间:2026/8/26 18:38:16
Apache Paimon 源码导读(一):MySQL CDC Action 从命令行到任务创建 目录一、这段源码处在整个同步流程的哪一步二、从命令行到 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. 创建 MySqlSyncDatabaseAction4. 解析 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.mainActionFactory.createActionFactoryUtil.discoverFactoryMySqlSyncDatabaseActionFactorySynchronizationActionFactoryBase.createSyncDatabaseActionFactoryBase 解析整库参数MySQL 工厂补充专属参数创建 MySqlSyncDatabaseActionaction.run 进入下一阶段如果不看类名可以把这条链路理解为接收命令 → 找到处理这条命令的工厂 → 分层解析参数 → 组装同步任务对象 → 运行任务其中最容易混淆的是“工厂”和“Action”Factory负责解析参数并创建对象Action保存任务配置并在后续构建、提交 Flink 作业。因此本篇看到的大部分 Factory 代码都属于任务准备阶段。三、六个核心类分别负责什么类主要职责所处阶段FlinkActionsAction 命令的统一 Java 入口接收参数并启动 ActionActionFactory根据 action 名称寻找具体工厂完成命令路由SynchronizationActionFactoryBase解析 CDC 同步任务共有的配置处理 Source、Catalog 和目标表配置SyncDatabaseActionFactoryBase解析整库同步共有的参数处理目标库、选表、表名和 Schema 规则MySqlSyncDatabaseActionFactory处理 MySQL 整库同步特有参数创建 MySqlSyncDatabaseActionMySqlSyncDatabaseAction表示本次 MySQL 整库同步任务后续负责构建 Flink 数据链路这几个类形成了一个清晰的分工FlinkActions 负责启动 ActionFactory 负责找到工厂 三层具体 Factory 负责解析和组装 MySqlSyncDatabaseAction 负责构建并运行任务四、FlinkActions统一入口只负责启动FlinkActions 的 main 方法很短核心逻辑可以概括为OptionalActionactionActionFactory.createAction(args);if(action.isPresent()){action.get().run();}它只完成三件事检查命令行中是否包含 action 名称调用 ActionFactory 创建具体的 ActionAction 创建成功后调用run()。为什么入口类要写得这么简单因为 Paimon 不只有mysql_sync_database。它还包含单表同步、其他数据库 CDC 同步以及表维护等多种 Action。如果所有分支都写进 main 方法入口很快就会变成一个庞大的 if-else。Paimon 的做法是FlinkActions 只保留统一启动流程具体差异交给各自的 ActionFactory。五、ActionFactory根据命令名称找到具体工厂假设用户执行下面的命令mysql_sync_database\--warehousehdfs:///paimon/warehouse\--databaseods\--mysql_confhostname127.0.0.1\--mysql_confusernameroot\--mysql_confpassword******\--mysql_confdatabase-namesource_db\--table_confbucket4数组args中第一个元素是 action 名称后面的内容才是该 Action 的参数。1. 统一 action 名称格式ActionFactory 首先进行名称标准化Stringactionargs[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_confhostname127.0.0.1--mysql_confusernameroot--mysql_confdatabase-namesource_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.catalogConfigcatalogConfigMap(params);this.cdcSourceConfigoptionalConfigMap(params,cdcConfigIdentifier());这两组配置有完全不同的用途配置表示什么后续交给谁使用--catalog_conf和--warehousePaimon Catalog 配置决定目标库表存放位置和元数据管理方式--mysql_confMySQL CDC Source 配置决定从哪个 MySQL 实例和数据库读取数据3. 创建 Action再补充目标表配置TactioncreateAction();action.withTableConfig(optionalConfigMap(params,TABLE_CONF));withParams(params,action);--table_conf表示目标 Paimon 表的公共配置例如 bucket 数量、changelog producer 或 Sink 并行度。同一次整库同步创建或使用的目标表会共享这组配置。这里有一个非常容易混淆的参数--database odsPaimon 的目标数据库--mysql_conf database-namesource_dbMySQL 的源数据库。虽然两者都出现了 database但它们分别属于目标端和源端不能混为一谈。八、SyncDatabaseActionFactoryBase补齐整库同步规则SyncDatabaseActionFactoryBase 在公共同步工厂的基础上增加了整库同步所需的参数。首先它读取 Paimon 目标数据库this.databaseparams.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 标识publicstaticfinalStringIDENTIFIERmysql_sync_database;FactoryUtil 正是通过这个 identifier把命令行中的 action 名称与当前工厂对应起来。2. 指定源端配置参数protectedStringcdcConfigIdentifier(){returnMYSQL_CONF;}这相当于告诉父类当前 Action 的 CDC Source 配置来自--mysql_conf。3. 创建 MySqlSyncDatabaseActionreturnnewMySqlSyncDatabaseAction(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 SourcerecordParse()如何把原始 CDC 消息转换成 Paimon 能处理的事件buildEventParserFactory()为什么返回解析器工厂而不是共享一个解析器实例buildSink()如何把 mode、tables 和全局 tableConfig 传给 Sink BuilderDataStream 最终如何连接到 Paimon Sink。等“Source → 事件解析 → Sink”这段主干打通后再继续阅读 Writer、PrepareCommit、Committer、ConflictDetection 和 Snapshot就能始终知道每个类在整条读写链路中的位置。源码版本Apache Paimon 1.4.2涉及模块paimon-flink-action、paimon-flink-common、paimon-flink-cdc