Apache Paimon 源码导读(二):Action 如何构建并提交 Flink 作业

发布时间:2026/8/26 18:36:41

Apache Paimon 源码导读(二):Action 如何构建并提交 Flink 作业
目录一、Action 对象已经有了Flink 作业还没有运行二、run() 到底定义在哪个类里三、创建 Action 时执行环境已经准备好了为什么开启 Object Reuse四、run() 实际只做三件事1. 没有配置 Checkpoint 时默认开启三分钟一次五、build() 不只是“画流程图”第一步检查 MySQL 必填配置第二步确保 Paimon 目标数据库存在第三步发现源表并准备目标表六、把 Source、Parse 和 Sink 串起来1. buildSource() 创建 MySQL CDC Source2. env.fromSource() 把 Source 加进拓扑3. flatMap() 添加 Parse 算子4. buildSink() 把 Paimon Sink 接到数据流后面七、为什么 DataStream 创建完了数据还没有流动八、env.execute() 才是真正的提交按钮九、源码里的三个 build 不要混淆十、一个容易忽略的分布式细节CatalogLoader十一、完整调用链再走一遍十二、读完这一段需要记住什么上一篇我们跟到了 MySqlSyncDatabaseAction 的创建过程。Factory 已经把命令行参数整理进 Action 对象但这时 Flink 作业还没有运行MySQL Binlog 也没有开始读取。接下来真正要弄清楚的是Action 什么时候开始搭建 Flink 作业又是哪一行代码把作业提交到集群这两个问题的答案分别是build()和env.execute()。它们看起来挨得很近作用却完全不同。一、Action 对象已经有了Flink 作业还没有运行先把上一篇和这一篇连起来命令行参数 ↓ ActionFactory 解析参数 ↓ 创建 MySqlSyncDatabaseAction ↓ action.run() ↓ build()搭建 Source → Parse → Sink ↓ env.execute()提交 Flink 作业 ↓ TaskManager 开始真正处理数据Action 对象可以理解为一份已经填好参数的任务配置。它知道数据从哪个 MySQL 数据库读取数据写到哪个 Paimon 数据库哪些表需要同步目标表使用什么配置使用 DIVIDED 还是 COMBINED 模式。但“知道怎么做”和“已经开始做”是两回事。只有执行到env.execute()Flink 才会真正提交作业。二、run() 到底定义在哪个类里MySqlSyncDatabaseAction 的继承关系如下MySqlSyncDatabaseAction ↓ SyncDatabaseActionBase ↓ SynchronizationActionBase ↓ ActionBase ↓ Action沿着 MySqlSyncDatabaseAction 往下看会发现它自己没有实现run()。Java 会沿继承关系向父类查找最终调用的是 SynchronizationActionBase 中的run()而不是 ActionBase 中那个更通用的实现。几个关键方法的实际归属如下方法实际实现位置作用run()SynchronizationActionBase控制 CDC 作业的整体启动顺序build()SynchronizationActionBase组织 Source、Parse 和 SinkbeforeBuildingSourceSink()MySqlSyncDatabaseAction发现 MySQL 表并准备 Paimon 表buildSource()MySqlSyncDatabaseAction创建 MySQL CDC SourcerecordParse()SyncDatabaseActionBase提供 CDC 消息解析函数buildSink()SyncDatabaseActionBase把 Paimon Sink 接到数据流后面execute()ActionBase调用 Flink 的env.execute()这套结构的特点是父类规定步骤子类提供每一步的具体做法。MySQL、PostgreSQL、Kafka 等数据源可以共用同一个启动框架只替换有差异的部分。三、创建 Action 时执行环境已经准备好了MySqlSyncDatabaseAction 的构造函数最终会调用 ActionBase 的构造函数publicActionBase(MapString,StringcatalogConfig){catalogOptionsOptions.fromMap(catalogConfig);initCatalog();initFlinkEnv(StreamExecutionEnvironment.getExecutionEnvironment());}这里主要准备了两类对象。第一类是 CatalogPaimon Catalog └── 负责访问和管理 Paimon 数据库、表及元数据 FlinkCatalog └── 让 Flink Table API 能识别这个 Paimon Catalog第二类是 Flink 执行环境StreamExecutionEnvironmentenv可以把env理解成搭建 Flink 作业的工作台。后面的 Source、数据转换和 Sink都会被逐个放到这个工作台上。需要注意的是获得StreamExecutionEnvironment并不代表已经连接集群并开始运行。它只是提供了一套 API用来描述作业应该长什么样。为什么开启 Object ReuseActionBase 初始化环境时还调用了env.getConfig().enableObjectReuse();流式任务会持续处理大量记录。如果每经过一个算子都创建新对象会增加内存分配和垃圾回收压力。开启对象复用后Flink 可以重复使用部分对象从而降低开销。对象复用也有风险下游算子不能随意长期保存输入对象的引用因为这个对象后面可能会被改写。源码注释说明Paimon 会在需要时自行复制不能安全复用的对象。四、run() 实际只做三件事SynchronizationActionBase 的run()很短Overridepublicvoidrun()throwsException{if(!env.getCheckpointConfig().isCheckpointingEnabled()){env.enableCheckpointing(DEFAULT_CHECKPOINT_INTERVAL);}build();execute(syncJobHandler.provideDefaultJobName());}它做了三件事检查 Checkpoint ↓ 构建 Flink 作业 ↓ 提交 Flink 作业1. 没有配置 Checkpoint 时默认开启三分钟一次默认值是privatestaticfinallongDEFAULT_CHECKPOINT_INTERVAL3*60*1000;也就是 180 秒。如果用户已经配置了 CheckpointPaimon 会沿用用户配置不会强行覆盖。只有完全没开启时才使用这个默认值。CDC 是持续运行的流式任务。Checkpoint 主要用于保存MySQL CDC Source 已经读到的 Binlog 位置Flink 有状态算子的运行状态Source 与 Sink 之间的一致性进度Paimon Sink 中等待提交的数据状态。这里顺便区分两个名字很像的概念Flink CheckpointFlink 作业的状态快照用于故障恢复和端到端一致性Paimon SnapshotPaimon 表提交成功后产生的数据版本。它们不是同一个东西但 Paimon 流式 Sink 通常会跟随 Flink Checkpoint 完成一致性提交。五、build() 不只是“画流程图”SynchronizationActionBase 的build()是这一阶段的核心Overridepublicvoidbuild()throwsException{syncJobHandler.checkRequiredOption();catalog.createDatabase(database,true);beforeBuildingSourceSink();DataStreamRichCdcMultiplexRecordinputbuildDataStreamSource(buildSource()).flatMap(recordParse()).name(Parse);EventParser.FactoryRichCdcMultiplexRecordparserFactorybuildEventParserFactory();buildSink(input,parserFactory);}这个方法可以分成两部分提交前准备 ├── 校验参数 ├── 创建 Paimon 数据库 └── 发现 MySQL 表准备 Paimon 表 构建 Flink 拓扑 ├── 添加 Source ├── 添加 Parse └── 添加 Sink第一步检查 MySQL 必填配置syncJobHandler.checkRequiredOption();MySQL 整库同步会检查 hostname、username、password 和 database-name 等必要配置。Factory 阶段已经把参数装进了 Action这里再检查一次是为了在提交作业前确认这些参数能满足具体 Source 的要求。整库同步还会拒绝 table-name 配置。因为指定单表属于mysql_sync_table的职责不应该混进mysql_sync_database。第二步确保 Paimon 目标数据库存在catalog.createDatabase(database,true);第二个参数为true表示数据库已经存在时不用报错。这一步发生在env.execute()之前所以它是当前进程直接执行的 Catalog 操作并不是一个 Flink 算子。第三步发现源表并准备目标表beforeBuildingSourceSink();MySQL 实现会在这里连接源库并读取表结构大致流程如下读取 MySQL 表结构 ↓ 应用 including / excluding 过滤规则 ↓ 排除没有主键的表 ↓ 转换成 Paimon Schema ↓ 目标表不存在创建 Paimon 表 目标表已存在检查 Schema 是否兼容 ↓ 得到本次真正需要同步的表列表这说明build()不完全是内存里的“建图”操作。它还会访问 MySQL 和 Paimon Catalog完成提交前的元数据准备。这样做的好处是能尽早发现问题。例如 MySQL 连接失败、找不到表或者 Schema 不兼容都可以在作业提交前直接报出来而不是等 TaskManager 启动以后才失败。六、把 Source、Parse 和 Sink 串起来提交前准备完成后代码开始搭建 Flink 数据链路DataStreamRichCdcMultiplexRecordinputbuildDataStreamSource(buildSource()).flatMap(recordParse()).name(Parse);拆开看会更清楚。1. buildSource() 创建 MySQL CDC SourceMySqlSyncDatabaseAction 的实现是OverrideprotectedMySqlSourceCdcSourceRecordbuildSource(){validateRuntimeExecutionMode();returnMySqlActionUtils.buildMySqlSource(cdcSourceConfig,tableList(...),typeMapping);}它首先确认 Flink 当前使用 STREAMING 模式然后根据 MySQL 配置和待同步表列表创建 MySqlSource。这里仍然只是创建 Source 对象没有开始读取 MySQL。2. env.fromSource() 把 Source 加进拓扑buildDataStreamSource()最终会调用env.fromSource(source,watermarkStrategy,syncJobHandler.provideSourceName());对 MySQL 来说Source 名称是MySQL Source这个名称会显示在 Flink Web UI 中。如果目标表配置启用了基于 Watermark 的自动 Tag 创建这里还会配置 CDC Watermark否则使用WatermarkStrategy.noWatermarks()。3. flatMap() 添加 Parse 算子.flatMap(recordParse()).name(Parse);MySQL CDC Source 输出的是CdcSourceRecord。它还是一条偏源端格式的 CDC 消息不能直接交给 Paimon Sink。Parse 算子负责把它转换成RichCdcMultiplexRecord这个对象会带上数据库、表名、字段和值等信息供后面的多表 Sink 继续处理。本篇先把 Parse 当作一个转换黑盒下一阶段再单独阅读recordParse()。4. buildSink() 把 Paimon Sink 接到数据流后面EventParser.FactoryRichCdcMultiplexRecordparserFactorybuildEventParserFactory();buildSink(input,parserFactory);buildSink()会根据 DIVIDED 或 COMBINED 模式把相应的 Paimon Sink 算子添加到拓扑中。到这里Flink 作业的主干已经完整MySQL Source ↓ Parse ↓ Paimon Sink但它仍然没有开始跑。七、为什么 DataStream 创建完了数据还没有流动Flink DataStream API 默认采用惰性执行。下面这些调用env.fromSource(...)stream.flatMap(...)stream.process(...)stream.sinkTo(...)主要是在记录算子以及它们之间的关系。可以把它们理解成画施工图Source 节点 ↓ Parse 节点 ↓ Sink 节点此时 Java 代码只是把这张图保存在env中。只有最后执行env.execute()Flink 才会拿着这张图去申请资源、部署 Task 并启动数据处理。这也是为什么调试时可能会遇到一种情况build()已经正常返回但 MySQL 还没有任何读取流量。因为任务还没有提交。八、env.execute() 才是真正的提交按钮build()返回后run()会继续调用execute(syncJobHandler.provideDefaultJobName());execute()定义在 ActionBaseprotectedvoidexecute(StringdefaultName)throwsException{ReadableConfigconfenv.getConfiguration();Stringnameconf.getOptional(PipelineOptions.NAME).orElse(defaultName);env.execute(name);}如果用户没有设置pipeline.nameMySQL 整库同步的默认作业名类似MySQL-Paimon Database Sync: ods其中ods是 Paimon 目标数据库名。执行env.execute(name)后Flink 会完成后续工作读取 env 中的所有 Transformation ↓ 生成 StreamGraph ↓ 转换成可提交的 JobGraph ↓ 通过当前 Flink Executor 提交作业 ↓ JobManager 接收并调度作业 ↓ TaskManager 启动 Source、Parse 和 Sink至于作业提交到哪里由 Flink 当前运行方式决定例如本地模式Standalone Session 集群YarnKubernetesApplication 模式。Paimon Action 不需要分别实现这些部署方式。它只调用标准的env.execute()具体提交流程交给 Flink。九、源码里的三个 build 不要混淆这一段源码里连续出现了多个名称带 build 的方法很容易误以为它们都会启动任务。方法真正作用是否启动作业Action.build()组织整条 Source、Parse、Sink 拓扑否buildSource()创建 Source并由后续代码加入拓扑否SinkBuilder.build()根据模式添加 Sink 相关算子否env.execute()将完整拓扑提交给 Flink是以后看到某个 Builder 的build()先问一句它是在创建对象、添加算子还是最终调用了env.execute()在这条链路中只有env.execute()跨过了“构建”和“运行”的边界。十、一个容易忽略的分布式细节CatalogLoader构建 Sink 时代码没有直接把当前的 Catalog 对象交给算子而是传入.withCatalogLoader(catalogLoader())catalogLoader()保存的是 Catalog 配置并提供“需要时重新创建 Catalog”的能力。为什么不直接把 Catalog 对象传到 TaskManager因为当前 Catalog 可能包含文件系统客户端、连接对象或缓存状态不适合跟随 Flink 作业一起序列化和发送。更稳妥的做法是客户端 └── 把可序列化的 Catalog 配置放进作业 TaskManager └── 根据配置创建自己需要的 Catalog这类 Loader 写法在分布式系统中很常见传递创建对象的方法而不是传递一个已经连好的客户端对象。十一、完整调用链再走一遍FlinkActions.main() ↓ action.run() ↓ SynchronizationActionBase.run() ├── 没有 Checkpoint │ └── 开启三分钟 Checkpoint │ ├── build() │ ├── 检查 MySQL 必填参数 │ ├── 创建 Paimon 目标数据库 │ ├── 发现 MySQL 表 │ ├── 创建或校验 Paimon 表 │ ├── 创建 MySqlSource │ ├── env.fromSource() │ ├── 添加 Parse 算子 │ └── 添加 Paimon Sink 算子 │ └── ActionBase.execute() └── env.execute(jobName) └── 提交到 Flink Runtime └── JobManager 调度 TaskManager十二、读完这一段需要记住什么MySqlSyncDatabaseAction 自己没有实现run()实际执行的是 SynchronizationActionBase 的实现。CDC Action 会在用户没有配置时默认开启三分钟一次的 Checkpoint。build()既会做 MySQL、Paimon 元数据准备也会搭建 Flink 算子拓扑。DataStream API 是惰性执行的调用fromSource()、flatMap()和buildSink()时数据还没有开始流动。env.execute()才是 Flink 作业真正的提交边界。作业提交后Source、Parse 和 Sink 才会作为 Flink Task 在集群中运行。下一篇可以继续进入 MySqlSyncDatabaseActionbeforeBuildingSourceSink()如何发现 MySQL 表为什么没有主键的表会被排除Paimon 目标表不存在时如何创建已有表的 Schema 如何做兼容性检查buildSource()最终让 MySQL CDC 监听哪些表。源码版本Apache Paimon 1.4.2涉及模块paimon-flink-common、paimon-flink-cdc

