Skip to content

S2 第 6 讲:读路径——为什么有的查询要归并、有的不用 ​

Apache Paimon 源码学习

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

主键表在两条读路径间分流
主键表在两条读路径间分流AI 生成配图

TL;DR:主键表有两条读路径——直接读(RawFileSplitRead,像 append 表一样顺序读文件)和合并读(MergeFileSplitRead,按 key 归并多个文件、同 key 取最新、最后去掉删除记录)。走哪条在规划阶段就由 split 上的 rawConvertible 决定:整个 bucket 都不在 L0、没有删除记录、且在同一层(或开启删除向量 / first-row 引擎)时整体直接读;否则切 section、装箱,只有一个文件且无删除记录的 split 才能直接读。合并读的代价不只是归并:重叠的 section 只能下推主键过滤;COUNT(*) 也没法只看元数据。

实验:5 张表、5 种文件状态 ​

labs/sql/lab08/read-path.sql 建了 5 张 bucket = 1 的主键表,每张造出一种典型的文件状态($files 实测):

表怎么写的文件
r_one写 1 次(key 1~100)1 个 L0 文件
r_disjoint写 2 次,key 不重叠(1~100、101~200)2 个 L0 文件
r_overlap写 2 次,key 重叠(1~100、51~150)2 个 L0 文件
r_compact同 r_overlap,再全量合并1 个 L5 文件(150 行)
r_delete写入、全量合并,再 DELETE WHERE k = 5L5(100 行)+ L0(1 行,deleteRowCount = 1)

每张表都用 SELECT v, COUNT(*) … GROUP BY v 查一次(为什么不用 COUNT(*),后面说),结果都正确:r_overlap 是 a 50 行、b 100 行,r_delete 是 99 行。

然后用 jdb 在规划(MergeTreeSplitGenerator)和两个 reader 上下断点:

bash
cd labs && PAIMON_VERSION=2.2-SNAPSHOT ./jdb-stacks.sh sql/lab08/read-path.sql jdb/s2-6-read-path.txt 300

一、两条读路径 ​

读路径由分流器逐个判断
读路径由分流器逐个判断AI 生成配图

两条读路径

KeyValueTableRead.reader(第 228~235 行)依次问每个 SplitReadProvider “这个 split 归你吗”:

  • PrimaryKeyTableRawFileSplitReadProvider:!forceKeepDelete && !isStreaming && split.rawConvertible(),再加上所有文件都有 deleteRowCount 统计——直接读;
  • MergeFileSplitReadProvider:兜底——合并读。
直接读 RawFileSplitRead合并读 MergeFileSplitRead
怎么读顺序读文件,和 append 表一样按 key 多路归并,同 key 交给合并函数
过滤全部过滤条件下推到文件读取重叠 section 只下推主键过滤
并行按文件装箱一个 section 不能拆
删除记录不需要处理(文件里没有)最后用 DropDeleteReader 过滤

二、5 张表的实录 ​

五张表

jdb 打印的读取方式:

r_one       RawFileSplitRead:169   split.dataFiles().size() = 1  rawConvertible = true
r_disjoint  MergeFileSplitRead:253 split.dataFiles().size() = 2  rawConvertible = false
              section.size() = 1  ×2
r_overlap   MergeFileSplitRead:253 split.dataFiles().size() = 2  rawConvertible = false
              section.size() = 2
r_compact   RawFileSplitRead:169   split.dataFiles().size() = 1  rawConvertible = true
r_delete    MergeFileSplitRead:253 split.dataFiles().size() = 2  rawConvertible = false
              section.size() = 2

前三个结论符合直觉:单文件直接读、重叠必须合并、全量合并后直接读。r_delete 也好理解:L0 里那条删除记录的 key 5 落在 L5 文件的 key 范围里,两者构成一个 2 个 run 的 section,删除要在读时生效。

r_disjoint 是反直觉的那个:两个文件 key 完全不重叠,却走了合并读。

三、为什么不重叠也要合并读:splitForBatch 的两步 ​

splitForBatch

MergeTreeSplitGenerator.splitForBatch(第 68~114 行)分两步:

第一步:整个 bucket 能不能免合并?

java
boolean rawConvertible = files.stream().allMatch(file -> file.level() != 0 && withoutDeleteRow(file));
boolean oneLevel = files 全在同一层;
if (rawConvertible && (deletionVectorsEnabled || mergeEngine == FIRST_ROW || oneLevel)) {
    return 按文件大小装箱,每组都是 rawConvertible;
}

jdb:r_compact 是 files.size() = 1, rawConvertible = true,直接返回。其余四张表都有 L0 文件,rawConvertible = false。

第二步:切 section,按 section 装箱。

java
List<List<DataFileMeta>> sections = new IntervalPartition(files, keyComparator).partition()...;
return packSplits(sections).stream()
        .map(f -> f.size() == 1 && withoutDeleteRow(f.get(0))
                ? SplitGroup.rawConvertibleGroup(f)      // 这个 split 只有一个文件 → 直接读
                : SplitGroup.nonRawConvertibleGroup(f))  // 否则 → 合并读

