00 学习计划:从入门到 Paimon 专家
Apache Paimon 源码学习
作者 X老师(DaemonforY),Paimon 2.0 / master 源码,按 CC BY-NC-SA 4.0 发布。配套实验和代码在 GitHub。

总体安排
- 节奏建议:每周 8~12 小时,约 14 周完成主线(阶段 0~6),之后进入长期的专家阶段(阶段 7)。
- 学习方法(每个主题都按这个循环):
- 用:先用 SQL 把功能跑起来,观察现象(系统表、目录文件)。
- 读:按本教程的“源码导读”顺序读代码,画调用链。
- 调:跑对应单测,在关键方法下断点,看变量。
- 改:做一个小改动(加日志 / 改参数 / 写测试),验证自己的理解。
- 讲:用自己的话写一份笔记或给别人讲一遍,讲不清楚就回到第 2 步。
- 笔记建议:在
paimon-learning/notes/下按章节记笔记,每章至少包含:一张调用链图、3 个“原来如此”的发现、1 个没搞懂的问题。
阶段总览
| 阶段 | 周数 | 主题 | 对应章节 | 里程碑 |
|---|---|---|---|---|
| 0 | 第 1 周 | 基础准备 + 当用户用起来 | 01 | 本地跑通 Flink SQL 读写 Paimon,看懂表目录 |
| 1 | 第 2~3 周 | 核心概念 + 写路径 | 02 | 能画出 snapshot→manifest→data 树,能讲清一次写入 |
| 2 | 第 4~6 周 | LSM 合并 + 提交 | 03、04、05 | 能手算 compaction 选择,能解释任一冲突异常 |
| 3 | 第 7~8 周 | 读路径 + DV + Changelog | 06、07 | 能解释“为什么这个 split 要 merge”,能推演 -U/+U |
| 4 | 第 9~10 周 | Flink 集成 | 08、09 | 能分析 checkpoint 慢、快照过期、重复/丢数问题 |
| 5 | 第 11~12 周 | Spark 集成 | 10 | 能判断一条 DML 走哪条路径、代价多大 |
| 6 | 第 13 周 | Catalog / REST | 11 | 能实现一个最小 REST Server 的 commit 接口 |
| 7 | 第 14 周起 | 专家进阶(长期) | 附录 + 社区 | 给 apache/paimon 提交被合入的 PR |
阶段 0(第 1 周):基础准备
前置知识自查(不熟的先补)
- [ ] LSM Tree:memtable、SSTable、level、compaction、读放大/写放大/空间放大。推荐读 RocksDB Wiki 的 Universal Compaction 一篇。
- [ ] Flink:Checkpoint 屏障对齐、两阶段提交(TwoPhaseCommitSink)、FLIP-27 Source(Enumerator/Reader/Split)。
- [ ] Spark:DataSource V2、Catalyst 分析规则、
SupportsRowLevelOperations。 - [ ] 列存格式:Parquet/ORC 的 row group、统计信息、谓词下推。
- [ ] Java 并发:线程池、Future、synchronized/双重检查。
任务
- [ ] 按 01 章 完成编译与 IDEA 导入(记得勾选 spark3、flink1 profile)。
- [ ] 读官方文档
docs/docs/concepts/下:Overview、Basic Concepts、Primary Key Table、Append Table。 - [ ] 本地起一个 Flink SQL Client(或写 Java main),完成:
- 建一张主键表、一张 append 表,插入数据、更新、删除。
- 查询系统表:
t$snapshots、t$files、t$manifests、t$options、t$schemas。 - 打开表目录,对照看
snapshot/、manifest/、schema/、bucket-N/。
验收标准
- 能说出主键表和 append 表的 3 个区别。
- 打开任意一个表目录,能说出每个子目录/文件的作用。
阶段 1(第 2~3 周):核心概念 + 写路径

任务
- [ ] 精读 02 章。
- [ ] 依次阅读核心数据结构并做笔记:
FileStore→KeyValue→io/DataFileMeta→manifest/ManifestEntry→manifest/ManifestFileMeta→Snapshot(在 paimon-api)。 - [ ] 跟读写路径:
table/sink/TableWriteImpl→operation/KeyValueFileStoreWrite→mergetree/MergeTreeWriter(write/flushWriteBuffer/prepareCommit)。 - [ ] 跑
PrimaryKeySimpleTableTest中任意一个写+读的用例,在MergeTreeWriter.write、flushWriteBuffer下断点。 - [ ] 实验:把
write-buffer-size调很小,观察 L0 文件数变化(t$files里的level列)。
验收标准
- 手画 snapshot → manifest list → manifest → data file 的树,并标出 base/delta/changelog 三类 manifest list 的区别。
- 讲清楚:一条记录从
TableWrite.write到落成 L0 文件经历了哪些类。 - 解释 sequenceNumber 的作用。
阶段 2(第 4~6 周):LSM 合并 + 提交