相关新闻

2026 年 AI Agent 框架横评:10 大框架优缺点对比 + 选型指南

2026 年 AI Agent 框架横评:10 大框架优缺点对比 + 选型指南

2026/8/26 18:26:40

本文由 GO FUNNY 出品。专注 AI Agent 协作方法论与开源工具深度解读。你想搭一个 AI agent,打开 GitHub 搜「agent framework」,按 star 排序,点进第一名,照着 quickstart 敲完,跑通了一个 demo——觉得这事稳了。然后…

信奥梯队选拔面试中如何评估学员的抗挫能力

信奥梯队选拔面试中如何评估学员的抗挫能力

2026/8/26 18:26:40

信奥梯队选拔面试中评估学员抗挫能力,核心是通过可控的“轻度受挫场景”,观察学员面对解题卡壳、思路错误时的真实反应,筛选出能适配信奥长期刷题、反复调试、赛事高压场景的优质苗子。 一、分层设计抗挫能力测试场景 1、‌预备梯队&#x…

★★★ 完全自动化!个人工作日志统计/效率分析工具 —— 时间都去哪了(TimeWhere)

★★★ 完全自动化!个人工作日志统计/效率分析工具 —— 时间都去哪了(TimeWhere)

2026/8/26 18:26:40

你是否也经历过这样的场景? 周五下午被要求交周报,你盯着屏幕发呆——这周到底干了什么?年终写总结,翻遍聊天记录和邮件,却凑不出一份像样的工作总结;每天下班时感觉"挺忙的",回头一…

