Skip to content

00 学习计划:从入门到 Paimon 专家 ​

Apache Paimon 源码学习

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

Yui和Kai沿Paimon路线图进阶学习
Yui和Kai沿Paimon路线图进阶学习AI 生成配图

总体安排 ​

  • 节奏建议:每周 8~12 小时,约 14 周完成主线(阶段 0~6),之后进入长期的专家阶段(阶段 7)。
  • 学习方法(每个主题都按这个循环):
    1. 用:先用 SQL 把功能跑起来,观察现象(系统表、目录文件)。
    2. 读:按本教程的“源码导读”顺序读代码,画调用链。
    3. 调:跑对应单测,在关键方法下断点,看变量。
    4. 改:做一个小改动(加日志 / 改参数 / 写测试),验证自己的理解。
    5. 讲:用自己的话写一份笔记或给别人讲一遍,讲不清楚就回到第 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 + Changelog06、07能解释“为什么这个 split 要 merge”,能推演 -U/+U
4第 9~10 周Flink 集成08、09能分析 checkpoint 慢、快照过期、重复/丢数问题
5第 11~12 周Spark 集成10能判断一条 DML 走哪条路径、代价多大
6第 13 周Catalog / REST11能实现一个最小 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 周):核心概念 + 写路径 ​

从快照到数据文件的写入链路
从快照到数据文件的写入链路AI 生成配图

任务 ​

  • [ ] 精读 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 合并 + 提交 ​

LSM分层合并与提交流程
LSM分层合并与提交流程AI 生成配图

第 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 三种策略分别在什么条件下使用。

第 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 Bucket bucket=-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

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