附录:配置速查、源码索引与调试技巧
Apache Paimon 源码学习
作者 X老师(DaemonforY),Paimon 2.0 / master 源码,按 CC BY-NC-SA 4.0 发布。配套实验和代码在 GitHub。

A. 关键配置速查

写入与合并
| 配置 | 默认 | 说明 | 章节 |
|---|---|---|---|
bucket | 主键表需指定;-1 动态桶;-2 postpone | 分桶数 | 02 |
write-buffer-size | — | 写缓冲大小,影响 L0 文件大小 | 02 |
target-file-size | 主键 128MB / append 256MB | 输出文件滚动大小 | 04 |
num-sorted-run.compaction-trigger | 5 | 触发合并的 run 数 | 03 |
num-sorted-run.stop-trigger | trigger + 3 | 写入反压的 run 数 | 03 |
num-levels | trigger + 1 | LSM 层数 | 03 |
compaction.size-ratio | 1 | 大小比例容忍度(%) | 03 |
compaction.max-size-amplification-percent | 200 | 空间放大阈值 | 03 |
compaction.small-file-ratio | 0.7 | 小文件阈值比例 | 04 |
compaction.optimization-interval | — | 定期全量合并 | 03 |
compaction.offpeak.start.hour / end.hour / offpeak-ratio | -1 / -1 / 0 | 低峰期更激进合并 | 03 |
compaction.force-up-level-0 | false | 强制 L0 上推 | 03 |
sort-spill-threshold | stop-trigger + 1 | 归并前溢写阈值 | 04 |
sort-engine | loser-tree | 多路归并算法 | 04 |
write-only | false | 只写不合并(配合独立 compaction 作业) | 05 |
Changelog 与 Lookup
| 配置 | 默认 | 说明 | 章节 |
|---|---|---|---|
changelog-producer | none | none / input / lookup / full-compaction | 07 |
lookup-wait | true | prepareCommit 等待 L0 合并 | 07 |
lookup-compact | RADICAL | RADICAL / GENTLE | 07 |
lookup-compact.max-interval | — | GENTLE 模式强制间隔 | 07 |
changelog-producer.row-deduplicate | false | 值不变不产 -U/+U | 07 |
changelog-producer.ignore-update-before / ignore-delete | false | 过滤 changelog | 07 |
删除向量与读取
| 配置 | 默认 | 说明 | 章节 |
|---|---|---|---|
deletion-vectors.enabled | false | 删除向量模式 | 06 |
deletion-vectors.merge-on-read | false | DV 模式下让 L0 可见(读变慢) | 06 |
scan.mode | default → latest-full | 流读起点 | 09 |
continuous.discovery-interval | 10s | 新快照发现间隔 | 09 |
scan.max-splits-per-task | 10 | enumerator 反压 | 09 |
scan.remove-normalize | false | 跳过 ChangelogNormalize | 09 |
提交与快照
| 配置 | 默认 | 说明 | 章节 |
|---|---|---|---|
commit.max-retries | 10 | 提交重试次数 | 05 |
commit.timeout | 无 | 提交重试总时长 | 05 |
commit.min-retry-wait / max-retry-wait | 10ms / 10s | 退避 | 05 |
commit.discard-duplicate-files | false | 幂等重复提交 | 05 |
snapshot.time-retained 等 | — | 快照保留 | 09 |
consumer-id | — | 消费者位点,防快照过期 | 09 |
consumer.mode | exactly-once | 消费位点一致性 | 09 |
consumer.expiration-time | — | 死 consumer 清理 | 09 |
Flink / REST
| 配置 | 说明 | 章节 |
|---|---|---|
sink.operator-uid.suffix | 稳定算子 UID,保证状态可恢复 | 08 |
sink.coordinator-commit.enabled | JobManager 内提交(仅 unaware-bucket append 表) | 08 |
metastore=rest、uri、warehouse、token.provider、token | REST Catalog | 11 |
B. 源码索引(按主题)
| 主题 | 入口类 |
|---|---|
| 总控 | FileStore、KeyValueFileStore、AppendOnlyFileStore |
| 元数据 | Snapshot(paimon-api)、manifest/ManifestEntry、ManifestFile、ManifestList、io/DataFileMeta |
| 写入 | table/sink/TableWriteImpl、operation/KeyValueFileStoreWrite、mergetree/MergeTreeWriter |
| LSM 层级 | mergetree/Levels、SortedRun |
| 合并选择 | mergetree/compact/UniversalCompaction、ForceUpLevel0Compaction、EarlyFullCompaction |
| 合并执行 | MergeTreeCompactManager、MergeTreeCompactTask、IntervalPartition、MergeTreeCompactRewriter |
| 归并读 | mergetree/MergeTreeReaders、MergeSorter、SortMergeReaderWithLoserTree |
| 合并引擎 | DeduplicateMergeFunction、PartialUpdateMergeFunction、aggregate/*、FirstRowMergeFunction |
| Lookup/Changelog | LookupChangelogMergeFunctionWrapper、LookupMergeFunction、LookupLevels、ChangelogMergeTreeRewriter |
| 提交 | operation/FileStoreCommitImpl、operation/commit/ConflictDetection、catalog/RenamingSnapshotCommit |
| 扫描 | table/source/DataTableScan、snapshot/SnapshotReaderImpl、operation/KeyValueFileStoreScan、MergeTreeSplitGenerator |
| 读取 | table/source/KeyValueTableRead、operation/MergeFileSplitRead、RawFileSplitRead |
| 删除向量 | deletionvectors/DeletionVector、BucketedDvMaintainer、ApplyDeletionVectorReader |
| 流读 | table/source/DataTableStreamScan、*FollowUpScanner、utils/NextSnapshotFetcher |
| 过期清理 | table/ExpireSnapshotsImpl、operation/SnapshotDeletion、OrphanFilesClean |
| Flink | flink/sink/FlinkSink、CommitterOperator、flink/source/ContinuousFileSplitEnumerator |
| Spark | spark/catalyst/analysis/PaimonMergeInto、commands/MergeIntoPaimonTable |
| REST | rest/RESTCatalog、rest/RESTApi、rest/RESTTokenFileIO、table/CatalogEnvironment |
C. 推荐测试(跟断点用)
| 主题 | 测试 |
|---|---|
| 端到端写读 | paimon-core/.../table/PrimaryKeySimpleTableTest |
| LSM 层级 | mergetree/LevelsTest、LookupLevelsTest、MergeSorterTest |
| 合并选择 | mergetree/compact/UniversalCompactionTest、ForceUpLevel0CompactionTest |
| 合并执行 | mergetree/compact/IntervalPartitionTest、ChangelogMergeTreeRewriterTest |
| Lookup changelog | mergetree/compact/LookupChangelogMergeFunctionWrapperTest |
| 提交 | operation/FileStoreCommitTest |
| 删除向量 | deletionvectors/DeletionVectorTest、BucketedDvMaintainerTest、flink DeletionVectorITCase |
| Flink 提交 | paimon-flink-common/.../sink/CommitterOperatorTest |
| Flink 流读 | .../source/ContinuousFileSplitEnumeratorTest、ContinuousFileStoreITCase、LookupChangelogWithAggITCase |
| Spark MERGE | paimon-spark-ut/.../sql/MergeIntoTableTestBase.scala |
| REST | paimon-core/.../rest/RESTCatalogTest(服务端 RESTCatalogServer) |
D. 调试技巧

- 单测优先:core 的测试都用本地文件系统,无需 Flink/Spark 集群,断点最方便。
- 看目录:测试里打印表路径(通常在
@TempDir下),直接cat snapshot/snapshot-N看 JSON。 - 系统表:
t$snapshots、t$files、t$manifests、t$schemas、t$consumers、t$table_indexes、t$partitions、t$tags、t$branches、t$options。 - 日志开关:
org.apache.paimon.mergetree.compactDEBUG → 合并选择原因;org.apache.paimon.operation.FileStoreCommitImplDEBUG → 提交的文件清单;org.apache.paimon.table.sourceDEBUG → 流读快照推进。
- 关键日志关键字:
| 关键字 | 含义 |
|---|---|
Universal compaction due to ... | 合并触发原因 |
Paimon compact task finished ... inputBytes / outputBytes | 合并结果 |
Atomic commit failed for snapshot | 抢快照号失败(会重试) |
File deletion conflicts detected | 删除冲突(多作业合并同一 bucket / 旧 savepoint) |
LSM conflicts detected | 层级重叠冲突 |
The wanted read snapshot with id X has expired | 流读快照过期 |
This exception is intentionally thrown after committing the restored checkpoints | 恢复补提交后的故意失败(正常) |
begin/end refresh data token | REST 数据 token 刷新 |
E. 常见问题 → 原因 → 解决
| 问题 | 原因 | 解决 |
|---|---|---|
| 小文件过多 | checkpoint 间隔短、bucket 多、合并跟不上 | 加大 checkpoint 间隔;调 num-sorted-run.*;独立 compaction 作业 |
| checkpoint 超时 | lookup/DV 表等 L0 合并;写入反压 | 调 lookup-compact=GENTLE;增加并行度;write-only + 独立合并 |
| 提交冲突持续 failover | 多作业合并同一 bucket;旧 savepoint 恢复 | 单写 + write-only;从最新 savepoint 恢复 / 回滚表 |
| 流读快照过期 | 保留时间短、消费慢、作业停太久 | consumer-id;加大保留时间 |
| 批读慢 | merge-on-read、value 过滤无法下推 | 开 DV;全量合并;file index |
| 下游 Flink 状态大 | changelog-producer=none 触发 ChangelogNormalize | 改用 lookup / input |
| JDK 17 编译失败找不到类 | activeByDefault profile 失效 | 加 -Pspark3,flink1 |