mermaid.cli自动化实践:如何在文档工程与pre-commit工作流中批量渲染.mmd图表

mermaid.cli自动化实践:如何在文档工程与pre-commit工作流中批量渲染.mmd图表

2026/8/26 19:56:44

mermaid.cli自动化实践:如何在文档工程与pre-commit工作流中批量渲染.mmd图表 【免费下载链接】mermaid.cli Development has been moved to https://github.com/mermaid-js/mermaid-cli 项目地址: https://gitcode.com/gh_mirrors/me/mermaid.cli 一、什么是…

不装ROS也能做机械臂逆运动学?PyKDL IK完整指南(mujoco-learning实战)

不装ROS也能做机械臂逆运动学?PyKDL IK完整指南(mujoco-learning实战)

2026/8/26 19:56:44

不装ROS也能做机械臂逆运动学?PyKDL IK完整指南(mujoco-learning实战) 【免费下载链接】mujoco-learning 项目地址: https://gitcode.com/gh_mirrors/mu/mujoco-learning 还在为配置 ROS 环境头秃吗?其实做**机械臂逆运动…

MiniMax-H3 Turbo-SLA新手避坑清单:10个常见问题一次讲透

MiniMax-H3 Turbo-SLA新手避坑清单:10个常见问题一次讲透

