Skip to content

附录:配置速查、源码索引与调试技巧 ​

Apache Paimon 源码学习

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

配置、源码与调试工具汇成排障地图
配置、源码与调试工具汇成排障地图AI 生成配图

A. 关键配置速查 ​

旋钮控制文件大小与分桶合并策略
旋钮控制文件大小与分桶合并策略AI 生成配图

写入与合并 ​

配置默认说明章节
bucket主键表需指定;-1 动态桶;-2 postpone分桶数02
write-buffer-size—写缓冲大小,影响 L0 文件大小02
target-file-size主键 128MB / append 256MB输出文件滚动大小04
num-sorted-run.compaction-trigger5触发合并的 run 数03
num-sorted-run.stop-triggertrigger + 3写入反压的 run 数03
num-levelstrigger + 1LSM 层数03
compaction.size-ratio1大小比例容忍度(%)03
compaction.max-size-amplification-percent200空间放大阈值03
compaction.small-file-ratio0.7小文件阈值比例04
compaction.optimization-interval—定期全量合并03
compaction.offpeak.start.hour / end.hour / offpeak-ratio-1 / -1 / 0低峰期更激进合并03
compaction.force-up-level-0false强制 L0 上推03
sort-spill-thresholdstop-trigger + 1归并前溢写阈值04
sort-engineloser-tree多路归并算法04
write-onlyfalse只写不合并(配合独立 compaction 作业)05

Changelog 与 Lookup ​

配置默认说明章节
changelog-producernonenone / input / lookup / full-compaction07
lookup-waittrueprepareCommit 等待 L0 合并07
lookup-compactRADICALRADICAL / GENTLE07
lookup-compact.max-interval—GENTLE 模式强制间隔07
changelog-producer.row-deduplicatefalse值不变不产 -U/+U07
changelog-producer.ignore-update-before / ignore-deletefalse过滤 changelog07

删除向量与读取 ​

配置默认说明章节
deletion-vectors.enabledfalse删除向量模式06
deletion-vectors.merge-on-readfalseDV 模式下让 L0 可见(读变慢)06
scan.modedefault → latest-full流读起点09
continuous.discovery-interval10s新快照发现间隔09
scan.max-splits-per-task10enumerator 反压09
scan.remove-normalizefalse跳过 ChangelogNormalize09

提交与快照 ​

配置默认说明章节
commit.max-retries10提交重试次数05
commit.timeout无提交重试总时长05
commit.min-retry-wait / max-retry-wait10ms / 10s退避05
commit.discard-duplicate-filesfalse幂等重复提交05
snapshot.time-retained 等—快照保留09
consumer-id—消费者位点,防快照过期09
consumer.modeexactly-once消费位点一致性09
consumer.expiration-time—死 consumer 清理09
配置说明章节
sink.operator-uid.suffix稳定算子 UID,保证状态可恢复08
sink.coordinator-commit.enabledJobManager 内提交(仅 unaware-bucket append 表)08
metastore=rest、uri、warehouse、token.provider、tokenREST Catalog11

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/ChangelogLookupChangelogMergeFunctionWrapper、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
Flinkflink/sink/FlinkSink、CommitterOperator、flink/source/ContinuousFileSplitEnumerator
Sparkspark/catalyst/analysis/PaimonMergeInto、commands/MergeIntoPaimonTable
RESTrest/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 changelogmergetree/compact/LookupChangelogMergeFunctionWrapperTest
提交operation/FileStoreCommitTest
删除向量deletionvectors/DeletionVectorTest、BucketedDvMaintainerTest、flink DeletionVectorITCase
Flink 提交paimon-flink-common/.../sink/CommitterOperatorTest
Flink 流读.../source/ContinuousFileSplitEnumeratorTest、ContinuousFileStoreITCase、LookupChangelogWithAggITCase
Spark MERGEpaimon-spark-ut/.../sql/MergeIntoTableTestBase.scala
RESTpaimon-core/.../rest/RESTCatalogTest(服务端 RESTCatalogServer)

D. 调试技巧 ​

本地测试台用放大镜定位故障线索
本地测试台用放大镜定位故障线索AI 生成配图
  1. 单测优先:core 的测试都用本地文件系统,无需 Flink/Spark 集群,断点最方便。
  2. 看目录:测试里打印表路径(通常在 @TempDir 下),直接 cat snapshot/snapshot-N 看 JSON。
  3. 系统表:t$snapshots、t$files、t$manifests、t$schemas、t$consumers、t$table_indexes、t$partitions、t$tags、t$branches、t$options。
  4. 日志开关:
    • org.apache.paimon.mergetree.compact DEBUG → 合并选择原因;
    • org.apache.paimon.operation.FileStoreCommitImpl DEBUG → 提交的文件清单;
    • org.apache.paimon.table.source DEBUG → 流读快照推进。
  5. 关键日志关键字:
关键字含义
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 tokenREST 数据 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

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