Skip to content

S2 第 1 讲:全景图——一张图看懂 paimon-core ​

Apache Paimon 源码学习

作者 X老师(DaemonforY),Paimon 2.0 / master 源码,按 CC BY-NC-SA 4.0 发布。配套实验和代码在 GitHub。

Yui和Kai站在四层数据城中俯瞰版本树与分桶文件流
Yui和Kai站在四层数据城中俯瞰版本树与分桶文件流AI 生成配图

TL;DR:paimon-core 可以看成四层——Catalog、Table、FileStore、存储结构。Paimon = 按 bucket 组织的 LSM 树(数据) + snapshot / manifest 元数据树(版本) + 乐观并发提交(ACID)。读源码从 FileStore.java 这个“总控台”开始:newWrite、newCommit、newScan、newRead 分别通向写入、提交、扫描、读取。本文所有调用链都是用 jdb 在真实运行中抓到的,不是凭记忆画的。

为什么要先看全景 ​

paimon-core 的 src/main/java/org/apache/paimon 下有 971 个 Java 文件。直接从某个类读起,很容易迷失在各种 Abstract*、*Impl、*Provider 里。

S1 的 8 期内容,我们从“表目录里有什么文件”的角度把 Paimon 拆了一遍。S2 换一个角度:这些文件是被哪些类、按什么顺序写出来和读进去的。第一讲先画地图。

一、四层结构 ​

Yui和Kai在四层数据建筑中追踪文件从表到存储的流转
Yui和Kai在四层数据建筑中追踪文件从表到存储的流转AI 生成配图

paimon-core 四层结构

先看各包的规模(括号内为 Java 文件数,d15d250cf):

包文件数作用
table/264对外的 Table 抽象:FileStoreTable、sink、source、系统表
mergetree/115LSM 核心:写缓冲、层级、合并、合并函数、lookup
operation/92FileStore 的操作:scan / write / commit / read / 过期 / 清理
utils/56SnapshotManager、TagManager 等工具
index/49动态桶、删除向量等索引文件
io/44数据文件读写、DataFileMeta
manifest/39元数据文件:ManifestEntry、ManifestFile、ManifestList
append/37append 表的写入与小文件合并
catalog/、jdbc/22、16Catalog 实现

rest/ 在 paimon-core 里只有 3 个文件,REST Catalog 的主体不在这个模块,第 12 讲再展开。

按依赖关系,可以分成四层:

  1. Catalog 层:管理库、表、Schema(FileSystemCatalog、JdbcCatalog、SchemaManager)。
  2. Table 层:给 Flink / Spark 的统一 API。FileStoreTable 有两个实现:建表时主键为空就是 AppendOnlyFileStoreTable,否则是 PrimaryKeyFileStoreTable(FileStoreTableFactory 第 221 行)。写走 TableWriteImpl → TableCommitImpl,读走 DataTableScan → KeyValueTableRead。
  3. FileStore 层:核心引擎。主键表的 PrimaryKeyFileStoreTable.store() 返回 KeyValueFileStore,append 表对应 AppendOnlyFileStore,二者都继承 AbstractFileStore。
  4. 存储结构层:mergetree/(LSM)、io/(数据文件)、manifest/(元数据)、append/、index/。

读源码的入口是 FileStore.java(第 65~118 行),它的方法和下面的包一一对应:

java
SnapshotManager snapshotManager();          // utils/
FileStoreScan newScan();                    // operation/ 扫描 manifest
SplitRead<T> newRead();                     // operation/ 读数据文件
FileStoreWrite<T> newWrite(String commitUser);          // operation/ → mergetree/
FileStoreCommit newCommit(String commitUser, FileStoreTable table);   // operation/
TagManager newTagManager();                 // tag/

主键表里存的不是普通的行,而是 KeyValue:key、sequenceNumber、valueKind(RowKind)、value,外加一个 level(KeyValue.java 第 47~53 行)。sequenceNumber 决定同一主键的新旧,S1 第 4、5 期的实验靠的就是它。

二、一棵元数据树 ​

快照像树根分出清单与数据文件,展示元数据版本关系
快照像树根分出清单与数据文件,展示元数据版本关系AI 生成配图

元数据树

一张表在磁盘上,就是一棵从 snapshot 出发的树:

table/
├── snapshot/snapshot-N            ← Snapshot(JSON,类在 paimon-api)
│     ├─ baseManifestList            之前的全部文件
│     ├─ deltaManifestList           本次提交新增 / 删除的文件
│     ├─ changelogManifestList       本次产生的 changelog(供流读)
│     └─ indexManifest               索引文件(动态桶索引、删除向量)
├── manifest/
│     ├─ manifest-list-xxx         ← ManifestList:一组 ManifestFileMeta
│     └─ manifest-xxx              ← ManifestFile:一组 ManifestEntry(ADD / DELETE + DataFileMeta)
├── schema/schema-N                ← SchemaManager
└── pt=xx/bucket-N/data-xxx.parquet ← 数据文件(DataFileMeta 记录 level、min/max key、seq 范围…)