2026/8/26 19:56:44

MiniMax-H3 Turbo-SLA新手避坑清单:10个常见问题一次讲透 【免费下载链接】Minimax-h3-Turbo-SLA 项目地址: https://ai.gitcode.com/hf_mirrors/lightx2v/Minimax-h3-Turbo-SLA MiniMax-H3 Turbo-SLA 是一个基于 4 步蒸馏 的图生视频(FL2V&…

DxWrapper 老游戏兼容快速上手指南:Windows 11 上四步跑通 DirectDraw 老游戏

DxWrapper 老游戏兼容快速上手指南:Windows 11 上四步跑通 DirectDraw 老游戏

2026/8/26 19:56:44

DxWrapper 老游戏兼容快速上手指南:Windows 11 上四步跑通 DirectDraw 老游戏 【免费下载链接】dxwrapper Fixes compatibility issues with older games running on Windows 10/11 by wrapping DirectX dlls. Also allows loading custom libraries with the file …

GPU显存测试完整指南:5分钟跑通 memtest_vulkan

GPU显存测试完整指南:5分钟跑通 memtest_vulkan

2026/8/26 19:56:44

GPU显存测试完整指南:5分钟跑通 memtest_vulkan 【免费下载链接】memtest_vulkan Vulkan compute tool for testing video memory stability 项目地址: https://gitcode.com/gh_mirrors/me/memtest_vulkan 游戏崩溃、训练任务中途中断、出图出现随机噪点——…