第 4 周:选文件
- [ ] 精读 03 章。
- [ ] 跑
UniversalCompactionTest,自己手写 3 个用例(对应章节里的例 A/B/C),先手算结果再跑测试验证。 - [ ] 读
ForceUpLevel0Compaction、EarlyFullCompaction、OffPeakHours。
第 5 周:执行合并
- [ ] 精读 04 章。
- [ ] 跑
IntervalPartitionTest,理解 section / sorted run 的切分。 - [ ] 在
MergeTreeCompactTask.doCompact下断点,观察哪些文件被 upgrade、哪些被 rewrite。 - [ ] 读
SortMergeReaderWithLoserTree,手画一次败者树的调整过程。
第 6 周:提交
- [ ] 精读 05 章。
- [ ] 跑
FileStoreCommitTest里带 conflict 的用例。 - [ ] 实验:写一个测试,两个不同 commitUser 对同一 bucket 做 compact 提交,复现 “File deletion conflicts detected”。
验收标准
- 给出任意一组 runs(level + size),能手算 UniversalCompaction 的选择结果和 outputLevel。
- 能解释为什么大文件“升级不重写”,以及哪 3 种情况会强制重写。
- 看到任何一条提交冲突异常,能说出是哪种冲突、谁和谁冲突、怎么解决。
阶段 3(第 7~8 周):读路径 + 删除向量 + Changelog
第 7 周:读路径与 DV
- [ ] 精读 06 章。
- [ ] 实验:建两张主键表,一张开
deletion-vectors.enabled,写相同数据,对比:t$files的 level 分布、t$table_indexes里的 DV 文件;- 在
KeyValueTableRead.createReader下断点,看各自走哪个 provider。
- [ ] 读
DeletionVector、BitmapDeletionVector、BucketedDvMaintainer,跑DeletionVectorTest。
第 8 周:Lookup Changelog
- [ ] 精读 07 章。
- [ ] 跑
LookupChangelogMergeFunctionWrapperTest、LookupLevelsTest。 - [ ] 实验:
changelog-producer=lookup的表,流读下游打印 RowKind,分别构造 +I、-U/+U、-D 场景,并验证changelog-producer.row-deduplicate的效果。
验收标准
- 能解释:DV 模式下读为什么不需要归并、为什么 L0 不可见、代价是什么。
- 给定 L0 和高层的数据,能推演出 lookup 产生的 changelog 和写出的新文件内容。
- 能讲清楚
upgradeStrategy三种策略分别在什么条件下使用。
阶段 4(第 9~10 周):Flink 集成
第 9 周:写入与提交
- [ ] 精读 08 章。
- [ ] 跑
CommitterOperatorTest,重点看“checkpoint 完成后重启”的用例。 - [ ] 实验:Flink 作业写 Paimon,checkpoint 完成后、提交前 kill TaskManager,观察重启后的日志(
RestoreAndFailCommittableStateManager故意失败一次)。
第 10 周:流式读取
- [ ] 精读 09 章。
- [ ] 跑
ContinuousFileSplitEnumeratorTest(含 consumer 进度相关用例)、ContinuousFileStoreITCase。 - [ ] 实验:
- 设置很短的
snapshot.time-retained,停掉流读作业一段时间再恢复,复现 “snapshot has expired”; - 加上
consumer-id再试一次,观察t$consumers和快照是否被保护。
- 设置很短的
验收标准
- 画出 checkpoint N 从 writer 到 committer 的完整时序。
- 解释 exactly-once 在 4 种故障时刻下是怎么保证的。
- 解释为什么 source 并行度超过 bucket 数没有意义、为什么 changelog-producer=none 会引入 ChangelogNormalize。
阶段 5(第 11~12 周):Spark 集成
- [ ] 精读 10 章。
- [ ] 本地 Spark 3.5 + Paimon,分别对主键表、append 表、append+DV 表执行同一条 MERGE INTO,用
EXPLAIN EXTENDED和 Spark UI 对比执行计划与 job 数。 - [ ] 跑
MergeIntoTableTestBase中的用例。 - [ ] 延伸阅读:
DeleteFromPaimonTableCommand、UpdatePaimonTableCommand、OptimizeMetadataOnlyDeleteFromPaimonTable。 - [ ] 读 Spark 读路径:
PaimonScan、PaimonBaseScanBuilder的谓词/聚合/limit 下推。
验收标准
- 任给一张表和一条 DML,能说出走哪条路径、产生什么类型的 commit、是否可能与流写冲突。
- 解释 CoW 为什么需要
filesToRead这类“只读不重写”的文件。
阶段 6(第 13 周):Catalog 与 REST
- [ ] 精读 11 章。
- [ ] 读
docs/static/rest-catalog-open-api.yaml,对照RESTApi。 - [ ] 跑
RESTCatalogTest,读RESTCatalogServer.commitSnapshot。 - [ ] 对比阅读
FileSystemCatalog、HiveCatalog(paimon-hive)的建表、提交、加锁方式。 - [ ] 动手项目:用 Spring Boot / Javalin 实现一个最小 REST Server,支持
config、getTable、createTable、commit(基于 base snapshot uuid 的 CAS),用 Paimon Java 客户端连上去写入数据。
验收标准
- 能讲清 REST 模式下 commit 的 CAS 协议,以及它和客户端幂等检查如何配合。
- 能讲清数据 token 下发的刷新与缓存机制。
阶段 7(第 14 周起,长期):专家进阶
主线之外,成为专家还需要覆盖以下主题(建议每个主题 1~2 周,按兴趣排序):
7.1 未覆盖的核心模块
- [ ] Schema Evolution:
schema/SchemaManager、字段 id 映射、读旧文件时的类型转换(CastExecutors)。 - [ ] 分桶模式:动态桶
bucket=-1(index/、crosspartition/)、Postpone Bucketbucket=-2(postpone/)、append 表 unaware bucket。 - [ ] 合并引擎:
PartialUpdateMergeFunction(sequence-group)、aggregate/下各种聚合函数、FirstRowMergeFunction。 - [ ] 快照生命周期:
ExpireSnapshotsImpl、SnapshotDeletion、OrphanFilesClean、Tag/Branch(tag/、utils/BranchManager)。 - [ ] Manifest 管理:
ManifestFileMerger、manifest 统计与过滤、ManifestEntryCache。 - [ ] 文件格式与索引:
paimon-format(Parquet/ORC/Avro 读写器)、file index(bloom filter、bitmap、BSI)、global index、向量/全文检索。 - [ ] Iceberg 兼容:
iceberg/生成 Iceberg 元数据。 - [ ] Data Evolution / Row Tracking:
MergeIntoPaimonDataEvolutionTable、_ROW_ID。 - [ ] CDC 入湖:
paimon-flink-cdc(MySQL/Kafka 整库同步、schema 变更)。 - [ ] Python / 多语言:
paimon-python/pypaimon的读写与 REST 客户端。
7.2 生产能力
- [ ] 性能调优:写入(buffer、并行度、
write-only+ 独立 compaction 作业)、合并(num-sorted-run.*、lookup-wait)、读取(DV、file index、split 大小)。 - [ ] 问题排查手册:整理你自己的 “异常 → 原因 → 解决” 表(参考附录)。
- [ ] 监控指标:
metrics/下 commit、compaction、scan 指标,接入 Prometheus。 - [ ] 压测:用
paimon-benchmark做写入/合并/查询基准,对比不同配置。
7.3 社区贡献(专家的标志)
- [ ] 订阅 dev 邮件列表,关注 GitHub issues/PR。
- [ ] 从
good first issue开始:修文档 → 补测试 → 修 bug → 做小特性。 - [ ] 每周 review 1~2 个 PR,学习 committer 的设计思路。
- [ ] 读 PIP(Paimon Improvement Proposal)设计文档,理解大特性的取舍。
- [ ] 目标:3 个月内有 PR 被合入;6~12 个月能独立设计并实现一个中等特性。
专家自检清单
- [ ] 能在白板上完整讲出一次写入→合并→提交→读取的全链路,并回答任意细节追问。
- [ ] 线上遇到小文件过多、checkpoint 超时、流读延迟、提交冲突、快照过期、读性能差,能快速定位并给出配置/架构方案。
- [ ] 能评估一个新需求(如新合并引擎、新索引)应该改哪些模块、影响哪些兼容性。
- [ ] 有被 apache/paimon 合入的代码。
每周例行模板
周一~周二:读本周章节 + 源码导读(2~4h)
周三~周四:跑测试、下断点、做实验(3~4h)
周五:写笔记 + 自测题(1~2h)
周末:复盘本周“没搞懂的问题”,补读相关代码;更新进度勾选(1~2h)进度记录
| 阶段 | 开始日期 | 完成日期 | 笔记链接 | 备注 |
|---|---|---|---|---|
| 0 | ||||
| 1 | ||||
| 2 | ||||
| 3 | ||||
| 4 | ||||
| 5 | ||||
| 6 | ||||
| 7 |