几个要点:

  • ManifestEntry = 文件类型(ADD / DELETE)+ 分区 + bucket + DataFileMeta。删文件也是“追加一条 DELETE 记录”,和数据层的 LSM 思路一致。
  • DataFileMeta 记着 level、minKey / maxKey、minSequenceNumber / maxSequenceNumber、rowCount、deleteRowCount 等,读的时候靠它裁剪、判断要不要归并。
  • bucket 是读写并行的最小单元:每个 bucket 一棵独立的 LSM 树、一个 writer。

S1 第 6 期用实验证明过:新快照复用旧快照的 manifest 和数据文件,时间旅行不复制数据。这就是这棵树的威力。

三、三条主线:用 jdb 抓真实调用栈 ​

画调用链最容易出错的地方,是凭记忆把类名和顺序写错。所以这一讲的调用链全部来自真实运行:用 JDK 自带的 jdb 连上实验程序,在关键方法下断点,命中时执行 where 打印调用栈。

仓库里新增了一个脚本,一条命令就能复现(需要先按教程 01 章在 d15d250cf 上编译安装 2.2-SNAPSHOT,否则行号对不上):

bash
cd labs && PAIMON_VERSION=2.2-SNAPSHOT ./jdb-stacks.sh sql/lab04/trace-write.sql jdb/s2-1-panorama.txt 30

断点清单 jdb/s2-1-panorama.txt:

stop at org.apache.paimon.mergetree.MergeTreeWriter:166
stop at org.apache.paimon.mergetree.MergeTreeWriter:255
stop in org.apache.paimon.operation.FileStoreCommitImpl.tryCommitOnce
stop in org.apache.paimon.catalog.RenamingSnapshotCommit.commit
stop in org.apache.paimon.operation.AbstractFileStoreScan.plan
stop at org.apache.paimon.operation.RawFileSplitRead:167
stop at org.apache.paimon.operation.MergeFileSplitRead:253

两个小坑:MergeTreeWriter.write 有重载,jdb 要求写参数类型,所以改用行号;断点打在方法声明那一行(没有字节码)会报 “No code at line”,要打在方法体第一行。

实验 SQL 是 S1 第 5 期用过的那条:INSERT INTO trace_orders VALUES (1,'A'), (2,'B'), (1,'C'),并行度 1,然后查 $files 和整表。

主线 ①:写入——只进内存 ​

写入与提交的真实调用栈

MergeTreeWriter.write 命中 3 次(3 行数据),调用栈(省略 Flink 框架帧和同名重载):

RowDataStoreWriteOperator.processElement   :55    Flink 算子 → Paimon
StoreSinkWriteImpl.write                   :143
TableWriteImpl.writeAndReturn              :247   Table 层
AbstractFileStoreWrite.write               :206   按 (分区, bucket) 找到 writer
MergeTreeWriter.write                      :166   分配 sequenceNumber,放进写缓冲

这一步不产生任何文件,S1 第 5 期已经用实验看到过。

主线 ②:提交——prepareCommit 刷盘,snapshot 落地 ​

批作业输入结束时,先在写端准备提交:

PrepareCommitOperator.endInput             :103   批作业输入结束(流作业是 checkpoint 前)
TableWriteOperator.prepareCommit           :143
TableWriteImpl.prepareCommit               :313
MemoryFileStoreWrite.prepareCommit         :155   KeyValueFileStoreWrite 的父类
AbstractFileStoreWrite.prepareCommit       :266
MergeTreeWriter.prepareCommit              :255   flushWriteBuffer → 写出 L0 文件

然后在提交端提交快照:

CommitterOperator.endInput                 :192
StoreCommitter.filterAndCommit             :123
TableCommitImpl.filterAndCommitMultiple    :330   幂等过滤
TableCommitImpl.commitMultiple             :281
FileStoreCommitImpl.commit                 :370
FileStoreCommitImpl.tryCommit              :885   重试循环
FileStoreCommitImpl.tryCommitOnce          :1314  冲突检测、写 manifest / manifest list
FileStoreCommitImpl.commitSnapshotImpl     :1737
RenamingSnapshotCommit.commit              :60    写 snapshot-N

RenamingSnapshotCommit 第 66 行调用 fileIO.tryToWriteAtomic:先写一个临时文件,再 rename 成 snapshot-N(FileIO.java 第 382~395 行)。rename 成功的那一刻,这次写入才对读者可见。rename 在不同文件系统上的原子性、并发提交时谁赢,是第 5 讲的内容。