3步给停更的Intel老Mac装上macOS Sequoia:OpenCore Legacy Patcher完整指南

3步给停更的Intel老Mac装上macOS Sequoia:OpenCore Legacy Patcher完整指南

2026/8/26 19:46:43

3步给停更的Intel老Mac装上macOS Sequoia:OpenCore Legacy Patcher完整指南 【免费下载链接】OpenCore-Legacy-Patcher Experience macOS just like before 项目地址: https://gitcode.com/GitHub_Trending/op/OpenCore-Legacy-Patcher 一台2015年初的MacBoo…

[光学原理与应用-521]:对光的错误理解与纠偏

[光学原理与应用-521]:对光的错误理解与纠偏

2026/8/26 1:50:39

首先光是一种能量的载体和形态,宏观上观察到的光是由无数个微观的光量子组成的,每个光子在产生的瞬间,其在真空的空间中以确定不变的速度沿着一个初始的方向一直向前,在微观层面,每个光量子的运动轨迹是以波函数所展现…

SIP通话转接原理与REFER方法实战解析

SIP通话转接原理与REFER方法实战解析

2026/8/26 1:49:16

1. 通话转接不是“挂断再拨号”,而是SIP会话的动态重定向你有没有遇到过这样的场景:客服坐席A正在和客户通电话,突然需要把这通对话无缝转给专家坐席B,客户完全感知不到中间的断连——既没听到忙音,也没被要求重新拨号…