jdb:r_one 是 sections.size() = 1,装出来的 split 只有一个文件——直接读。r_disjoint 是 sections.size() = 2,但 packSplits 按 source.split.target-size(默认 128 MB)装箱,每个 section 至少按 source.split.open-file-cost(默认 4 MB)计——两个小 section 被装进了同一个 split,split 里有 2 个文件,于是打上 nonRawConvertible。

所以 r_disjoint 走的是 MergeFileSplitRead,但它的两个 section 各只有 1 个 run(jdb:section.size() = 1 两次)——每个 section 用的是“非重叠” reader,不需要多路归并,只是绕了一圈 KeyValue 的读取流程。

判断的单位是 split,不是“文件之间有没有重叠”。

四、合并读的隐藏代价:value 过滤不能下推 ​

合并读会阻止部分过滤下推
合并读会阻止部分过滤下推AI 生成配图

过滤下推

MergeFileSplitRead.withFilter(第 202~238 行)把过滤条件分成两组:

java
filtersForAll = 全部条件;
filtersForKeys = 只涉及主键的条件;
// 重叠的 section 用 filtersForKeys,单 run 的 section 用 filtersForAll(第 318~321 行)

源码注释给了原因。两个重叠的 run:

run1: (seq=1, k1, 100), (seq=2, k2, 200)
run2: (seq=3, k1, 10),  (seq=4, k2, 20)

如果把 value >= 100 下推下去,run2 整个被过滤,归并后 k1、k2 只剩旧版本 100、200——错的。

实验里对应的查询:

sql
SELECT * FROM r_overlap WHERE v = 'a' AND k BETWEEN 58 AND 62;
-- Empty set

k 58~62 在第二次写入中被更新成了 'b'。只有先归并拿到最新版本、再按 v = 'a' 过滤,结果才是正确的空集。代价是:对频繁更新、文件大量重叠的表,按非主键字段过滤的查询要读更多数据。这正是第 7 讲删除向量要解决的问题之一。

五、直接读的额外红利:COUNT(*) 只看元数据 ​

COUNT 下推

第一次做这个实验时,我用的是 SELECT COUNT(*) FROM r_one,结果 jdb 在 r_one、r_compact 上一次 reader 都没有命中。原因:Paimon 的 Flink source 实现了聚合下推(BaseDataTableSource.applyAggregates),所有 split 都能直接读时,COUNT(*) 直接用文件元数据里的行数相加(DataSplit.mergedRowCount,第 139~154 行,前提是 rawConvertible),根本不打开文件。

EXPLAIN SELECT COUNT(*) FROM r_compact;
  HashAggregate(isMerge=[true], select=[Final_COUNT(count1$0) AS EXPR$0])
  +- Exchange(distribution=[single])
     +- TableSourceScan(table=[[paimon, default, r_compact,
          aggregates=[grouping=[], aggFunctions=[Count1AggFunction()]]]], fields=[count1$0])

EXPLAIN SELECT COUNT(*) FROM r_overlap;
  HashAggregate(isMerge=[true], select=[Final_COUNT(count1$0) AS EXPR$0])
  +- Exchange(distribution=[single])
     +- LocalHashAggregate(select=[Partial_COUNT(*) AS count1$0])
        +- TableSourceScan(table=[[paimon, default, r_overlap]], fields=[k, v])

(节选自 Optimized Execution Plan;为排版把 r_compact 的 TableSourceScan 折成两行。)

r_compact 的计数已经在 TableSourceScan 里算好,Flink 只做最后一步汇总;r_overlap 则要把 k, v 真正读出来,由 Flink 的 LocalHashAggregate 逐行计数。同一条 SQL 在 Flink 1.20.1 上的计划写作 project=[k]、只读 k 一列,下推与否的结论相同。

r_overlap 两个文件行数之和是 200,实际只有 150 行(50 个 key 有两个版本),元数据算不出来,只能读出来归并再数。这也是为什么后面的读取实验改成了 GROUP BY v(按非分区字段分组不会被下推)。

小结:一张表判断读路径 ​

文件状态读路径本讲的表
整个 bucket:无 L0、无删除、同一层直接读(可按文件并行,COUNT(*) 下推)r_compact
开启删除向量 / first-row 引擎,且无 L0直接读第 7 讲
split 里只有一个文件、无删除直接读r_one
split 里多个文件(哪怕不重叠)合并读r_disjoint
key 重叠 / 有删除记录合并读,重叠 section 只下推主键过滤r_overlap、r_delete

生产启示 ​

  1. 合并(compaction)不只是为了减少文件数:全量合并后文件在同一层,读取可以整体走直接读、过滤全下推、COUNT(*) 秒出。
  2. 频繁更新的表,按非主键字段过滤会变慢:重叠越多,下推越少。考虑删除向量模式(第 7 讲)或定期全量合并。
  3. 用 EXPLAIN 看 COUNT(*) 是否下推:出现 aggregates=[...] 才是只看元数据。

下一讲:删除向量怎么把“读时去重”变成“写时标记”,让有更新的表也能直接读?

自测题:把 r_disjoint 的 source.split.target-size 调小,让两个 section 装进两个 split(例如 'source.split.target-size' = '4 mb',配合默认 4 MB 的 open-file-cost),读取会走哪条路?(答案:每个 split 只有一个文件、且无删除记录,两个 split 都直接读。按源码推导,未单独运行。)


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