Skip to content
P · Apache Paimon 源码第 6 课⏱ 15 分钟

06 读路径:MergeFileSplitRead 与删除向量如何跳过合并

Apache Paimon 源码学习

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

读路径在分岔口选择直读或归并
读路径在分岔口选择直读或归并AI 生成配图

源码:table/source/KeyValueTableRead.java、table/source/MergeTreeSplitGenerator.java、table/source/splitread/PrimaryKeyTableRawFileSplitReadProvider.java、operation/MergeFileSplitRead.java、operation/RawFileSplitRead.java、deletionvectors/*、table/source/AbstractBatchTableScan.java

学习目标 ​

  • 知道“是否合并”在哪里决定(rawConvertible)。
  • 理解 merge-on-read 的实现与代价(尤其是 value 过滤不能下推)。
  • 理解删除向量模式的不变式、写入端如何维护、读取端如何利用。

一句话结论 ​

合并与否在 scan 规划阶段就决定了(DataSplit.rawConvertible)。删除向量模式下,写入端保证“每个 key 在 L≥1 只有一条存活记录”,读取端就能像 append 表一样直接读文件,用位图过滤已删除行。

1. 路由 ​

KeyValueTableRead 按顺序尝试 SplitReadProvider,第一个 match() 的胜出:

java
new PrimaryKeyIndexedSplitReadProvider(...)        // 索引加速
new PrimaryKeyTableRawFileSplitReadProvider(...)   // ← 直接读文件,不合并
new MergeFileSplitReadProvider(...)                // ← merge-on-read
new IncrementalChangelogReadProvider(...)
new IncrementalDiffReadProvider(...)

PrimaryKeyTableRawFileSplitReadProvider.match:

java
!context.forceKeepDelete() && !dataSplit.isStreaming() && dataSplit.rawConvertible()
&& 所有文件都有 deleteRowCount 统计   // 老版本文件缺统计时保守走合并

2. rawConvertible 从哪来:MergeTreeSplitGenerator.splitForBatch ​

扫描阶段决定能否直接读
扫描阶段决定能否直接读AI 生成配图
java
boolean rawConvertible = files.stream().allMatch(f -> f.level() != 0 && withoutDeleteRow(f));
boolean oneLevel = 所有文件在同一层;
if (rawConvertible && (deletionVectorsEnabled || mergeEngine == FIRST_ROW || oneLevel)) {
    // 整个 bucket 免合并 → 按文件大小装箱,每组是一个 raw split
    return BinPacking.packForOrdered(files, max(fileSize, openFileCost), targetSplitSize)...;
}
// 否则:IntervalPartition 切 section → 按 section 装箱
// → 一个 split 恰好只有一个无 DELETE 的文件时才是 raw,其余 nonRaw
情况为什么能免合并
开启 DV + 无 L0 + 无 DELETE 行重复版本已被 DV 标记删除
first-row 引擎 + 无 L0写入时 lookup 已过滤重复
所有文件在同一层(如全量合并后)同层不重叠,每个 key 只出现一次

raw split 按文件装箱,并行度不受 key 范围限制;merge split 必须把重叠文件放一起(一个 section 不能拆),热点 bucket 只能一个 task 读。

3. 兜底:MergeFileSplitRead 的 merge-on-read ​

java
createReader(DataSplit split)
  ├─ 流读 / postpone bucket → createNoMergeReader   // 流读 delta/changelog,不合并
  └─ 批读 → createMergeReader(files, deletionFiles)
        dvFactory = DeletionVector.factory(fileIO, files, deletionFiles)   // 合并时也应用 DV
        for (section : new IntervalPartition(files).partition())
            readerForSection(section,
                section.size() > 1 ? overlappedSectionFactory      // 只下推主键过滤
                                   : nonOverlappedSectionFactory,  // 下推全部过滤
                ... ReducerMergeFunctionWrapper, mergeSorter)
        → ConcatRecordReader → DropDeleteReader → projectKey / projectOuter

为什么重叠 section 只下推主键过滤(源码注释的例子):

run1: (seq=1, k1, 100), (seq=2, k2, 200)
run2: (seq=3, k1, 10),  (seq=4, k2, 20)
若下推 "value >= 100",只读到 run1,新版本 k1=10、k2=20 被过滤,结果错误地返回旧值。

merge-on-read 的代价:多路归并(CPU + 多 reader 内存)、value 过滤无法下推到格式层、并行度受 section 限制。

4. 删除向量如何做到免合并 ​

删除向量用位图屏蔽旧行实现直读
删除向量用位图屏蔽旧行实现直读AI 生成配图

4.1 写入端维护的不变式 ​

在 level ≥ 1 的文件中,应用 DV 之后,每个 key 只有一条存活记录,且不存在 DELETE 记录。

三处维护:

  1. L0 上推时标记旧版本:LookupChangelogMergeFunctionWrapper.getResult() lookup 到旧值时拿到文件名 + 行号:
    java
    deletionVectorsMaintainer.notifyNewDeletion(fileName, rowPosition);
  2. DELETE 不进高层:DV 模式 forceDropDelete=true,dropDelete 条件含 || dvMaintainer != null;有删除行的 L0 文件升级时强制 CHANGELOG_WITH_REWRITE。
  3. 重写后清理旧 DV:notifyRewriteCompactBefore 调 dvMaintainer.removeDeletionVectorOf(file)。

提交一致性:BucketedDvMaintainer.writeDeletionVectorsIndex() 写 bucket 的 DV 索引文件,随 CompactResult 一起进入同一快照的 index manifest,读者不会看到“新文件有了、DV 还没有”。

4.2 扫描:跳过 L0 + 开启 value 过滤 ​

java
// CoreOptions
public boolean batchScanSkipLevel0() {
    if (deletionVectorsEnabled()) return !deletionVectorsMergeOnRead();
    return mergeEngine() == FIRST_ROW;
}
// AbstractBatchTableScan
snapshotReader.withLevelFilter(level -> level > 0).enableValueFilter();
  • 必须跳过 L0:L0 还没经过 lookup,高层旧版本没打 DV,直接读会重复。
  • enableValueFilter():每个 key 只有一条存活记录,按 value 统计裁剪整个文件是安全的。
  • SnapshotReaderImpl.generateSplits 一次性读 bucket 的 DV 索引(scanDvIndex),给 DataSplit 挂上每个文件的 DeletionFile(path, offset, length, cardinality)。

4.3 读取:RawFileSplitRead 直接读 + DV 过滤 ​

java
DeletionVector dv = dvFactory.get();
fileIndexResult = FileIndexEvaluator.evaluate(..., dv);     // 文件索引,感知 DV
if (!fileIndexResult.remain()) return new EmptyFileRecordReader<>();
reader = new DataFileRecordReader(...)                      // 原生格式 reader:向量化、谓词全下推
if (dv != null && !dv.isEmpty()) return new ApplyDeletionVectorReader(reader, dv);

过滤(ApplyDeletionFileRecordIterator.next):

java
while (true) {
    InternalRow next = iterator.next();
    if (next == null) return null;
    if (!deletionVector.isDeleted(returnedPosition())) return next;
}

DV 本体是 RoaringBitmap(BitmapDeletionVector 32 位 / Bitmap64DeletionVector 超大文件),序列化头部魔数区分。

5. 完整例子 ​

初始:L5 文件 f1,第 3 行 = (k1, v1)
① 写入 (k1, v2) → L0 文件 f0。此时批读跳过 L0 → 读到 v1(尚不可见)
② 合并把 L0 上推,outputLevel = L4
   lookup(k1, 从 L5) → 命中 (f1, pos=3) → DV[f1] = {3}
   写出 f2@L4: (k1, v2)
   提交:ADD f2、DELETE f0;index manifest ADD dv-index
③ 批读:split = [f1 + DV(f1), f2],rawConvertible = true
   f1 跳过第 3 行,f2 直接读 → k1 = v2 恰好一次,零归并

不开 DV 时,f1@L5 和 f2@L4 重叠 → non-raw split → 走归并。

6. 权衡 ​

维度merge-on-read(默认)删除向量
批读路径多路归并直接读 + 位图过滤
value 过滤下推仅单 run section全部文件,且可按统计裁剪文件
并行度受 section 限制按文件任意拆
数据可见性L0 立即可见L0 在合并前不可见(lookup-wait 下约一个 checkpoint);deletion-vectors.merge-on-read=true 可见但读变慢
写入代价低每次 L0 合并需 lookup
limit 下推支持limitPushdownEnabled() 显式关闭
流读不受影响不受影响

本质:把“每次读都去重”换成“写入/合并时一次性 lookup 标记”。

动手实验 ​

  • [ ] 建两张主键表(一张 'deletion-vectors.enabled'='true'),写相同数据;查 t$files、t$table_indexes;在 KeyValueTableRead.createReader 下断点对比 provider。
  • [ ] 断点:MergeTreeSplitGenerator.splitForBatch、PrimaryKeyTableRawFileSplitReadProvider.match、notifyNewDeletion。
  • [ ] 跑 DeletionVectorTest、BucketedDvMaintainerTest、DeletionVectorITCase。

自测题 ​

  1. rawConvertible 在哪里计算?有哪三种情况整个 bucket 免合并?
  2. 为什么 merge 模式不能下推 value 过滤?
  3. DV 模式为什么必须跳过 L0?代价是什么?怎么缓解?
  4. DV 文件和数据文件如何保证在同一个快照中可见?

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