Kolla-ansible单节点OpenStack部署实战:从环境准备到排坑指南

Kolla-ansible单节点OpenStack部署实战:从环境准备到排坑指南

2026/8/26 17:50:58

1. 为什么选择Kolla-ansible来部署单节点OpenStack?如果你正在寻找一种能把OpenStack从“概念”快速变成“可用的实验环境”的方法,那么Kolla-ansible几乎是当前最主流、最省心的选择。我见过太多人卡在手动编译依赖、配置服务、处理版本冲突的泥潭里&am…

Python random 模块常用函数详解:从入门到实战

Python random 模块常用函数详解:从入门到实战

2026/8/26 0:05:45

目录 1. 引言2. 准备工作3. 基础随机函数4. 序列相关函数5. 随机种子与复现6. 实战案例7. 注意事项8. 常见问题与排查9. 总结 1. 引言 摘要: 本文系统介绍 Python 标准库 random 模块中最常用的随机数生成函数。内容涵盖基础随机函数(random()、unifor…

Hermes接入团队协作后,我推翻了三个效率假设

Hermes接入团队协作后,我推翻了三个效率假设

2026/8/26 0:05:45

聊《Hermes真能提效吗?先看流程里最慢的那一步》之前,先说一句实在的:别急着背概念,先看它在真实项目里到底解决什么问题。摘要团队把 Hermes 接进项目三个月后,交付速度没有提升反而慢了。复盘后发现,最先…

免费AI大模型调教指南:打造专属网文写作助手

免费AI大模型调教指南:打造专属网文写作助手

2026/8/26 0:05:45

1. 先搞清楚“AI小说扩展模式”到底能帮你做什么如果你是一个刚开始写网文、或者卡在L3级别以下的作者,最头疼的可能是情节推进不下去、人物对话干瘪,或者世界观设定不够丰满。自己对着空白文档硬憋,效率很低。这时候,一个能理解你…

摆脱论文困扰!盘点2026年全网爆红的的AI论文写作工具

摆脱论文困扰!盘点2026年全网爆红的的AI论文写作工具

2026/8/22 2:02:26

一天写完毕业论文在2026年已不再是天方夜谭。2026年最炸裂、实测能大幅提速的AI论文写作工具,覆盖选题构思、文献整理、内容生成、格式排版等核心场景,真正帮你高效搞定论文难题。 一、全流程王者:一站式搞定论文全链路(一天定稿首…

导师推荐!2026最新AI论文工具测评与实用推荐

导师推荐!2026最新AI论文工具测评与实用推荐

2026/8/26 18:07:30

2026年真正好用的AI论文工具,核心看生成的论文质量、低AI味、格式正确、学术适配四大指标。综合实测,千笔AI、ThouPen、豆包、DeepSeek、Grammarly 是当前最值得推荐的梯队,覆盖从免费到付费、从中文到英文、从文科到理工的全场景需求。 一、…

告别游戏崩溃:XCOM 2模组管理器的智能革命

告别游戏崩溃:XCOM 2模组管理器的智能革命

2026/8/26 17:57:52

告别游戏崩溃:XCOM 2模组管理器的智能革命 【免费下载链接】xcom2-launcher The Alternative Mod Launcher (AML) is a replacement for the default game launchers from XCOM 2 and XCOM Chimera Squad. 项目地址: https://gitcode.com/gh_mirrors/xc/xcom2-lau…