一个细节:上面三个断点停在同一个线程里,线程名是 Writer : trace_orders -> Global Committer : trace_orders -> end: Writer (1/1)#0。并行度为 1 时,Flink 把 Writer 和 Committer 两个算子 chain 在了一起。

主线 ③:读取——规划在 JobManager,读文件在 TaskManager ​

读取的真实调用栈

读取分两步。规划(读 manifest、生成 split)发生在 JobManager 的 SourceCoordinator 里,线程是 flink-pekko.actor.default-dispatcher-*:

SourceCoordinator.start                    (Flink)
StaticFileStoreSource.getSplits            :176
AbstractDataTableScan.plan                 :156   Table 层:TableScan
SnapshotReaderImpl.read                    :394
AbstractFileStoreScan.plan                 :301   读 manifest、按分区 / bucket / 统计信息裁剪

读取发生在 TaskManager 的 Source Data Fetcher 线程:

FileStoreSourceSplitReader.fetch           :124
KeyValueTableRead.createReader             :167   Table 层:TableRead
KeyValueTableRead.reader                   :230   依次问 SplitReadProvider:这个 split 归谁读?
RawFileSplitRead.createReader              :167   ← trace_orders 实测命中这里

这里有个出乎意料的结果:trace_orders 是主键表,文件还在 level 0,但读的时候没有走合并读 MergeFileSplitRead,而是直接读 RawFileSplitRead——MergeFileSplitRead 这个类在整个作业里都没有被加载。

原因在 MergeTreeSplitGenerator.splitForBatch(第 108~114 行):文件先按 key 范围做区间划分,一个区间里只有一个文件、且没有删除记录时,这个 split 被标记为 rawConvertible,PrimaryKeyTableRawFileSplitReadProvider.match 就会接手,直接读文件。只有一个文件,自然不用归并。

换成 S1 第 5 期的 wb_two 表(分两次写入,两个文件的 key 重叠),再抓一次(jdb/s2-1-merge-read.txt),这次命中了合并读:

KeyValueTableRead.reader                   :230
SplitRead$1.createReader                   :103
MergeFileSplitReadProvider (lambda)        :53
MergeFileSplitRead.createReader            :253   区间划分 → 败者树归并 → 合并函数 → DropDeleteReader

同一次运行里还有一个对照:读 audit_log 系统表时,即使只有一个文件也走了 MergeFileSplitRead——因为 audit_log 要保留 -D 记录(forceKeepDelete),直接读的 provider 不接这种 split。

合并读内部的几个组件:IntervalPartition 做区间划分(第 338 行)、归并默认用败者树(sort-engine 默认 LOSER_TREE)、最后 DropDeleteReader 过滤掉删除记录(第 487 行)。什么时候能跳过归并、删除向量怎么帮忙,是第 6、7 讲的内容。

四、读源码路线 ​

读源码路线

沿着一条数据的生命周期走:

顺序从哪个类读看什么本系列
1FileStore.java总控台,方法和包一一对应本讲
2MergeTreeWriter写缓冲、刷盘、触发合并第 2 讲
3UniversalCompaction、MergeTreeCompactTask选哪些文件合并、怎么合并第 3、4 讲
4FileStoreCommitImpl冲突检测、重试、原子提交第 5 讲
5MergeFileSplitRead、删除向量读路径、什么时候不用归并第 6、7 讲
6Lookup changelog、Flink Sink / Sourcechangelog、checkpoint、流读第 8~10 讲

建议的读法:先跑一遍 jdb-stacks.sh 拿到调用栈,再对着栈一层层点进源码。调用栈告诉你“谁调用了谁”,源码告诉你“为什么这样调用”,两者对着看,比从入口硬读快得多。

小结 ​

  • 四层:Catalog / Table / FileStore / 存储结构;主键表 PrimaryKeyFileStoreTable + KeyValueFileStore,append 表 AppendOnlyFileStoreTable + AppendOnlyFileStore。
  • 一棵树:snapshot → manifest list → manifest → 数据文件;bucket 是并行单元。
  • 三条主线:写进内存 → prepareCommit 刷成 L0 文件 → 提交写下 snapshot 才可见;读时先规划 split,再按 split 选择直接读或合并读。

下一讲:一条数据的写入之旅——MergeTreeWriter 里的写缓冲、刷盘和合并触发。

自测题:trace_orders 如果连续写两次(不 DROP 表),第二次查询整表时会命中 RawFileSplitRead 还是 MergeFileSplitRead?(答案:取决于两个文件的 key 范围是否重叠。实验里两次都写订单 1,key 重叠,会被划进同一个区间,走 MergeFileSplitRead;如果第二次写的是全新的订单号且 key 范围不重叠,两个区间各一个文件,仍然可以直接读。未单独运行,可以用 jdb-stacks.sh 自己验证。)


代码示例在页面里运行时使用 HiveGPT 的模型接口。延伸阅读来自 JavaGuide(Apache-2.0),版权归原作者。