1. 1. 1. Spill 保存什么状态,以及怎样恢复执行
  2. 2. 2. 整体架构
    1. 2.1. 2.1 从算子语义看分层
  3. 3. 3. 一次聚合怎样写出状态并恢复结果
  4. 4. 4. HashAggregation Spill 实现
    1. 4.1. 4.1 两种 Spiller
    2. 4.2. 4.2 HashAggregation:写出中间态,读回时合并
    3. 4.3. 4.3 GroupingSet::spill() — input spill
    4. 4.4. 4.4 GroupingSet::spill(rowIterator) — output spill
    5. 4.5. 4.5 reclaim 触发路径
    6. 4.6. 4.6 getOutputWithSpill — 归并输出
    7. 4.7. 4.7 mergeNextWithAggregates — 有序归并聚合
    8. 4.8. 4.8 mergeNextWithoutAggregates — distinct 归并
  5. 5. 5. HashJoin Spill 实现
    1. 5.1. 5.1 整体流程
    2. 5.2. 5.2 HashJoin:从 build spill 到下一轮恢复
    3. 5.3. 5.3 HashBuildSpiller
    4. 5.4. 5.4 setupSpiller
    5. 5.5. 5.5 ensureInputFits — 内存压力检测
    6. 5.6. 5.6 spillInput — 正在 spill 时的直通路径
    7. 5.7. 5.7 finishHashBuild — 收尾与传递
    8. 5.8. 5.8 reclaim — 内存仲裁触发路径
    9. 5.9. 5.9 postHashBuildProcess / setupSpillInput — 递归恢复
    10. 5.10. 5.10 processSpillInput — 从 spill 文件重建
    11. 5.11. 5.11 Probe 侧 Spill 详解
      1. 5.11.1. 5.11.1 HashBuildResult — Build/Probe 信息传递载体
      2. 5.11.2. 5.11.2 路径 1:input spill(探测侧主路径)
      3. 5.11.3. 5.11.3 路径 2:restore 轮的 probe(从 spill 文件读取 probe 行)
      4. 5.11.4. 5.11.4 路径 3:output spill(内存仲裁触发 reclaim)
    12. 5.12. 5.12 Probe 阶段回收:只存剩余输入还不够
      1. 5.12.1. 5.12.1 prepareForSpillRestore — 多轮恢复的关键
      2. 5.12.2. 5.12.2 Probe 侧 Spill 全生命周期
  6. 6. 6. OrderBy Spill 实现
    1. 6.1. 6.1 两种 Spiller
    2. 6.2. 6.2 OrderBy:输入 run 与剩余输出
    3. 6.3. 6.3 SortBuffer::spillInput
    4. 6.4. 6.4 SortBuffer::spillOutput
    5. 6.5. 6.5 SortBuffer::finishSpill
    6. 6.6. 6.6 SortBuffer::getOutputWithSpill
  7. 7. 7. Spiller 层详解
    1. 7.1. 7.1 SpillerBase — 抽象基类
    2. 7.2. 7.2 SpillRun — 内存中的 spill 缓冲
    3. 7.3. 7.3 fillSpillRuns — 数据填充
    4. 7.4. 7.4 SpillRun 为何切分 / maxSpillRunRows 的作用
    5. 7.5. 7.5 SpillRun、batch 和 file 是三个尺度
    6. 7.6. 7.6 行列转换:extractSpill 与 extractSpillVector
    7. 7.7. 7.7 RowContainer 到列向量:为何要提取而不是 memcpy
      1. 7.7.1. 7.7.1 extractSpillVector — 控制每批大小
      2. 7.7.2. 7.7.2 extractSpill — 行式 → 列式核心转换
    8. 7.8. 7.8 RowContainer 行布局与列提取原理
    9. 7.9. 7.9 Accumulator 中间态提取
    10. 7.10. 7.10 writeSpill 内的 Batch 切分(kTargetBatchRows / kTargetBatchBytes)
    11. 7.11. 7.11 runSpill — 并发写入
    12. 7.12. 7.12 有序写出与异步任务的收尾
    13. 7.13. 7.13 finishSpill — 收尾
    14. 7.14. 7.14 各子类 Spiller 详解
      1. 7.14.1. 7.14.1 HashBuildSpiller
      2. 7.14.2. 7.14.2 AggregationInputSpiller
      3. 7.14.3. 7.14.3 AggregationOutputSpiller
      4. 7.14.4. 7.14.4 SortInputSpiller
      5. 7.14.5. 7.14.5 SortOutputSpiller
      6. 7.14.6. 7.14.6 NoRowContainerSpiller / MergeSpiller
    15. 7.15. 7.15 各子类对比
  8. 8. 8. 基础设施层详解
    1. 8.1. 8.1 SpillConfig — 配置层
    2. 8.2. 8.2 配置与内存触发条件
    3. 8.3. 8.3 SpillStats — 统计层
    4. 8.4. 8.4 SpillFile — 文件 I/O 层
    5. 8.5. 8.5 文件与分区:当前接口的变化
      1. 8.5.1. 8.5.1 SpillFileInfo — 元数据
      2. 8.5.2. 8.5.2 SpillWriter — 写路径
      3. 8.5.3. 8.5.3 SpillReadFile — 读路径
    6. 8.6. 8.6 SpillPartitionId — 层级分区标识
    7. 8.7. 8.7 SpillState — 分区状态管理
    8. 8.8. 8.8 SpillPartition / SpillPartitionSet — 读回路径
    9. 8.9. 8.9 两种读回路径,以及有限归并 fan-in
    10. 8.10. 8.10 读流层:SpillMergeStream / BatchStream
  9. 9. 9. 递归 Spill 流程图
  10. 10. 10. 关键设计决策总结
    1. 10.1. 10.1 限制、统计和清理
    2. 10.2. 10.2 有序 vs. 无序 Spill
    3. 10.3. 10.3 两阶段 Spill(Input/Output Spiller 分离)
    4. 10.4. 10.4 内存预留机制(proactive spill)
    5. 10.5. 10.5 HashBuild 多线程 Spill 协调
    6. 10.6. 10.6 递归 Spill(HashJoin 专属)
    7. 10.7. 10.7 probed flag Spill(Right/Full Outer Join)
  11. 11. 11. C++ 编码品味:值得学习的细节
    1. 11.1. 11.1 从当前实现学到的几个编程约束
    2. 11.2. 11.2 SCOPE_EXIT:让清理逻辑紧贴触发逻辑
    3. 11.3. 11.3 folly::makeGuard:异步任务的"排水阀"
    4. 11.4. 11.4 跨线程异常传递:exception_ptr 模式
    5. 11.5. 11.5 "1 + N" 并发策略:第一个任务不切线程
    6. 11.6. 11.6 CHECK vs DCHECK:按热度分配断言代价
    7. 11.7. 11.7 FOLLY_LIKELY / FOLLY_UNLIKELY:把预期写进代码
    8. 11.8. 11.8 TestValue::adjust:无侵入式测试注入
    9. 11.9. 11.9 NanosecondTimer:RAII 计时的最小实现
    10. 11.10. 11.10 SpillPartitionId:把所有接口设施配齐
    11. 11.11. 11.11 withWLock / withRLock:锁的最小化持有
    12. 11.12. 11.12 SpillRun::sorted:防止双重排序的状态位
    13. 11.13. 11.13 pool_->release():主动归还,而非等待回收
    14. 11.14. 11.14 succinctBytes():有信息量的日志
    15. 11.15. 11.15 testing* 前缀:测试可观测性的统一约定
    16. 11.16. 11.16 constexpr 局部常量:让魔数有名字
    17. 11.17. 11.17 编码品味小结
  12. 12. 12. 从恢复需求推导 Spill 的分层
    1. 12.1. 12.1 资源有限时扩大可执行工作集的范围
    2. 12.2. 12.2 分层:每层只管自己的事
    3. 12.3. 12.3 分区:把"全量恢复"变成"分批恢复"
    4. 12.4. 12.4 有序与无序的分叉:根据"读回时的需求"决定
    5. 12.5. 12.5 两阶段 Spill:尊重算子的生命周期
    6. 12.6. 12.6 SpillRun 切分:在多个约束之间找到平衡点
    7. 12.7. 12.7 批量行列转换:峰值内存与 I/O 效率的双重优化
    8. 12.8. 12.8 主动预留 vs. 被动仲裁:两道防线
    9. 12.9. 12.9 NoRowContainerSpiller:让接口适应现实,而非让现实适应接口
    10. 12.10. 12.10 HashJoinBridge:异步协调的最小化接口
    11. 12.11. 12.11 DictionaryVector wrap:零拷贝分区
    12. 12.12. 12.12 小结:设计模式归纳
    13. 12.13. 12.13 代码阅读与验证边界
  13. 13. 13. 以恢复工作集检验 Spill 的设计取舍
Macduan Notes

Velox Spiller

Spill 的目标是在保留查询语义的前提下,把暂时放不下的执行状态移到外部存储,再按算子需要的方式读回。真正困难的部分是:写出哪种状态、何时能释放原内存、怎样恢复、如何与其他 Driver 协调,以及异常时谁负责清理。

执行流程与相关实现核对于 2026-09-29,Velox 源码版本为 48883e8521b2。下文源码节选、历史资料和宿主集成引用各自注明版本;教学输入用于解释状态变化,未作为性能基准运行。

1. Spill 保存什么状态,以及怎样恢复执行

Spill 是算子的外部执行机制:在内存不足时,把仍需使用的状态按可恢复的形式写出,释放相应内存,再分批恢复处理。它需要保持算子的结果语义,也需要控制恢复工作集,因此文件写入只是其中一段。

使用场景 需要保存的状态 恢复后的主要工作
HashJoin build 行、相关 probe 行及必要匹配状态 按轮次恢复相应分区,重建表并继续 Join
HashAggregation 输入阶段 分组 key 与当前聚合状态 读取有序 spill 数据,归并相同 group 的状态
OrderBy 输入阶段 按排序键组织的有序 run 多路归并,按要求产生全局有序结果
算子输出阶段 尚未返回给下游的结果状态或剩余有序行 按具体算子的输出协议读回,避免重放已输出部分

算子决定哪些状态可写、何时可回收和怎样恢复;Spiller 负责收集、分区、排序、批量提取和写出;文件层负责序列化页及其 I/O。需要有序读回的路径与无序分区恢复的路径不同,不能把所有 spill 都解释成“写一个文件,再读回来”。

Spill 也不会消除资源限制。序列化缓冲、恢复内存、磁盘容量、递归深度和数据倾斜都可能限制执行。下面从基础设施走到各算子的完整路径,再展开资源保护、异常传播和并行写出的工程细节。

图 1:算子决定语义,Spiller 组织写出,文件层负责序列化页;读回方式由算子需求决定。
图 1:算子决定语义,Spiller 组织写出,文件层负责序列化页;读回方式由算子需求决定。

本文于 2026-09-19 对照 Velox 1d1b76567870 重写,按基础设施、写出机制、HashJoin、Aggregation、OrderBy 和资源生命周期展开。文中的“当前”均指这一提交。

代码片段分别标明源码节选或流程示意;流程示意省略统计、异常包装与无关分支,不是可独立编译的程序。历史资料与宿主集成保留各自版本,不能据此推断它们组成了经过构建验证的发行版本。

2. 整体架构

2.1 从算子语义看分层

层次 对象 责任
算子 HashBuild/HashProbe、GroupingSet、SortBuffer 决定 spill 时机、数据含义和恢复流程
写出策略 SpillerBase 及各子类 收集行、分区、排序、提取 batch
分区状态 SpillState、SpillPartitionId、SpillPartition 组织 writer、文件和读回入口
文件 I/O SpillWriter、SpillReadFile 序列化页的写入与读取
资源 SpillConfig、spill pool、SpillStats 配置、临时内存、限制和观测

当前 HashBuildSpiller 定义在 HashBuild.h;AggregationInput/OutputSpiller 定义在 GroupingSet.h;通用 NoRowContainerSpiller、SortInput/OutputSpiller 等在 Spiller.h。不要按旧的单文件类图找全部子类。参见 SpillerBase、HashBuildSpiller、Aggregation spillers。

Spill 并不承诺查询一定成功。磁盘限制、不可回收状态、递归层数和恢复内存仍可能让查询失败。

Velox 的 spill 系统是一个分层架构,每层职责明确:

Spill 各层职责
图 2:Spill 各层职责。已按当前实现修正标注,具体约束见相邻正文。
Fig. Spiller 六层架构:算子触发 → Spiller 抽取分区排序 → State / Partition 管理 → 文件 I/O

核心设计理念:

  • HashJoin 与聚合输入可以按 hash key 分区,使恢复按较小范围推进;OrderBy 的有序 run 则使用单分区并依赖外部归并。分区与排序解决不同问题,不能把“所有 spill 都按 hash 分区”作为统一规则。
  • 递归 spill:若恢复某分区后仍 OOM,可对该分区再次 spill(最多 maxSpillLevel 层)
  • 有序与无序分离:sort-based 算子(Aggregation/OrderBy)spill 时保序,以支持归并输出;HashJoin 无需排序
  • 两阶段 spill:Aggregation/OrderBy 区分 input 阶段与 output 阶段,各用不同的 Spiller

3. 一次聚合怎样写出状态并恢复结果

用 SELECT k, sum(v) FROM t GROUP BY k 贯穿写出与恢复。第一批输入 [(1,10), (2,7)] 得到两个聚合状态;内存压力触发 spill 后,后续输入 [(1,5), (2,3)] 产生新的状态。最终必须得到 (1,15)、(2,10),不能把两个 run 简单拼接成四行,也不能把已经输出的组再次输出。

阶段调用路径与状态恢复所需的不变量
输入期间积累HashAggregation → GroupingSet → RowContainer保留 group keys 与 accumulator 状态;真实 aggregate 的中间态可能不是最终 SQL 值。
选择可回收状态主动内存检查或 reclaimer → GroupingSet::spill回收必须处于算子允许的阶段;input spiller 和 output spiller 处理的状态不同。
分区并形成 runAggregationInputSpiller → fillSpillRuns → runSpill同一个 group 必须落在同一 hash 分区;每个 run 按 group keys 排序。
提取与写出extractSpill → key vectors / extractAccumulators → writeSpill → SpillState行式状态被转换为可序列化的列向量;不是把 RowContainer 的指针和原始内存块写到盘上。
收束异步写任务取得全部已提交任务结果,finishFile / finishSpill 形成文件集合异常路径也必须等任务不再访问 run / container 后,才能安全释放它们。
释放状态并继续输入GroupingSet 清理已落盘的表状态,再接后续批次后续相同 key 可形成新的中间态;spill 文件仍由对应分区的生命周期管理。
完成并读回noMoreInput → getOutputWithSpill → ordered reader / TreeOfLosers同分区各有序 run 归并,同 key 的 accumulator 经 addSingleGroupIntermediateResults 合并,最终提取结果。

sum 的中间态在此便于直观展示;avg 等聚合还需要计数或复合状态,具体形状由 aggregate 的 intermediate type 与提取 / 合并接口决定。文件、run、序列化 batch 也不是一一对应:run 控制排序边界,batch 控制临时工作集,file 受写出与切分策略约束。

源码主线:GroupingSet::spill → SpillerBase::runSpill → getOutputWithSpill → mergeNextWithAggregates。下文先解释各算子的恢复要求,再进入共同的 run、序列化、文件与流实现。

4. HashAggregation Spill 实现

4.1 两种 Spiller

4.2 HashAggregation:写出中间态,读回时合并

GroupingSet::spill() 为输入阶段建立 AggregationInputSpiller,按 grouping keys 分区、按 key 排序,写出 key 和 accumulator 的 spillType,然后清掉当前表。相同 key 在不同 spill run 中出现是正常情况,读回时必须再次合并。

GroupingSet 对 HashStringAllocator 使用 freezeAndExecute,防止并行提取过程中 aggregate 修改非线程安全 allocator。它不是给 allocator 加一把通用锁,而是要求这个阶段不要做被禁止的可变操作。

getOutputWithSpill 为每个分区建立 ordered reader。nextWithEquals 告诉合并器后续 key 是否相等,GroupingSet 维护当前 group 的 merge state,通过 aggregate 的中间态合并逻辑更新,直到 key 结束才形成最终结果。DISTINCT 则走 withoutAggregates 分支,只需要按 key 去重。

图 9:每个 run 中的同 key 状态在 ordered merge 时重新聚合,不能简单拼接 spill 文件。
图 3:每个 run 中的同 key 状态在 ordered merge 时重新聚合,不能简单拼接 spill 文件。

输出阶段还有 AggregationOutputSpiller:已经聚合完成但尚未输出的组可以从 rowIterator 起写出。这条路径与输入阶段重建/合并的目标不同,不能重放已经交给下游的结果。当前实现用 input/output spiller 的互斥条件约束两条路径。

HashAggregation 的 spill 逻辑集中在 GroupingSet 中,区分两个阶段:

inputSpiller_  (AggregationInputSpiller):
  - 在 input 处理阶段触发(hash 表太大)
  - needSort = true:按 group key 排序
  - 分区:按 hash(group key)
  - 目的:spill 后 hash 表清空,腾出内存继续处理输入

outputSpiller_ (AggregationOutputSpiller):
  - 在 output 处理阶段触发(输出时内存不足)
  - needSort = false:已是输出顺序(hash 表 scan 顺序),不需要重排
  - 无分区(HashBitRange 为空)
  - 目的:将剩余输出行先 spill 到磁盘,之后边读边输出

两者互斥:同一个 GroupingSet 要么用 inputSpiller_,要么用 outputSpiller_,不会同时存在。

4.3 GroupingSet::spill() — input spill

由 HashAggregation::reclaim() 或 ensureInputFits 调用:

流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

void GroupingSet::spill() {
  if (table_ == nullptr || table_->numDistinct() == 0) return;

  auto* rows = table_->rows();
  VELOX_CHECK_NULL(outputSpiller_);

  // 首次 spill:初始化 inputSpiller_
  if (inputSpiller_ == nullptr) {
    // sortingKeys:按所有 group key 列排序(升序,nulls 在末尾)
    const auto sortingKeys = SpillState::makeSortingKeys(
        std::vector<CompareFlags>(rows->keyTypes().size()));
    inputSpiller_ = std::make_unique<AggregationInputSpiller>(
        rows,
        makeSpillType(),  // key 列 + accumulator 列的类型
        HashBitRange(startPartitionBit, startPartitionBit + numPartitionBits),
        sortingKeys,
        spillConfig_,
        spillStats_);
  }

  // 冻结 HashStringAllocator:spill 可能多线程执行,防止 accumulator 并发分配
  rows->stringAllocator().freezeAndExecute([&]() {
    inputSpiller_->spill();
  });

  // distinct 聚合:记录本次 spill 生成的文件数(用于后续归并去重)
  if (isDistinct() && numDistinctSpillFilesPerPartition_.empty()) {
    // 记录每个 partition 的文件数,文件 ID < 该数的是 distinct 文件
    for (int partition = 0; partition < maxPartitions; ++partition) {
      numDistinctSpillFilesPerPartition_[partition] =
          inputSpiller_->state().numFinishedFiles(SpillPartitionId(partition));
    }
  }

  // 清空 hash 表(所有行已 spill 到磁盘)
  table_->clear(/*freeTable=*/true);
}

makeSpillType() 生成 spill 类型:

RowTypePtr GroupingSet::makeSpillType() const {
  std::vector<TypePtr> types;
  // Group key 列
  for (auto& hasher : hashers_) types.push_back(hasher->type());
  // Accumulator 的中间状态类型(不是最终结果类型)
  for (auto& accumulator : accumulators(false)) {
    types.push_back(accumulator.spillType());
  }
  return ROW(names, types);
}

4.4 GroupingSet::spill(rowIterator) — output spill

由 HashAggregation::reclaim() 在 output 阶段调用:

流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

void GroupingSet::spill(const RowContainerIterator& rowIterator) {
  VELOX_CHECK(!hasSpilled()); // 确保没有 inputSpiller_

  auto* rows = table_->rows();
  outputSpiller_ = std::make_unique<AggregationOutputSpiller>(
      rows, makeSpillType(), spillConfig_, spillStats_);

  rows->stringAllocator().freezeAndExecute([&]() {
    outputSpiller_->spill(rowIterator); // 从 rowIterator 指定的位置开始 spill
  });
  table_->clear(/*freeTable=*/true);
}

AggregationOutputSpiller::spill(rowIterator) 内部:

void AggregationOutputSpiller::spill(const RowContainerIterator& startRowIter) {
  // 只标记单一分区(partition 0),无 hash 分区
  state_.setPartitionSpilled(SpillPartitionId(0));
  SpillerBase::spill(&startRowIter); // 从 rowIterator 位置开始迭代
}

AggregationOutputSpiller 的收尾遵循输出阶段的已有顺序,只有满足对应最后一个 run 的条件时才结束文件。它不需要把每个内部提取批次或 run 都作为独立的待排序输入文件;具体边界与 AggregationInputSpiller 的有序 run 分开理解。

4.5 reclaim 触发路径

流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

void HashAggregation::reclaim(uint64_t targetBytes,
                               memory::MemoryReclaimer::Stats& stats) {
  // 1. 先尝试 compact(清理 accumulator 的字符串存储),代价低
  const auto compactedBytes = groupingSet_->compact();
  if (compactedBytes >= targetBytes) return;

  if (!canSpill()) { /* ... */ return; }

  if (groupingSet_->hasSpilled()) {
    if (isOutputProcessing) {
      // 已在输出阶段且已 spill:无法再次 spill,报告无法回收
      return;
    }
  }

  if (isOutputProcessing) {
    // 在输出阶段:spill 从 resultIterator_ 位置起的剩余行
    groupingSet_->spill(resultIterator_);
  } else {
    // 在输入阶段:spill 所有行(hash 表全量 dump)
    groupingSet_->spill();
  }
}

4.6 getOutputWithSpill — 归并输出

当 groupingSet_->hasSpilled() 为 true 时,getOutput 路径切换为:

流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

bool GroupingSet::getOutputWithSpill(int32_t maxOutputRows,
                                      int32_t maxOutputBytes,
                                      const RowVectorPtr& result) {
  if (outputSpillPartition_ == -1) {
    // 首次进入:初始化归并结构
    mergeRows_ = std::make_unique<RowContainer>(keyTypes, ...); // 临时 RowContainer
    initializeAggregates(*mergeRows_, false);

    // 将所有 spill 文件信息收集到 spillPartitionSet_
    inputSpiller_->finishSpill(spillPartitionSet_);
    // 或:outputSpiller_->finishSpill(spillPartitionSet_);

    removeEmptyPartitions(spillPartitionSet_);
    prepareNextSpillPartitionOutput(); // 打开第一个 partition 的 TreeOfLosers
  }

  // 逐 partition 归并输出
  return mergeNext(maxOutputRows, maxOutputBytes, result);
}

bool GroupingSet::prepareNextSpillPartitionOutput() {
  merge_ = nullptr;
  if (spillPartitionSet_.empty()) return false;

  auto it = spillPartitionSet_.begin();
  outputSpillPartition_ = it->first.partitionNumber();
  // createOrderedReader:构建 N 路败者树(含文件数超限时的预归并)
  merge_ = it->second->createOrderedReader(*spillConfig_, pool_, spillStats_);
  spillPartitionSet_.erase(it);
  return true;
}

4.7 mergeNextWithAggregates — 有序归并聚合

流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

bool GroupingSet::mergeNextWithAggregates(
    int32_t maxOutputRows, int32_t maxOutputBytes, const RowVectorPtr& result) {

  bool nextKeyIsEqual{false};
  for (;;) {
    // 从败者树取下一行(带 equal 标志:与上一行 key 是否相同)
    auto [stream, equal] = merge_->nextWithEquals();

    if (stream == nullptr) {
      // 当前 partition 耗尽:输出已积累的行
      extractSpillResult(result);
      if (result->size() > 0) return true;
      // 换下一个 partition
      if (!prepareNextSpillPartitionOutput()) return false;
      continue;
    }

    if (!nextKeyIsEqual) {
      // 新的 group key:在 mergeRows_ 中创建新行,初始化 accumulators
      mergeState_ = mergeRows_->newRow();
      initializeRow(*stream, mergeState_); // 复制 key 列 + 初始化 accumulator
    }
    // 用当前行更新 accumulator(partial merge)
    updateRow(*stream, mergeState_);
    nextKeyIsEqual = equal;
    stream->pop();

    // 达到行数/字节数上限:输出
    if (!nextKeyIsEqual && (mergeRows_->numRows() >= maxOutputRows
                            || mergeRowBytes() >= maxOutputBytes)) {
      extractSpillResult(result);
      return true;
    }
  }
}

关键:merge_->nextWithEquals() 利用已排序的特性,保证相同 key 的行连续出现,从而可以在流式归并的同时完成聚合(无需将整个 group 缓存在内存中)。

4.8 mergeNextWithoutAggregates — distinct 归并

对于 SELECT DISTINCT 场景,spill 文件分为两类:

  • distinct 文件(文件 ID < numDistinctSpillFilesPerPartition_):已输出过的 distinct 行
  • input 文件(文件 ID >= 阈值):新输入的行

归并时:若某个相同 key 的 run 中包含来自 distinct 文件的行,则跳过该 run(因为已经输出过);否则输出该 run 中的一行。


5. HashJoin Spill 实现

5.1 整体流程

5.2 HashJoin:从 build spill 到下一轮恢复

build peers 的当前汇合入口是 OperatorCtx::allPeersFinished,内部仍使用 Task 的 Driver barrier,再返回保活 Driver 的 peer Operator 引用。HashBuild 在最后到达者中接管各局部表与 spill 文件;变更的是接口与访问粒度,不是“两侧按相同分区逐轮恢复”的协议。见 finishHashBuild。

HashBuildSpiller 按 join key hash 分区,数据不需要排序。第一次 spill 写出当前 RowContainer;后续输入可直接分区写出。finishHashBuild 收集各 build peer 的分区,将结果交给 HashJoinBridge。

Probe 获取 spilled partition IDs 后,把对应输入保存到自己的 input spiller。当前主要通过 NoRowContainerSpiller 写 RowVector;不要凭旧命名寻找一个独立的 HashProbeSpiller 类。

图 6:build 和 probe 必须沿同一分区路径保存数据;恢复顺序由 bridge 协调。
图 4:build 和 probe 必须沿同一分区路径保存数据;恢复顺序由 bridge 协调。

一轮 probe 完成后,bridge 选择待恢复分区并按文件分 shard 给 build peers。HashBuild 通过 unordered reader 重建表,Probe 读取同一 restoredPartitionId 的输入。若仍放不下,setupSpiller 推进 hash 位窗口,创建子分区。

当前 join 是整表回收路径:bridge 明确约束非空 table 与本轮 spilledPartitionIds 的互斥关系。不能画成任意保留一部分 bucket 的混合表,除非另一个代码路径明确支持它。

达到递归限制时,setupSpiller 禁用进一步 spill 并记录指标,查询若仍不能获取内存就可能失败。数据倾斜尤其需要关注:同一 key 的大量重复行不会因为增加 hash bits 自动分散。

HashJoin spill 同时涉及 build 状态、对应的 probe 输入,以及 probe 阶段尚未排出的结果。Build 侧按 join key 保存表中的行;Probe 侧依据 bridge 发布的分区信息保存相同分区的输入。若内存仲裁发生在 probe 阶段,还要先保存当前输入尚未产出的结果,防止恢复后丢行或重复输出。

Build / Probe / Restore 三阶段
图 5:Build / Probe / Restore 三阶段。已按当前实现修正标注,具体约束见相邻正文。
Fig. HashJoin spill 三阶段:Build 落盘 → Probe 同步分区 → Restore 逐分区递归重建

5.3 HashBuildSpiller

流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

class HashBuildSpiller : public SpillerBase {
  const bool spillProbeFlag_;   // 是否需要 spill probed 标记(right/full join)
  bool spillTriggered_{false};  // 一旦触发 spill,后续所有输入直接 spill

  // spill() 重载1:spill RowContainer 中的所有行(内存仲裁触发)
  void spill() {
    spillTriggered_ = true;
    SpillerBase::spill(nullptr);  // → fillSpillRuns → runSpill → writeSpill
  }

  // spill() 重载2:spill 一个 partition 的输入向量(spillTriggered 后的直通路径)
  void spill(const SpillPartitionId& partitionId, const RowVectorPtr& spillVector) {
    VELOX_CHECK(spillTriggered_);
    if (!state_.isPartitionSpilled(partitionId)) {
      state_.setPartitionSpilled(partitionId);
    }
    state_.appendToPartition(partitionId, spillVector);
  }
};

**重写 extractSpill**:对于 right/full outer join,需要在 spill 数据中附加一列 bool probed_flag,记录该行是否已被 probe 侧命中:

void HashBuildSpiller::extractSpill(folly::Range<char**> rows, RowVectorPtr& resultPtr) {
  // 提取普通列
  for (auto i = 0; i < types.size(); ++i) {
    container_->extractColumn(rows.data(), rows.size(), i, result->childAt(i));
  }
  // 若需要 probed flag,从 RowContainer 的标记位提取
  if (spillProbeFlag_) {
    auto* flagVector = result->childAt(types.size())->asFlatVector<bool>();
    for (auto i = 0; i < rows.size(); ++i) {
      flagVector->set(i, container_->hasProbedFlag(rows[i]));
    }
  }
}

5.4 setupSpiller

在 HashBuild 初始化时和每次 restore 时调用:

流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

void HashBuild::setupSpiller(SpillPartition* spillPartition) {
  if (!canSpill()) return;

  // 首次调用:构造 spillType_(在普通列后可能追加 bool 列用于 probed flag)
  if (spillType_ == nullptr) {
    spillType_ = hashJoinTableSpillType(tableType_, joinType_);
    if (needProbedFlagSpill_) {
      spillProbedFlagChannel_ = spillType_->size() - 1;
      // 初始化 false 常量向量(build 侧 spill 时所有行未被 probe)
      spillProbedFlagVector_ = std::make_shared<ConstantVector<bool>>(
          pool(), 0, /*isNull=*/false, BOOLEAN(), false);
    }
  }

  uint8_t startPartitionBit = config->startPartitionBit;

  if (spillPartition != nullptr) {
    // Restore 场景:从 spill 文件中读取数据重建 hash 表
    spillInputReader_ = spillPartition->createUnorderedReader(
        config->readBufferSize, pool(), spillStats_.get());
    restoringPartitionId_ = spillPartition->id();

    // 递归 spill:bit offset 右移一层
    startPartitionBit = partitionBitOffset(
        spillPartition->id(), startPartitionBit, numPartitionBits) + numPartitionBits;

    // 若超过最大 spill 层数,禁用本轮 spill(只能尽力处理)
    if (config->exceedSpillLevelLimit(startPartitionBit)) {
      exceededMaxSpillLevelLimit_ = true;
      return;
    }
  }

  spiller_ = std::make_unique<HashBuildSpiller>(
      joinType_, restoringPartitionId_,
      table_->rows(), spillType_,
      HashBitRange(startPartitionBit, startPartitionBit + config->numPartitionBits),
      config, spillStats_.get());
}

5.5 ensureInputFits — 内存压力检测

每次 addInput 前调用,是 proactive spill 的入口:

流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

void HashBuild::ensureInputFits(RowVectorPtr& input) {
  if (!canSpill() || spiller_ == nullptr || spiller_->spillTriggered()) return;

  // 估算本批 input 引起的内存增量
  const auto tableIncrementBytes = table_->hashTableSizeIncrease(input->size());
  const auto rowContainerIncrementBytes = rows->sizeIncrement(input->size(), ...);
  const auto incrementBytes = rowContainerIncrementBytes + tableIncrementBytes;

  // 若当前预留充足,直接返回
  if (availableReservationBytes >= minReservationBytes) {
    if (freeRows > input->size() && outOfLineFreeBytes >= flatBytes) return;
    if (pool()->availableReservation() > 2 * incrementBytes) return;
  }

  // 尝试扩容 reservation(扩容本身可能触发内存仲裁,进而触发 spill)
  const auto targetIncrementBytes = std::max<int64_t>(
      incrementBytes * 2,
      currentUsage * spillConfig_->spillableReservationGrowthPct / 100);

  {
    Operator::ReclaimableSectionGuard guard(this); // 标记为可被仲裁
    if (pool()->maybeReserve(targetIncrementBytes)) {
      if (spiller_->spillTriggered()) {
        // 扩容过程中仲裁器触发了 spill,不需要这份预留了
        pool()->release();
      }
      return;
    }
  }
  // maybeReserve 失败(OOM),后续仲裁器会调用 reclaim()
}

5.6 spillInput — 正在 spill 时的直通路径

一旦 spiller_->spillTriggered() 为 true,所有新输入不再进入 hash 表,而是直接按 partition 分发到 spill 文件:

流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

void HashBuild::spillInput(const RowVectorPtr& input) {
  if (!canSpill() || spiller_ == nullptr || !spiller_->spillTriggered()
      || !activeRows_.hasSelections()) return;

  // 1. 重新计算所有行的 hash(与 addInput 中的计算可能重复,代价可接受)
  computeSpillPartitions(input); // 结果存入 spillPartitions_[row]

  // 2. 按 partition 分组,并从 activeRows_ 中移除(防止重复插入 hash 表)
  for (auto row = 0; row < numInput; ++row) {
    if (!activeRows_.isValid(row)) continue;
    activeRows_.setValid(row, false); // 从 active 中移除
    rawSpillInputIndicesBuffers_[spillPartitions_[row]][numSpillInputs_[partition]++] = row;
  }

  // 3. 为非 spill 原始输入构建 spillChildVectors_(key + dependent + probed_flag)
  maybeSetupSpillChildVectors(input);

  // 4. 按 partition 逐批 spill
  for (uint32_t partition = 0; partition < numSpillInputs_.size(); ++partition) {
    spillPartition(partition, numSpillInputs_[partition],
                   spillInputIndicesBuffers_[partition], input);
  }
}

5.7 finishHashBuild — 收尾与传递

所有 build 线程完成后,最后一个线程执行:

流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

bool HashBuild::finishHashBuild() {
  // ...(等待所有 peer 线程完成)

  SpillPartitionSet spillPartitions;

  // 收集所有 peer 的 spill 文件(多线程 build 的各自 spill 合并)
  for (auto* build : otherBuilds) {
    if (build->spiller_ != nullptr) {
      build->spiller_->finishSpill(spillPartitions);
    }
  }
  if (spiller_ != nullptr) {
    spiller_->finishSpill(spillPartitions);
    removeEmptyPartitions(spillPartitions);
  }

  // 构建 hash 表(只包含未 spill 的行)
  table_->prepareJoinTable(std::move(otherTables), ...);

  // 注册 tableSpillFunc:允许 join bridge 在 probe 后 spill hash 表
  HashJoinTableSpillFunc tableSpillFunc;
  if (canReclaim()) {
    tableSpillFunc = [hashBitRange = spiller_->hashBits(), ...](auto table) {
      return spillHashJoinTable(table, ...);
    };
  }

  // 将 hash 表 + spill partitions 一起传给 join bridge
  joinBridge_->setHashTable(table, std::move(spillPartitions),
                            joinHasNullKeys_, std::move(tableSpillFunc));
}

5.8 reclaim — 内存仲裁触发路径

内存仲裁器调用 reclaim() 强制 spill:

流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

void HashBuild::reclaim(uint64_t, memory::MemoryReclaimer::Stats& stats) {
  // 状态检查:只有 kRunning/kWaitForBuild/kYield 状态且不在非可回收区才能 spill
  if (nonReclaimableState()) {
    ++stats.numNonReclaimableAttempts;
    return;
  }

  // 暂停所有 peer driver(task->pauseRequested() 保证)
  const std::vector<Operator*> operators =
      task->findPeerOperators(pipelineId, this);

  // 检查所有 peer 都可被回收
  for (auto* op : operators) {
    if (static_cast<HashBuild*>(op)->nonReclaimableState()) {
      ++stats.numNonReclaimableAttempts;
      return;
    }
  }

  // 收集所有 peer 的 spiller,统一 spill 整个 hash 表
  std::vector<HashBuildSpiller*> spillers;
  for (auto* op : operators) {
    spillers.push_back(static_cast<HashBuild*>(op)->spiller_.get());
  }
  spillHashJoinTable(spillers, config); // 核心:spill 所有 peer 的 hash 表

  // 清空 hash 表(已 spill,不再需要)+ 释放 memory pool reservation
  for (auto* op : operators) {
    static_cast<HashBuild*>(op)->table_->clear(true);
    static_cast<HashBuild*>(op)->pool()->release();
  }
}

spillHashJoinTable 内部对每个 spiller 调用 spiller->spill()(即 HashBuildSpiller::spill()),将所有行按 hash 分区写到磁盘。

5.9 postHashBuildProcess / setupSpillInput — 递归恢复

Build 完成后,询问 join bridge 是否有 spill 数据需要 restore:

流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

void HashBuild::postHashBuildProcess() {
  if (!canSpill()) {
    setState(State::kFinish);
    return;
  }

  // 向 join bridge 请求下一个要恢复的 spill partition
  auto spillInput = joinBridge_->spillInputOrFuture(&future_);
  if (!spillInput.has_value()) {
    setState(State::kWaitForProbe); // 等待 probe 侧完成并传递 spill partition
    return;
  }
  setupSpillInput(std::move(spillInput.value()));
}

void HashBuild::setupSpillInput(HashJoinBridge::SpillInput spillInput) {
  if (spillInput.spillPartition == nullptr) {
    setState(State::kFinish); // 没有更多 spill partition,结束
    return;
  }

  // 重置所有状态,为本次 restore 重新初始化
  table_.reset(); spiller_.reset(); spillInputReader_.reset();
  restoringPartitionId_.reset();

  // 列顺序重置(spill 文件中列已按 key + dependent 排好)
  std::iota(keyChannels_.begin(), keyChannels_.end(), 0);
  std::iota(dependentChannels_.begin(), dependentChannels_.end(), keyChannels_.size());

  setupTable();
  setupSpiller(spillInput.spillPartition.get()); // 建立新的 spiller(可能是递归层)
  processSpillInput();                           // 开始读取并重建
}

5.10 processSpillInput — 从 spill 文件重建

流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

void HashBuild::processSpillInput() {
  while (spillInputReader_->nextBatch(spillInput_)) {
    // 复用 addInput 路径:重新 hash + 插入 hash 表(或再次 spill 到下一层)
    addInput(std::move(spillInput_));
    if (!isRunning()) return;
    if (shouldYield()) {
      state_ = State::kYield;
      future_ = ContinueFuture{folly::Unit{}};
      return;
    }
  }
  noMoreInputInternal(); // → finishHashBuild → postHashBuildProcess(循环)
}

这形成了递归循环:processSpillInput → addInput → ensureInputFits → [spill] → noMoreInput → finishHashBuild → postHashBuildProcess → setupSpillInput → processSpillInput。

5.11 Probe 侧 Spill 详解

Probe 侧(HashProbe)的 spill 涉及三条路径,通过 HashJoinBridge 与 build 侧协调:

HashJoinBridge 的交接事件
图 6:HashJoinBridge 的交接事件。已按当前实现修正标注,具体约束见相邻正文。
Fig. HashJoinBridge:setHashTable / probeFinished / append… 三个方法协调 build 与 probe 的 spill

5.11.1 HashBuildResult — Build/Probe 信息传递载体

流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

struct HashBuildResult {
  std::shared_ptr<BaseHashTable> table;                  // 刚建好的 hash 表
  std::optional<SpillPartitionId> restoredPartitionId;   // 非 null 表示本次是 restore 轮
  SpillPartitionIdSet spillPartitionIds;                 // 本轮 build 中被 spill 的 partition ID
  bool hasNullKeys;
};

spillPartitionIds 非空 ⟺ build 侧触发了 spill,probe 侧必须同步将对应行 spill 到磁盘。 restoredPartitionId 非 null ⟺ 本次 hash 表是从 spill 文件 restore 得来的,probe 侧需要从对应的 spill 文件中读取之前被 spill 的 probe 行。


5.11.2 路径 1:input spill(探测侧主路径)

触发时机:asyncWaitForHashTable 拿到 HashBuildResult 后,发现 spillPartitionIds 非空。

流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

void HashProbe::asyncWaitForHashTable() {
  auto hashBuildResult = joinBridge_->tableOrFuture(&future_);
  table_ = hashBuildResult->table;
  initializeResultIter();

  // 1. 若本次是 restore 轮,打开对应 partition 的 probe spill 文件
  maybeSetupSpillInputReader(hashBuildResult->restoredPartitionId);

  // 2. 若 build 侧有 spill partition,初始化 inputSpiller_
  maybeSetupInputSpiller(hashBuildResult->spillPartitionIds);
  checkMaxSpillLevel(hashBuildResult->restoredPartitionId);
}

**maybeSetupInputSpiller**:

void HashProbe::maybeSetupInputSpiller(
    const SpillPartitionIdSet& spillPartitionIds) {

  spillInputPartitionIds_ = spillPartitionIds;
  if (spillInputPartitionIds_.empty()) return;

  // 计算 bit offset(与 build 侧一致,保证 hash 分区对齐)
  const auto bitOffset = partitionBitOffset(
      *spillInputPartitionIds_.begin(),
      spillConfig()->startPartitionBit,
      spillConfig()->numPartitionBits);

  // NoRowContainerSpiller:无 RowContainer,直接将外部 RowVector 写到磁盘
  inputSpiller_ = std::make_unique<NoRowContainerSpiller>(
      probeType_,
      restoringPartitionId_,  // 递归时携带父 partition ID
      HashBitRange(bitOffset, bitOffset + spillConfig()->numPartitionBits),
      spillConfig(),
      spillStats_.get());

  // 只 spill build 侧已 spill 的那些 partition(其余 partition 可以正常 probe)
  inputSpiller_->setPartitionsSpilled(spillInputPartitionIds_);

  // SpillPartitionFunction:计算每行属于哪个 partition(与 build 侧使用相同的 hash bit 范围)
  spillPartitionFunction_ = std::make_unique<SpillPartitionFunction>(
      SpillPartitionIdLookup(spillInputPartitionIds_, ...),
      probeType_, keyChannels_);
}

**spillInput**(在 addInput 中调用):

void HashProbe::spillInput(RowVectorPtr& input) {
  const auto numInputRows = input->size();
  // 预分配各 partition 的 indices buffer(复用避免重复 malloc)
  prepareInputIndicesBuffers(numInputRows,
      inputSpiller_->state().spilledPartitionIdSet());

  // 为每行计算所属 partition
  spillPartitionFunction_->partition(*input, spillPartitions_);

  vector_size_t numNonSpillingInput = 0;
  for (auto row = 0; row < numInputRows; ++row) {
    const auto& partitionId = spillPartitions_[row];
    if (!inputSpiller_->state().isPartitionSpilled(partitionId)) {
      // 该 partition 未被 build 侧 spill,行留在内存继续 probe
      rawNonSpillInputIndicesBuffer_[numNonSpillingInput++] = row;
    } else {
      // 该 partition 已被 build 侧 spill,行需要 spill 到磁盘
      rawSpillInputIndicesBuffers_.at(partitionId)
          [numSpillInputs_.at(partitionId)++] = row;
    }
  }

  if (numNonSpillingInput == numInputRows) return; // 无行需要 spill

  // 强制加载所有 lazy 列(spill 序列化需要实际数据)
  for (int32_t i = 0; i < input->childrenSize(); ++i) {
    input->childAt(i)->loadedVector();
  }

  // 按 partition 分批 spill(wrap 创建 DictionaryVector,不拷贝数据,只共享 indices)
  for (const auto& [partitionId, numRows] : numSpillInputs_) {
    if (numRows == 0) continue;
    inputSpiller_->spill(
        partitionId,
        wrap(numRows, spillInputIndicesBuffers_.at(partitionId), input));
  }

  if (numNonSpillingInput == 0) {
    input = nullptr;     // 所有行都 spill 了,无需继续 probe
  } else {
    input = wrap(numNonSpillingInput, nonSpillInputIndicesBuffer_, input);
    // 只保留未 spill 的行继续 probe
  }
}

wrap 用 DictionaryVector 与 indices 表示属于某个分区的行,避免在分区步骤复制整行 payload;base vector 仍需保活。Lazy 值的加载时机、提取、序列化和文件缓冲各有路径,最终写文件仍会读取并编码这些值,因此应称“分区视图复用”,而不是从输入到磁盘全程零拷贝。

**noMoreInputInternal**(探测完成,收尾 input spill):

void HashProbe::noMoreInputInternal() {
  noMoreSpillInput_ = true;
  if (!spillInputPartitionIds_.empty()) {
    // 关闭 inputSpiller_,将所有分区的 SpillFiles 收集到 inputSpillPartitionSet_
    inputSpiller_->finishSpill(inputSpillPartitionSet_);
    // 注:NoRowContainerSpiller 不排序,spillSortTimeNanos 永远为 0
  }

  // 等待所有 peer probe 线程完成(防止一个线程先收尾而另一个还在写同一 partition)
  if (!operatorCtx_->task()->allPeersFinished(...)) {
    setState(ProbeOperatorState::kWaitForPeers);
    return;
  }

  lastProber_ = true; // 最后一个完成的 probe 线程负责通知 build 侧
}

5.11.3 路径 2:restore 轮的 probe(从 spill 文件读取 probe 行)

当 build 侧从某个 spill partition 重建 hash 表后,probe 侧必须把之前 spill 到磁盘的对应 probe 行读回重新 probe。

**maybeSetupSpillInputReader**:

流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

void HashProbe::maybeSetupSpillInputReader(
    const std::optional<SpillPartitionId>& restoredPartitionId) {
  if (!restoredPartitionId.has_value()) return;

  // 从 inputSpillPartitionSet_ 中取出对应 partition 的 SpillPartition
  auto iter = inputSpillPartitionSet_.find(restoredPartitionId.value());
  auto partition = std::move(iter->second);
  restoringPartitionId_ = restoredPartitionId;

  // 创建无序读取器(probe 行之间无顺序要求,顺序读即可)
  spillInputReader_ = partition->createUnorderedReader(
      spillConfig_->readBufferSize, pool(), spillStats_.get());
  inputSpillPartitionSet_.erase(iter);
}

**addSpillInput**(在 isBlocked 的 kRunning 状态时轮询调用):

void HashProbe::addSpillInput() {
  if (input_ != nullptr || noMoreSpillInput_) return;

  if (!spillInputReader_->nextBatch(input_)) {
    // spill 文件读完,走正常的 noMoreInput 流程
    noMoreInputInternal();
    return;
  }

  // 复用 addInput 路径(包含 spillInput 检查,支持递归 spill)
  addInput(std::move(input_));
}

5.11.4 路径 3:output spill(内存仲裁触发 reclaim)

5.12 Probe 阶段回收:只存剩余输入还不够

HashProbe 已经开始输出时,reclaim 必须考虑当前输入中尚未排出的 join 结果。HashProbe::reclaim 协调 peers,先 spill 剩余输出,再视是否仍有 probe 输入决定是否 spill 表,最后清 buffers、设置后续输入路由并登记分区。

right/full 等语义需要 probed flag;有些行此前已匹配,即使恢复后不再重复输出那些匹配,也必须防止它们最后被误当作未匹配 build 行。这个标志属于语义状态,不是可随意丢弃的缓存。

Build 阶段、bridge 已发布但 probe 未开始的窗口、probe 阶段各有不同 reclaim 入口。cached hash table、counting join 和部分执行模式又限制 canSpill。完整状态协作见 HashJoin。

触发条件:probe 正在处理某批 input 时,输出过大导致内存仲裁器调用 reclaim()。

**reclaim**:

流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

void HashProbe::reclaim(uint64_t, memory::MemoryReclaimer::Stats& stats) {
  if (exceededMaxSpillLevelLimit_ || nonReclaimableState()) {
    ++stats.numNonReclaimableAttempts; return;
  }

  const auto probeOps = findPeerOperators();
  bool hasMoreProbeInput = false;
  for (auto* probeOp : probeOps) {
    if (probeOp->nonReclaimableState()) { ++stats.numNonReclaimableAttempts; return; }
    hasMoreProbeInput |= !probeOp->noMoreSpillInput_;
  }

  // 步骤 1:并发将所有 peer 的 pending input 产出的 output 先 spill 到磁盘
  spillOutput(probeOps);

  // 步骤 2:若还有 probe input 未处理,spill hash 表(释放 build 侧内存)
  SpillPartitionSet spillPartitionSet;
  if (hasMoreProbeInput) {
    spillPartitionSet = spillHashJoinTable(
        table_, restoringPartitionId_, tableSpillHashBits_,
        joinNode_, spillConfig(), spillStats_.get());
  }

  // 步骤 3:通知各 probe 算子设置 inputSpiller_(后续输入按新 partitions spill)
  const auto spillPartitionIdSet = toSpillPartitionIdSet(spillPartitionSet);
  for (auto* probeOp : probeOps) {
    probeOp->clearBuffers();
    if (!spillPartitionSet.empty()) {
      probeOp->maybeSetupInputSpiller(spillPartitionIdSet);
    }
    probeOp->pool()->release();
  }

  // 步骤 4:清空 hash 表 + 将 spill partitions 注册到 join bridge(供 build 侧后续恢复)
  table_->clear(true);
  if (!spillPartitionIdSet.empty()) {
    joinBridge_->appendSpilledHashTablePartitions(std::move(spillPartitionSet));
  }
}

spillOutput(单个算子):

void HashProbe::spillOutput() {
  if (input_ == nullptr && !needLastProbe()) return;

  // 用 NoRowContainerSpiller 将 pending input 产生的所有 output 写到磁盘(单一 partition 0)
  auto outputSpiller = std::make_unique<NoRowContainerSpiller>(
      outputType_, std::nullopt, HashBitRange{}, spillConfig(), spillStats_.get());
  outputSpiller->setPartitionsSpilled({SpillPartitionId(0)});

  // 调用 getOutputInternal(toSpillOutput=true):产出 output 但直接 spill 而非返回
  for (;;) {
    auto output = getOutputInternal(/*toSpillOutput=*/true);
    if (output != nullptr) {
      for (int32_t i = 0; i < output->childrenSize(); ++i) {
        output->childAt(i)->loadedVector(); // 强制物化 lazy 列
      }
      outputSpiller->spill(SpillPartitionId(0), output);
      continue;
    }
    if (input_ == nullptr) break; // 当前 input 批次已处理完
    // right semi join 特殊:input_ 不为 null 但 output 为 null 属于正常状态
    break;
  }

  // 收尾:将 spill 文件信息存入 spillOutputPartitionSet_(供后续读回)
  outputSpiller->finishSpill(spillOutputPartitionSet_);
  removeEmptyPartitions(spillOutputPartitionSet_);
}

读回阶段:下次 getOutput 调用时,maybeReadSpillOutput() 将磁盘数据读回:

bool HashProbe::maybeReadSpillOutput() {
  if (spillOutputReader_ == nullptr) {
    maybeSetupSpillOutputReader(); // 首次:打开 UnorderedStreamReader
  }
  if (spillOutputReader_ == nullptr) return false;

  if (!spillOutputReader_->nextBatch(output_)) {
    spillOutputReader_.reset(); // 读完,重置
    return false;
  }
  return true; // 将读回的 output 直接返回给调用方
}

output spill 特殊注意:getOutputInternal(toSpillOutput=true) 被调用时,若发现 input_ 为 null(即没有 pending input),会提前返回 null——这保证了 reclaim 路径不会误进入"probe 完成逻辑"(如 prepareForSpillRestore),避免状态混乱。


5.12.1 prepareForSpillRestore — 多轮恢复的关键

当一轮 probe 完成时,peer barrier 先保证必要的探测工作已经结束;协调者再通过 prepareForSpillRestore 与 bridge 推进下一轮。build 侧补行可以根据配置并行输出,因此不能把“最后一个 probe 负责协调”扩大为“所有 right/full build 侧输出永远由它单线程完成”。

流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

void HashProbe::prepareForSpillRestore() {
  // 重置本轮相关状态
  noMoreSpillInput_ = false;
  if (lastProber_) {
    table_->clear(true); // 清空当前 hash 表,释放内存
  }
  table_.reset();
  inputSpiller_.reset();
  spillInputReader_.reset();
  restoringPartitionId_.reset();
  spillInputPartitionIds_.clear();
  spillOutputReader_.reset();

  if (!lastProber_) return;

  // 通知 join bridge:probe 侧已完成,请准备下一个 spill partition
  joinBridge_->probeFinished();
  // bridge 内部:从 spillPartitionSets_ 取下一个 partition,均分为 N shard,
  // 放入 restoringSpillShards_,build 侧 spillInputOrFuture() 返回各自的 shard

  // 唤醒等待的 peer probe 线程(之前在 kWaitForPeers 状态阻塞)
  wakeupPeerOperators();
  lastProber_ = false;
}

5.12.2 Probe 侧 Spill 全生命周期

asyncWaitForHashTable()
    │
    ├─ spillPartitionIds 非空 → maybeSetupInputSpiller()
    │      创建 NoRowContainerSpiller,标记需要 spill 的 partition 集合
    │
    ├─ restoredPartitionId 非 null → maybeSetupSpillInputReader()
    │      打开 UnorderedStreamReader,顺序读取之前 spill 的 probe 行
    │
addInput(input)
    │
    └─ needToSpillInput()? → spillInput(input)
           按 hash 分 partition:
             spilled partition → DictionaryVector wrap → inputSpiller_.spill()
             non-spilled partition → 继续正常 probe(无拷贝)
    │
getOutput() → ensureOutputFits()
    │
    ├─ maybeReadSpillOutput():读取上次 spillOutput 写到磁盘的 output batch
    │
    └─ 内存不足 → reclaim():
           1. spillOutput(probeOps)        并发 spill 各 probe 的 pending output
           2. spillHashJoinTable(table)     将 hash 表 spill 到磁盘
           3. maybeSetupInputSpiller(ids)   各 probe 设置新 inputSpiller_
           4. appendSpilledHashTablePartitions() → 通知 bridge
    │
isBlocked(kRunning) + spillInputReader_ 非 null → addSpillInput()
    │    逐 batch 从 spill 文件读取 probe 行 → addInput() → spillInput() → probe
    │
noMoreInputInternal()
    │
    ├─ inputSpiller_.finishSpill(inputSpillPartitionSet_)
    └─ allPeersFinished() → lastProber_ = true
             │
             ▼
        hasMoreSpillData()? → prepareForSpillRestore()
             │                    joinBridge_->probeFinished()
             │                    wakeupPeerOperators()
             ▼
        asyncWaitForHashTable()  ← 循环(等待 build 侧重建 hash 表)
             │
             └─ 所有 spill partition 处理完 → setState(kFinish)

6. OrderBy Spill 实现

6.1 两种 Spiller

6.2 OrderBy:输入 run 与剩余输出

SortBuffer 先把输入保存在 RowContainer。无 spill 时,noMoreInput 后排序行指针;发生 input spill 时,把内存中的行排成有序 run 写出并清空,输入结束后还要写出剩余状态,最后归并。参见 noMoreInput、spillInput。

SortOutputSpiller 处理另一种情况:已经在内存排好序、也可能输出了一部分,但剩余结果需要回收。它从 sortedRows_[numOutputRows_] 开始写剩余指针,不再重排,也不能把已输出前缀再次写出。参见 spillOutput。

排序规则包含升降序和 null 顺序,spill 与 merge 必须使用一致规则。对于相等键,不能从用了某种排序算法就推导整个分布式查询提供未声明的稳定顺序保证。

OrderBy 的 spill 通过 SortBuffer 实现,同样区分两个阶段:

inputSpiller_  (SortInputSpiller):
  - input 阶段:RowContainer 行数过多时触发
  - needSort = true:按 sort key 排序后写入
  - 无 hash 分区(HashBitRange 为空)
  - 写入一个 partition(partition 0)

outputSpiller_ (SortOutputSpiller):
  - output 阶段:已完成内排序,输出时内存不足
  - needSort = false:外部传入已排序的 SpillRows
  - 接收 sortedRows_ 中未输出部分,直接写入

6.3 SortBuffer::spillInput

流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

void SortBuffer::spillInput() {
  if (inputSpiller_ == nullptr) {
    // sortingKeys:将所有 sort column + 其余列按 sort 顺序排列
    const auto sortingKeys = SpillState::makeSortingKeys(sortCompareFlags_);
    inputSpiller_ = std::make_unique<SortInputSpiller>(
        data_.get(),         // RowContainer
        spillerStoreType_,   // sort 列在前,其余列在后的类型
        sortingKeys,
        spillConfig_,
        spillStats_);
  }
  inputSpiller_->spill();  // fillSpillRuns → sort → writeSpill → 写磁盘
  data_->clear();           // 清空 RowContainer,腾出内存
}

spillerStoreType_ 的列顺序重排(构造时完成):

// sort 列排在最前面(归并时只需比较前几列)
for (auto i = 0; i < numSortKeys; ++i) {
  sorted_types.push_back(input->childAt(sortColumnIndices[i])->type());
  sorted_names.push_back(input->nameOf(sortColumnIndices[i]));
}
// 其余列追加在后面
for (auto i = 0; i < input->size(); ++i) {
  if (!isSortKey(i)) {
    sorted_types.push_back(input->childAt(i)->type());
    sorted_names.push_back(input->nameOf(i));
  }
}

SortBuffer 将排序键调整到 spill schema 的前部,便于排序与归并比较,并通过 columnMap_ 恢复输出列顺序。这种布局减少了比较路径的列定位工作,但不表示其他 payload 列可以省略反序列化。

6.4 SortBuffer::spillOutput

流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

void SortBuffer::spillOutput() {
  if (hasSpilled()) return;                     // 已经 spill 过
  if (numOutputRows_ == sortedRows_.size()) return; // 全部输出完了

  outputSpiller_ = std::make_unique<SortOutputSpiller>(
      data_.get(), spillerStoreType_, spillConfig_, spillStats_);

  // 将 sortedRows_ 中未输出的部分([numOutputRows_, end))直接传给 outputSpiller_
  auto spillRows = SpillerBase::SpillRows(
      sortedRows_.begin() + numOutputRows_, sortedRows_.end(),
      *memory::spillMemoryPool());
  outputSpiller_->spill(spillRows);

  data_->clear();
  sortedRows_.clear(); sortedRows_.shrink_to_fit();
  finishSpill(); // output spill 只触发一次,立即收尾
}

SortOutputSpiller::spill(SpillRows) 内部:

void SortOutputSpiller::spill(SpillRows& rows) {
  auto& spillRun = createOrGetSpillRun(SpillPartitionId(0));
  spillRun.rows = SpillRows(rows.begin(), rows.end(), ...);
  // 计算 numBytes
  for (const auto* row : rows) spillRun.numBytes += container_->rowSize(row);
  markSeenPartitionsSpilled();

  runSpill(true); // rows 已按 sort key 有序,跳过排序直接写入
}

SortOutputSpiller::runSpill 重写:在父类 runSpill 后立即调用 finishFile(保证文件封尾)。

6.5 SortBuffer::finishSpill

流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

void SortBuffer::finishSpill() {
  // inputSpiller_ 和 outputSpiller_ 互斥,只有一个会被初始化
  if (inputSpiller_ != nullptr) {
    inputSpiller_->finishSpill(spillPartitionSet_);
  } else {
    outputSpiller_->finishSpill(spillPartitionSet_);
  }
  // OrderBy 无 hash 分区,只有一个 partition
  VELOX_CHECK_EQ(spillPartitionSet_.size(), 1);
}

6.6 SortBuffer::getOutputWithSpill

流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

void SortBuffer::prepareOutputWithSpill() {
  if (spillMerger_ != nullptr) return; // 已初始化

  VELOX_CHECK_EQ(spillPartitionSet_.size(), 1);
  // 创建 TreeOfLosers:归并所有 spill 文件(可能包含内存中剩余行)
  spillMerger_ = spillPartitionSet_.begin()->second->createOrderedReader(
      *spillConfig_, pool(), spillStats_);
  spillPartitionSet_.clear();
}

void SortBuffer::getOutputWithSpill() {
  int32_t outputRow = 0, outputSize = 0;
  while (outputRow + outputSize < output_->size()) {
    SpillMergeStream* stream = spillMerger_->next();
    VELOX_CHECK_NOT_NULL(stream);

    spillSources_[outputSize] = stream->current().get();
    spillSourceRows_[outputSize] = stream->currentIndex(&isEndOfBatch);
    ++outputSize;

    if (isEndOfBatch) {
      // batch 到边界时先复制出来,再 pop 触发加载下一 batch
      gatherCopy(output_.get(), outputRow, outputSize,
                 spillSources_, spillSourceRows_, columnMap_);
      outputRow += outputSize; outputSize = 0;
    }
    stream->pop();
  }
  // 处理最后一批
  if (outputSize != 0) {
    gatherCopy(output_.get(), outputRow, outputSize, ...);
  }
  numOutputRows_ += output_->size();
}

gatherCopy 使用 columnMap_ 将 spill 文件中的列顺序(sort 列在前)映射回原始输出列顺序。


7. Spiller 层详解

7.1 SpillerBase — 抽象基类

文件:velox/exec/Spiller.h:29, Spiller.cpp:31

流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

class SpillerBase {
 protected:
  RowContainer* const container_;    // 持有行数据的容器(nullptr 表示无容器)
  folly::Executor* const executor_;  // 非 null 时并发写 spill
  const HashBitRange bits_;          // hash 分区 bit 范围
  const RowTypePtr rowType_;
  const uint64_t maxSpillRunRows_;   // 单 run 最大行数
  const std::optional<SpillPartitionId> parentId_; // 递归 spill 时的父 ID

  bool finalized_{false};            // finishSpill 后置 true,不允许再写
  SpillState state_;                 // 管理所有 partition writer
  folly::F14FastMap<SpillPartitionId, SpillRun> spillRuns_; // 内存中的行缓冲

  // 纯虚方法(子类决定行为差异)
  virtual bool needSort() const = 0;         // 是否在写前排序
  virtual std::string type() const = 0;      // 类型名称(日志用)
  virtual void extractSpill(folly::Range<char**>, RowVectorPtr&); // 从 container 提取
  virtual void runSpill(bool lastRun);        // 执行一次 spill run 的写入
};

7.2 SpillRun — 内存中的 spill 缓冲

流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

struct SpillRun {
  SpillRows rows;         // char* 指针数组,指向 RowContainer 内部的行
  uint64_t numBytes{0};   // 所有行的估算总字节数(通过 rowSize() 累加)
  bool sorted{false};     // 是否已排序(避免同一 run 被重复排序)

  void clear() {
    rows.clear();
    rows.shrink_to_fit();  // 主动归还内存,spill 后立即回收指针数组
    numBytes = 0;
    sorted = false;
  }
};

SpillRun::rows 存的是 RowContainer 内部的原始行指针(char*),不拷贝行数据本身,仅持有指针。真正的行数据仍在 RowContainer 的 Arena 内存中,直到 extractSpill 时才物化为列式 RowVector。每个分区(SpillPartitionId)对应一个 SpillRun,存在 spillRuns_ map 中。


7.3 fillSpillRuns — 数据填充

这是 spill 流程的第一步,负责扫描 RowContainer、计算 hash、将行指针按分区放入对应 SpillRun。

流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

bool SpillerBase::fillSpillRuns(RowContainerIterator* iterator) {
  constexpr int32_t kHashBatchSize = 4096; // 每次从 RowContainer 取多少行

  bool lastRun{false};
  uint64_t totalRows{0};

  for (;;) {
    // 1. 从 RowContainer 按迭代器顺序取下一批行指针(迭代器原地推进)
    const auto numRows =
        container_->listRows(iterator, kHashBatchSize, rows.data());
    if (numRows == 0) { lastRun = true; break; }  // 全部取完,这是最后一轮

    auto rowSet = folly::Range<char**>(rows.data(), numRows);

    // 2. 计算 hash(逐列迭代,第 2 列起 combine=true,追加 XOR 到已有 hash)
    if (!isSinglePartition) {
      for (auto i = 0; i < container_->keyTypes().size(); ++i) {
        container_->hash(i, rowSet, /*combine=*/i > 0, hashes.data());
      }
    }

    // 3. 按 hash 将行指针分配到对应分区的 SpillRun
    for (auto i = 0; i < numRows; ++i) {
      const auto partitionNum =
          isSinglePartition ? 0 : bits_.partition(hashes[i]);
      auto& spillRun = createOrGetSpillRun(SpillPartitionId(partitionNum));
      spillRun.rows.push_back(rows[i]);
      spillRun.numBytes += container_->rowSize(rows[i]); // 估算行字节数
    }

    totalRows += numRows;
    // 4. 达到 maxSpillRunRows_ 上限时暂停,进入 runSpill,之后再继续
    if (maxSpillRunRows_ > 0 && totalRows >= maxSpillRunRows_) break;
  }

  markSeenPartitionsSpilled(); // 标记本次出现的所有 partition 为已 spill
  return lastRun;
}

hash() 对第一列以 combine=false 初始化 hashes,后续列以 combine=true 把该列 hash 混合进已有结果;具体混合由 RowContainer 的类型化 hash 路径决定,并不是简单地把所有列的 hash 做裸 XOR。每列仍需遍历对应的行指针,批量接口减少了逐行上层调用。


7.4 SpillRun 为何切分 / maxSpillRunRows 的作用

7.5 SpillRun、batch 和 file 是三个尺度

fillSpillRuns 以最多 4096 行为一批从 RowContainer 取行指针,计算 key hash,把行指针放到各分区的 SpillRun 中。SpillRun 保存 char* 指针数组和估算字节数,不立即复制出完整列向量。

maxSpillRunRows 限制本轮收集的总行数,不是“每个分区各自有这么多行”的硬上限。检查在一批候选处理后发生,因此可以越过目标行数一小批;0 表示此项不限制。

图 4:行指针 run、序列化 batch 和磁盘文件的边界不同,不能用同一个大小参数解释。
图 7:行指针 run、序列化 batch 和磁盘文件的边界不同,不能用同一个大小参数解释。

writeSpill 的默认提取目标为 最多 64 行、约 256 KiB。extractSpillVector 会包含让字节数越过目标的那一行,保证至少能写一行;因此 256 KiB 不是单行或 batch 的严格最大值,也不是压缩后磁盘字节上限。

文件边界由 writer 与目标文件大小管理;有序 run 结束后会主动结束当前文件,保证后续新 run 不被误当成前一 run 的有序延续。

spill() 的主循环是 do { fillSpillRuns; runSpill; } while (!lastRun)。maxSpillRunRows_ 控制每轮 fillSpillRuns 最多处理多少行,从而将整个大 RowContainer 的 spill 分成多个小 run 依次执行。这样做有四个原因:

1. 指针数组本身的内存代价

SpillRun::rows 是 std::vector<char*>,每个元素 8 字节。若 RowContainer 有 1 亿行,不切分时仅指针数组就需要 800 MB。maxSpillRunRows_ 通常设为数十万行,指针数组只占几 MB。

2. 排序的内存代价(仅 sorted spill)

切小 run 限制参与排序的指针数组和辅助空间,可能改善工作集局部性。行指针仍会间接访问 RowContainer,宽行、长字符串以及多个并行分区都可能超出 cache;不能保证每轮排序完全在缓存中完成。

3. 分步释放,配合内存仲裁

每轮写入成功后,SpillRun::clear 清理本轮行指针数组和相关统计;这会减少指针缓冲的占用,但不代表对应 RowContainer 的源行已经逐批释放。真正的行状态由调用算子在写出完成后 clear,释放了 used/reservation 后,仲裁层仍需 shrink 才归还 root capacity。Spiller 的循环不会仅凭“本轮达到回收目标”自动中止尚未完成的全量 spill。

4. PrestoPage 序列化限制

Presto 序列化页存在相应的长度表示和校验限制,但控制 run 的行指针数量不能单独保证页字节大小。run、extractSpillVector 提取 batch、writer 序列化缓冲以及磁盘文件各有自己的边界;超大单行仍可能触发大小限制,应分别检查。

maxSpillRunRows_ 的取值:
  HashBuildSpiller          → spillConfig->maxSpillRunRows(由当前 spillConfig 与宿主配置决定)
  AggregationInputSpiller   → spillConfig->maxSpillRunRows
  AggregationOutputSpiller  → spillConfig->maxSpillRunRows
  SortInputSpiller          → spillConfig->maxSpillRunRows
  SortOutputSpiller         → spillConfig->maxSpillRunRows
  NoRowContainerSpiller     → 0(无上限,直接 appendToPartition,无 RowContainer)

7.6 行列转换:extractSpill 与 extractSpillVector

7.7 RowContainer 到列向量:为何要提取而不是 memcpy

RowContainer 存的是执行状态,可能包含 key、dependent、null 标志、变长数据引用和 aggregate accumulator。它不是可直接持久化的 RowVector 内存镜像。

extractSpill 按列调用 RowContainer::extractColumn,把行状态恢复成列向量;accumulator 则通过 extractForSpill 输出其专用 spill 表示。临时 RowVector 来自 spillMemoryPool,循环中可以 prepareForReuse 后复用。

例如 avg 的状态需要保留足够信息恢复 sum/count;不能写出最终 avg 数值再把多个 avg 当作原始输入合并。更复杂 aggregate 的 spillType 与提取方法必须彼此匹配。

HashBuildSpiller 还有专门的 probed flag 提取。对于需要输出 build 侧匹配/未匹配行的 join,恢复后必须知道哪些行已经在此前被 probe 过。参见 HashBuildSpiller::extractSpill。

SpillerBase::runSpill 清理的是本轮指针数组;底层 RowContainer 的状态由调用者在成功写出后 clear/释放。不要依据旧注释写成“每写一个 batch 就立刻删除对应源行”。

spill 的核心转换:行式 RowContainer → 列式 RowVector。这发生在 writeSpill → extractSpillVector → extractSpill 调用链中。

7.7.1 extractSpillVector — 控制每批大小

流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

int64_t SpillerBase::extractSpillVector(
    SpillRows& rows,
    int32_t maxRows,      // 上限:64 行
    int64_t maxBytes,     // 上限:256 KB(估算)
    RowVectorPtr& spillVector,
    size_t& nextBatchIndex) {   // 游标,指向 rows[] 中下一个未处理的行

  // 1. 确定本批行数:不超过 maxRows,且估算总字节不超过 maxBytes
  int32_t numRows = 0;
  int64_t bytes = 0;
  auto limit = std::min(rows.size() - nextBatchIndex, (size_t)maxRows);

  for (; numRows < limit; ++numRows) {
    bytes += container_->rowSize(rows[nextBatchIndex + numRows]);
    if (bytes > maxBytes) {
      ++numRows;  // 超了也要包含这一行(至少保证提取 1 行,避免死循环)
      break;
    }
  }

  // 2. 真正执行行 → 列转换
  extractSpill(folly::Range(&rows[nextBatchIndex], numRows), spillVector);
  nextBatchIndex += numRows;
  return bytes;
}

rowSize() 的计算:

uint32_t RowContainer::rowSize(const char* row) const {
  return fixedRowSize_ +
      (rowSizeOffset_
           ? *reinterpret_cast<const uint32_t*>(row + rowSizeOffset_)
           : 0);
}

fixedRowSize_ 是 RowContainer 根据运行时 schema、字段宽度、对齐和 accumulator 布局计算的固定行部分。rowSizeOffset_ 对应的计数字段参与估算变长部分;这个估算用于 batch/内存策略,不应等同于序列化后的精确字节数。

7.7.2 extractSpill — 行式 → 列式核心转换

流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

void SpillerBase::extractSpill(
    folly::Range<char**> rows,   // 本批 char* 指针
    RowVectorPtr& resultPtr) {

  // 1. 分配或复用 RowVector
  if (resultPtr == nullptr) {
    resultPtr = BaseVector::create<RowVector>(
        rowType_, rows.size(), memory::spillMemoryPool());
  } else {
    resultPtr->prepareForReuse();  // 复用已有分配,避免重复 malloc
    resultPtr->resize(rows.size());
  }

  auto* result = resultPtr.get();

  // 2. 提取各普通列(key + dependent 列)
  for (auto i = 0; i < container_->columnTypes().size(); ++i) {
    container_->extractColumn(rows.data(), rows.size(), i, result->childAt(i));
  }

  // 3. 提取各 accumulator 列(聚合中间状态)
  const auto& accumulators = container_->accumulators();
  for (auto i = 0; i < accumulators.size(); ++i) {
    accumulators[i].extractForSpill(
        rows, result->childAt(i + container_->columnTypes().size()));
  }
}

spillMemoryPool() 是专用的 spill 内存池,与算子的 operator pool 独立,防止 spill 写入过程中分配的临时内存被算子内存统计。


7.8 RowContainer 行布局与列提取原理

理解行列转换,需要先了解 RowContainer 的行内存布局:

RowContainer 与 spill 提取
图 8:RowContainer 与 spill 提取。已按当前实现修正标注,具体约束见相邻正文。
Fig. RowContainer 行布局:normalizedKey 前缀 → null bits → 列值 → accumulator → 元信息

extractColumn(rows[], numRows, columnIndex, result) 的工作:

1. 确定该列在行内的 byte offset(通过 RowColumn 对象,同一运行时 schema 布局内固定)
2. 读取 null bit 填充 result 的 null bitmap
3. 对每行读取 offset 处的值:
   - 固定宽度列:直接 memcpy 4/8 字节
   - StringView 列:读取 ptr+size,若数据是 inline(≤12B)直接复制;
     否则从 HashStringAllocator 的堆上读取实际字符串,写入 result 的 string buffer

关键性能特征:

  • 同一 RowContainer 内,同一列的 offset 对所有行相同;布局在创建容器时依据运行时 schema 确定,不是所有查询共享一个编译期常量。列提取可以复用该 offset,但仍需按类型、null 与变长数据表示读取。
  • 变长列(StringView)在行中只存引用,实际字符串在 HashStringAllocator 堆上;提取时需要一次间接访问,可能有 cache miss
  • freezeAndExecute 约束 spill 提取期间对 HashStringAllocator 的可变操作,避免并行提取时继续通过这个非线程安全 allocator 分配或释放。它不是一把会自动串行化任意访问的通用锁;Task pause、reclaim 协议和 aggregate 的提取实现仍承担安全责任。

7.9 Accumulator 中间态提取

对于聚合算子,extractForSpill 提取的不是最终聚合结果,而是中间状态(partial intermediate):

流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

void Accumulator::extractForSpill(folly::Range<char**> groups, VectorPtr& result) const {
  spillExtractFunction_(groups, result);
  // 实际调用:aggregate->extractAccumulators(groups.data(), groups.size(), &result)
}

extractAccumulators 是 Aggregate 虚函数,每种聚合函数自己实现。例如:

  • sum(int64) → 提取 int64 partial sum
  • avg(double) → 提取 (sum: double, count: int64) row struct
  • array_agg(T) → 提取序列化的数组(out-of-line 存储,需要 copy)

恢复(unspill)时,updateRow 调用 aggregate->addSingleGroupIntermediateResults,将磁盘读回的中间状态与内存中的行合并:

void GroupingSet::updateRow(SpillMergeStream& input, char* row) {
  mergeSelection_.setValid(input.currentIndex(), true);
  for (auto i = 0; i < aggregates_.size(); ++i) {
    mergeArgs_[0] = input.current()->childAt(i + keyChannels_.size());
    // 将 spill 文件中的中间状态合并到 row 的 accumulator
    aggregates_[i].function->addSingleGroupIntermediateResults(
        row, mergeSelection_, mergeArgs_, false);
  }
}

这实现了"流式归并聚合":相同 key 的行按排好的顺序依次到来,逐行 merge accumulator,避免将整个 group 缓存在内存中。


7.10 writeSpill 内的 Batch 切分(kTargetBatchRows / kTargetBatchBytes)

writeSpill 在一个 SpillRun 内继续按 batch 提取列向量。当前默认行数目标为 64,字节目标约 256 KiB;提取逻辑会纳入使估算字节数越过目标的那一行,并保证至少推进一行。因此字节目标不是硬上限,超大单行不能被假定为小于 256 KiB。

SpillRun.rows = [ptr0, ptr1, ..., ptr_N]   // N 可能有数十万
                    ↓
writeSpill 循环:
  batch 0: [ptr0  .. ptr63]   → extractSpill → RowVector(64行) → appendToPartition
  batch 1: [ptr64 .. ptr127]  → extractSpill → RowVector(64行) → appendToPartition
  ...
  batch K: [ptr_last_start .. ptr_N]        → extractSpill → RowVector(M行) → appendToPartition

为什么要切成 64 行/256KB 的小 batch?

① RowVector 内存峰值控制

每次 extractSpill 分配(或复用)一个 RowVector,其内存由 spillMemoryPool 分配。单批 64 行的 RowVector 内存占用可控(通常几 KB 到几百 KB),而 prepareForReuse() 使相邻批次复用同一块内存,峰值内存 ≈ 单批大小,而非整个 run 大小。

② 与 SpillWriter 写缓冲的配合

SpillWriter 将序列化内容累积到自己的写缓冲,再按文件接口触发实际写入。提取 batch 的目标与 writeBufferSize 是不同尺度,小 batch 可以帮助控制物化峰值,但不保证每次写入恰好贴近缓冲大小,更不等于一次 batch 对应一次系统调用。

③ "至少 1 行"的保证防止死循环

流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

if (bytes > maxBytes) {
  ++numRows;  // 即使超了也把这一行纳入
  break;
}

若某行本身 > 256 KB(如包含超大字符串),不加这一行就会无限循环(numRows=0 → written 不推进)。++numRows 确保每次循环至少推进一行。

④ 64 行是何依据?

64 是当前实现的行数目标。以单列 int64 为例,64 项数据本身为 512 B;但完整 RowVector 还包含其他列、null、变长值和对象开销,不能据此证明整个 batch 位于 L1。这个参数对提取、序列化和缓存的实际影响应按数据形态测量。

  • 单个含 64 项的 int64 数据缓冲为 512 B,但 RowVector 的完整工作集并不只有这个缓冲;列数、字符串和并行任务都会增加容量需求。
  • extractColumn 的内层循环(每列 64 次写入)具有良好的向量化潜力
  • SpillMergeStream 按 batch 补充数据,归并树仍需根据逐行选出的最小值维护状态。批大小影响读取和解码开销,不能直接推导为“每批才重平衡一次”。

7.11 runSpill — 并发写入

7.12 有序写出与异步任务的收尾

需要排序的子类通过 ensureSorted 对行指针排序,使用 TimSort 或配置的 PrefixSort;数据行本身不按排序顺序整体搬动。SpillRun::sorted 防止同一 run 被重复排序。参见 ensureSorted。

runSpill 为非空分区创建 AsyncSource 任务:

  1. 第一个任务留给当前线程通过 move() 执行/取得结果。
  2. 若有 executor,后续任务提交 prepare()。
  3. 当前线程逐个取得结果,检查异常并清理 run。
  4. guard 在异常退出时继续消费已发出的任务,避免悬空工作引用 spiller。

这是“当前线程也参加、返回前收齐结果”的并发策略,不是 fire-and-forget。executor 为空时也能在当前线程完成;配置 executor 不等于所有写入都转移到后台。

异步任务通过 exception_ptr 把文件/序列化错误带回调用者。只有确认写出成功之后,算子才能清掉待恢复的内存状态。release() 释放未用 reservation,不能替代清理仍存活的行。

流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

void SpillerBase::runSpill(bool lastRun) {
  spillStats_->spillRuns.fetch_add(1);

  std::vector<std::shared_ptr<AsyncSource<SpillStatus>>> writes;
  for (const auto& [id, spillRun] : spillRuns_) {
    if (spillRun.rows.empty()) continue;
    writes.push_back(
        memory::createAsyncMemoryReclaimTask<SpillStatus>(
            [partitionId = id, this]() { return writeSpill(partitionId); }));
    // 第 2 个及之后的任务提交到 executor(后台线程),第 1 个在当前线程执行
    if ((writes.size() > 1) && executor_ != nullptr) {
      executor_->add([source = writes.back()]() { source->prepare(); });
    }
  }

  // 折叠守卫:无论成功失败,确保所有后台任务都被 move()(drain)
  auto sync = folly::makeGuard([&]() {
    for (auto& write : writes) { try { write->move(); } catch (...) {} }
  });

  // 主线程阻塞等待所有任务(含第 1 个本线程任务)
  for (auto& write : writes) {
    results.push_back(write->move());
  }

  // 检查错误 + 清空 SpillRun(归还指针数组内存)
  for (auto& result : results) {
    if (result->error) std::rethrow_exception(result->error);
    spillRuns_.at(result->partitionId).clear();
    if (needSort()) {
      state_.finishFile(result->partitionId); // sorted spill:每 run 一个文件
    }
  }
}

并发策略细节:

  • 第 1 个分区的写任务在调用线程上通过 write->move() 触发(AsyncSource::prepare() 从未被调用,move() 内联执行)
  • 第 2+ 个分区提交到 executor_(后台 CPU 线程池)并发执行
  • createAsyncMemoryReclaimTask 把回收任务关联到内存回收执行上下文,以支持相应状态与线程语义。这个包装不等于任意 spill 写任务都可在中途安全取消;runSpill 仍必须取得已发出任务的结果并处理异常,才能释放被任务引用的状态。

state_.finishFile(partitionId) 调用 SpillWriter::finishFile(),flush 当前文件并记录 SpillFileInfo,下次 appendToPartition 会创建新文件。这使每个 sorted run 独立成文,后续归并时文件间已有序。


7.13 finishSpill — 收尾

流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

void SpillerBase::finishSpill(SpillPartitionSet& partitionSet) {
  finalizeSpill();  // finalized_ = true,之后不再允许写入

  for (const auto& partitionId : state_.spilledPartitionIdSet()) {
    // 递归 spill 时,要将本层 partition ID 包装进父 ID 形成完整路径
    auto wholeId = parentId_.has_value()
        ? SpillPartitionId(parentId_.value(), partitionId.partitionNumber())
        : partitionId;

    // state_.finish(partitionId) 调用 SpillWriter::finish(),
    // 关闭所有文件,返回 SpillFiles(文件路径、大小、排序键等元数据)
    if (partitionSet.count(wholeId) == 0) {
      partitionSet.emplace(wholeId,
          std::make_unique<SpillPartition>(wholeId, state_.finish(partitionId)));
    } else {
      // 多个 Spiller(多线程 HashBuild)的文件合并到同一 partition
      partitionSet[wholeId]->addFiles(state_.finish(partitionId));
    }
  }
}

partitionSet 是 std::map<SpillPartitionId, ...>,按 SpillPartitionId 有序排列(DFS 前序),确保递归 spill 场景下父 partition 先于子 partition 被处理。


7.14 各子类 Spiller 详解

7.14.1 HashBuildSpiller

needSort      : false(hash join 恢复时重建 hash 表,顺序无关)
HashBitRange  : [startPartitionBit, startPartitionBit + numPartitionBits)(分区存储)
targetFileSize: spillConfig->maxFileSize
maxSpillRunRows: spillConfig->maxSpillRunRows
parentId_     : restoringPartitionId_(递归时指定父 partition)

差异点——重写 extractSpill 附加 probed flag 列:

流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

void HashBuildSpiller::extractSpill(folly::Range<char**> rows, RowVectorPtr& resultPtr) {
  // 先提取普通列(key + dependent)
  for (auto i = 0; i < container_->columnTypes().size(); ++i) {
    container_->extractColumn(rows.data(), rows.size(), i, result->childAt(i));
  }
  // 若是 right/full join,额外提取每行的 probed 标志位
  if (spillProbeFlag_) {
    auto* flagVector = result->childAt(types.size())->asFlatVector<bool>();
    for (auto i = 0; i < rows.size(); ++i) {
      flagVector->set(i, container_->hasProbedFlag(rows[i]));
      // hasProbedFlag 直接读行内 probedFlagOffset_ 处的 bit
    }
  }
}

spill() 重载1(spillTriggered_=true + 全量 spill):

void HashBuildSpiller::spill() {
  spillTriggered_ = true;
  SpillerBase::spill(nullptr);  // fillSpillRuns(nullptr) → 从头扫全部行
}

spill(partitionId, spillVector) 重载2(直通路径,已 triggered 后):

void HashBuildSpiller::spill(const SpillPartitionId& id, const RowVectorPtr& vec) {
  // 无需经过 RowContainer,直接 appendToPartition
  if (!state_.isPartitionSpilled(id)) state_.setPartitionSpilled(id);
  state_.appendToPartition(id, vec);
}

7.14.2 AggregationInputSpiller

needSort      : true(按所有 group key 排序,支持流式归并聚合)
HashBitRange  : [startPartitionBit, startPartitionBit + numPartitionBits)
targetFileSize: uint64_t::max(文件大小不限,由 maxFileSize 配置控制)
maxSpillRunRows: spillConfig->maxSpillRunRows
sortingKeys   : 所有 key 列,升序,nulls last

继承默认 extractSpill:提取 key 列 + 各 accumulator 的 extractAccumulators 中间状态。

每次 runSpill 后(needSort=true),调用 state_.finishFile(id) 关闭当前文件——每个 run 的数据写到独立的 sorted 文件。后续 createOrderedReader 用败者树对所有文件做 N 路归并,还原全局有序流。

关键约束:rows->stringAllocator().freezeAndExecute([&]() { inputSpiller_->spill(); }) 在 spill 期间冻结 HashStringAllocator,防止 accumulator 在 spill 时并发分配/释放变长内存。


7.14.3 AggregationOutputSpiller

needSort      : false(hash 表 scan 顺序即为输出顺序,不需要重排)
HashBitRange  : {}(空,单 partition 0)
targetFileSize: uint64_t::max
maxSpillRunRows: spillConfig->maxSpillRunRows

spill(startRowIter) 从 rowIterator 指定的偏移开始扫描(支持从已输出的位置继续):

流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

void AggregationOutputSpiller::spill(const RowContainerIterator& startRowIter) {
  state_.setPartitionSpilled(SpillPartitionId(0)); // 单分区
  SpillerBase::spill(&startRowIter);  // 从 startRowIter 开始迭代
}

重写 runSpill:在最后一次 run 后强制 finishFile:

void AggregationOutputSpiller::runSpill(bool lastRun) {
  SpillerBase::runSpill(lastRun);
  if (lastRun) {
    for (const auto& [id, _] : spillRuns_) {
      state_.finishFile(id); // 确保最后的数据落盘
    }
  }
}

7.14.4 SortInputSpiller

needSort      : true(按 sort key 排序)
HashBitRange  : {}(空,单 partition 0,OrderBy 无 hash 分区)
targetFileSize: uint64_t::max
maxSpillRunRows: spillConfig->maxSpillRunRows
sortingKeys   : 所有 sort 列,对应 CompareFlags(方向 + nulls 位置)

spillerStoreType_ 把 sort 列排在最前面,这样归并时只需比较前 N 列(减少 cache miss)。

spill() 直接调用 SpillerBase::spill(nullptr),无参数——从 RowContainer 头部开始扫全部行。


7.14.5 SortOutputSpiller

needSort      : false(调用方已传入排好序的 SpillRows)
HashBitRange  : {}(单 partition 0)
targetFileSize: uint64_t::max
maxSpillRunRows: spillConfig->maxSpillRunRows

接口与其他 spiller 不同——接受外部已排序的行指针数组:

流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

void SortOutputSpiller::spill(SpillRows& rows) {
  auto& run = createOrGetSpillRun(SpillPartitionId(0));
  // 直接接管外部已排序的行指针(sortedRows_ 的剩余部分)
  run.rows = SpillRows(rows.begin(), rows.end(), run.rows.get_allocator());
  for (const auto* row : rows) run.numBytes += container_->rowSize(row);
  markSeenPartitionsSpilled();
  runSpill(/*lastRun=*/true);  // 一次性完成,不分多 run
}

重写 runSpill:执行完后立即 finishFile,因为这是唯一一次写入。


7.14.6 NoRowContainerSpiller / MergeSpiller

needSort      : false
container_    : nullptr(无 RowContainer)
maxSpillRunRows: 0(不使用)

无 RowContainer,通过 spill(partitionId, spillVector) 直接将外部 RowVector 写入磁盘,跳过整个 fillSpillRuns → runSpill → extractSpill 流程:

流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

void NoRowContainerSpiller::spill(
    const SpillPartitionId& id, const RowVectorPtr& spillVector) {
  if (!state_.isPartitionSpilled(id)) state_.setPartitionSpilled(id);
  state_.appendToPartition(id, spillVector); // 直接写磁盘
}

MergeSpiller 继承 NoRowContainerSpiller 并携带排序键信息,供需要有序 spill/merge 的路径复用。普通 HashProbe 的输入和输出 spill 使用 NoRowContainerSpiller;不能仅凭 MergeSpiller 的名字就把它解释为 HashProbe 专用实现。


7.15 各子类对比

子类 needSort HashBitRange container_ 特殊行为
HashBuildSpiller false 多分区 ✓ 重写 extractSpill:附加 probed flag 列;双重 spill() 入口
AggregationInputSpiller true 多分区 ✓ 每 run 对应独立 sorted 文件;freeze allocator
AggregationOutputSpiller false 单分区(空) ✓ 从迭代器偏移开始;重写 runSpill 收尾
SortInputSpiller true 单分区(空) ✓ sort 列前置;每 run 独立 sorted 文件
SortOutputSpiller false 单分区(空) ✓ 接受外部已排序 SpillRows;重写 runSpill
NoRowContainerSpiller false 多分区 nullptr 直接 appendToPartition,无 fill/extract 流程
MergeSpiller false 多分区 nullptr 同上,额外保留 sortingKeys

8. 基础设施层详解

8.1 SpillConfig — 配置层

8.2 配置与内存触发条件

Velox 的全局 spill_enabled 默认 false;aggregation/join/order_by 的子开关默认 true,但只有全局开启且算子计划允许时才生效。宿主可以覆盖这些默认值。参见 spill 开关。

SpillConfig 汇集目录回调、文件大小、读写缓冲、压缩、hash 分区位、最大层数、spill executor、统计和 query 累计字节限制回调。它描述执行所需策略,并不自行选择哪个 Operator 成为回收对象。

常见触发有两类:

  • 算子在 ensureInputFits / ensureOutputFits 等位置估算下一步所需预留。增长可能触发仲裁和回收,随后继续或进入 spill 路径。
  • 内存仲裁通过 Task / Operator reclaimer 请求释放可回收状态。Task 暂停、不可回收区间和 peer 状态决定能否执行。

因此,不应把“预留失败”机械理解成一个完全独立、与仲裁无关的同步 if 分支。具体算子检查不同;配置百分比是用于保留/增长预留的策略参数,不是“used 达到固定百分比就必然 spill”的全局规则。

文件:velox/common/base/SpillConfig.h

SpillConfig 是所有 spill 行为的控制中枢。

流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

struct SpillConfig {
  // 回调函数(避免与具体 Task/Driver 实现耦合)
  GetSpillDirectoryPathCB   getSpillDirPathCb;         // 惰性获取 spill 目录
  UpdateAndCheckSpillLimitCB updateAndCheckSpillLimitCb; // 更新全局 quota 并检查是否超限

  // 文件参数
  std::string fileNamePrefix;          // spill 文件名前缀
  uint64_t    maxFileSize;             // 单文件大小上限,超过后自动 rotate
  uint64_t    writeBufferSize;         // 写缓冲区,攒够再 flush
  uint64_t    readBufferSize;          // 读缓冲区,支持 async 双缓冲预读

  // 执行参数
  folly::Executor* executor;           // 非 null 时 spill 写入在后台线程并发执行

  // 内存触发门槛
  int32_t  minSpillableReservationPct;      // 可 spill reservation 不低于 used 的 X%
  int32_t  spillableReservationGrowthPct;   // 一次 reserve 扩容 used 的 X%
  uint64_t maxSpillRunRows;                 // 单 run 最多行数(防止 OOM)
  uint64_t writerFlushThresholdSize;        // TableWriter flush 门槛

  // 分区参数(决定 spill 时的并行度)
  uint8_t  startPartitionBit;    // hash 分区起始 bit(e.g. 29)
  uint8_t  numPartitionBits;     // 每层分区 bit 数(e.g. 3 → 8 分区)
  int32_t  maxSpillLevel;        // 最大递归 spill 层数,-1 表示不限

  // 归并参数
  uint32_t numMaxMergeFiles;     // 归并时最多同时 open 的文件数(防 FD 耗尽)

  // 排序优化
  std::optional<PrefixSortConfig> prefixSortConfig; // 若设置则用 PrefixSort 代替 TimSort

  // 其他
  common::CompressionKind compressionKind;  // 压缩方式
  std::string fileCreateConfig;             // 传递给底层 FileSystem 的创建选项
  uint32_t windowMinReadBatchRows;          // Window 算子读取 batch 最小行数
};

关键计算方法:

// 返回当前 bit offset 对应的 spill level
int32_t spillLevel(uint8_t startBitOffset) const;

// 判断是否已超过最大 spill 层数
bool exceedSpillLevelLimit(uint8_t startBitOffset) const;

递归 spill 时,每下一层 startPartitionBit 向右移 numPartitionBits 位, spillLevel 就是 (startBitOffset - config.startPartitionBit) / numPartitionBits。


8.3 SpillStats — 统计层

文件:velox/exec/SpillStats.h

SpillStats 的主要计数与耗时字段使用 atomic,另有 IoStats ioStats 成员,不能把整个结构说成所有字段都是 std::atomic。单个计数可并发更新,也不意味着多个字段的读取组成同一时刻的一致快照。

流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

struct SpillStats {
  // 计数
  std::atomic_uint64_t spillRuns;           // spill 触发次数
  std::atomic_uint64_t spilledRows;         // spill 的行数
  std::atomic_uint64_t spilledInputBytes;   // spill 前内存中的原始字节数
  std::atomic_uint64_t spilledBytes;        // 写到磁盘的字节数(含压缩)
  std::atomic_uint32_t spilledPartitions;   // spill 的分区数
  std::atomic_uint64_t spilledFiles;        // 生成的 spill 文件数
  std::atomic_uint64_t spillMaxLevelExceededCount; // 超过最大 spill 层数次数

  // 写路径时间(纳秒)
  std::atomic_uint64_t spillFillTimeNanos;           // 填充 SpillRun(哈希分区)
  std::atomic_uint64_t spillSortTimeNanos;           // 排序
  std::atomic_uint64_t spillExtractVectorTimeNanos;  // 从 RowContainer 提取 RowVector
  std::atomic_uint64_t spillSerializationTimeNanos;  // PrestoPage 序列化
  std::atomic_uint64_t spillFlushTimeNanos;          // 压缩 + 拷贝到写缓冲
  std::atomic_uint64_t spillWriteTimeNanos;          // 写到文件系统
  std::atomic_uint64_t spillWrites;                  // write() 调用次数

  // 读路径时间(纳秒)
  std::atomic_uint64_t spillReadBytes;
  std::atomic_uint64_t spillReads;
  std::atomic_uint64_t spillReadTimeNanos;
  std::atomic_uint64_t spillDeserializationTimeNanos;
  IoStats ioStats;
};

每个算子持有自己的 SpillStats 对象(通过 spillStats_.get() 传递),同时每次更新都调用对应的 updateGlobal*() 函数汇总到进程级全局统计(globalSpillStats()),供 query profile 使用。


8.4 SpillFile — 文件 I/O 层

8.5 文件与分区:当前接口的变化

SpillWriter 和 SpillReadFile 建立在 SerializedPageFile 的写读接口上。文件中是可反序列化的批次,schema、压缩和排序信息通过相应文件描述传递;不能把 spill 文件当作 RowContainer 的直接 dump,也不应默认它是 Parquet。

SpillState::appendToPartition 为已标记分区追加数据;finishFile 结束当前文件,finish 将分区文件交给后续读取。SpillerBase::finishSpill 标记 finalized,并把本地分区与 parentId 组合成完整层级 ID,收集到 SpillPartitionSet。参见 finishSpill。

当前 SpillPartitionId 编码一条分区路径,每层最多 3 个 partition bits,支持 level 0 到 3。它不等同于旧模型中的“hash bit offset + partition number”;hash 位窗口由配置与恢复层次共同计算。比较顺序按祖先路径逐层排序,让相关子分区相邻。

配置允许值、64 位 hash 剩余位数和 ID 表示能力都是独立边界。即使某个低层配置方法用 -1 表示不限制,也不能据此声称整体可以无限递归。

文件:velox/exec/SpillFile.h/.cpp

8.5.1 SpillFileInfo — 元数据

流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

struct SpillFileInfo {
  uint32_t id;                              // 单调递增的文件 ID
  RowTypePtr type;                          // 数据类型
  std::string path;                         // 文件路径
  uint64_t size;                            // 文件字节数
  std::vector<SpillSortKey> sortingKeys;    // 排序键(无序 spill 则为空)
  common::CompressionKind compressionKind;  // 压缩方式
};
using SpillFiles = std::vector<SpillFileInfo>;

8.5.2 SpillWriter — 写路径

继承 SerializedPageFileWriter(底层用 PrestoPage 格式)。

构造参数:
  type, sortingKeys            : 数据类型与排序键
  compressionKind              : 压缩方式
  pathPrefix                   : 文件路径前缀,实际路径为 prefix + "_" + fileId
  targetFileSize               : 单文件目标大小
  writeBufferSize              : 写缓冲区大小
  updateAndCheckSpillLimitCb   : quota 回调

关键流程:
  write(rows, ranges)   → 序列化 RowVector 片段,写入缓冲
  finishFile()          → 强制 flush + 关闭当前文件,记录 SpillFileInfo
  finish()              → 收尾所有文件,返回 SpillFiles 列表

文件切分逻辑(父类 SerializedPageFileWriter 实现):
  每次 write 后,若当前文件大小 >= targetFileSize → 自动 finishFile()

每次 updateFileStats() 回调会调用 updateAndCheckSpillLimitCb,若全局 spill 字节超过 query limit 则抛 VELOX_SPILL_LIMIT_EXCEEDED 异常。

8.5.3 SpillReadFile — 读路径

create(fileInfo, bufferSize, pool, stats) → 打开文件,设置读缓冲
nextBatch(batch)                          → 反序列化下一 batch RowVector
                                          → 若文件系统支持 async IO,双缓冲预读

8.6 SpillPartitionId — 层级分区标识

文件:velox/exec/Spill.h:277

这是支持递归 spill 的核心数据结构。用单个 uint32_t 编码完整的层级路径:

Bit 布局(低位在前):
  bits [2:0]   : Level 1 分区号(最多 8 个)
  bits [5:3]   : Level 2 分区号
  bits [8:6]   : Level 3 分区号
  bits [11:9]  : Level 4 分区号
  bits [28:12] : 未使用
  bits [31:29] : 当前所在层级(0-3)

常量:
  kMaxSpillLevel  = 3   (最多 4 层,层级编号 0~3)
  kMaxPartitionBits = 3 (每层最多 8 分区)

构造方式:

流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

SpillPartitionId(uint32_t partitionNumber);           // 创建 level 0 分区
SpillPartitionId(SpillPartitionId parent, uint32_t n); // 创建子分区

有序性(operator<):

排序规则:
  1. 相同父节点的子分区按分区号从小到大排列
  2. 不同层级时,按层级路径字典序比较
  3. 浅层节点排在其子节点前面(DFS 前序)

示例:p_0 < p_1 < p_2_0 < p_2_1 < p_2_2 < p_3

这个有序性使 SpillPartitionSet(std::map<SpillPartitionId, ...>)能以 DFS 前序遍历,确保处理完父 partition 后再处理其子 partition(递归 spill 场景)。


8.7 SpillState — 分区状态管理

文件:velox/exec/Spill.h:581, Spill.cpp:131

每个 Spiller 持有一个 SpillState,负责管理所有分区的写入器:

流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

class SpillState {
  // 已 spill 分区的 ID 集合(用于快速判断某 partition 是否已开始 spill)
  SpillPartitionIdSet spilledPartitionIdSet_;

  // partition -> SpillWriter 的映射,folly::Synchronized 保证线程安全
  // 不同 partition 可以并发写入(不同线程写各自的 partition)
  folly::Synchronized<SpillPartitionWriterSet> partitionWriters_;
};

**核心方法 appendToPartition**:

uint64_t SpillState::appendToPartition(
    const SpillPartitionId& id, const RowVectorPtr& rows) {

  // 1. 惰性获取 spill 目录(每次调用都通过 callback,目录可能在首次写时才创建)
  auto spillDir = getSpillDirPathCb_();

  // 2. 若该 partition 还没有 writer,在 wLock 内创建(路径格式:dir/prefix-spill-encodedId)
  partitionWriters_.withWLock([&](auto& lockedWriters) {
    if (!lockedWriters.contains(id)) {
      lockedWriters.emplace(id, std::make_unique<SpillWriter>(
          rows->type(), sortingKeys_,
          compressionKind_,
          fmt::format("{}/{}-spill-{}", spillDir, fileNamePrefix_, id.encodedId()),
          targetFileSize_, writeBufferSize_, ...));
    }
  });

  // 3. 验证序列化大小不超过 2GB(PrestoPage 限制)
  validateSpillBytesSize(rows->estimateFlatSize());
  updateSpilledInputBytes(bytes);

  // 4. 写入(不持有 wLock,不同 partition 并发写互不干扰)
  IndexRange range{0, rows->size()};
  return partitionWriter(id)->write(rows, folly::Range<IndexRange*>(&range, 1));
}

关键设计:partitionWriters_ 只在创建 writer 时加 wLock,实际 write() 不加锁——因为每个 partition 只有一个 writer,不同 partition 的写入天然不冲突。


8.8 SpillPartition / SpillPartitionSet — 读回路径

8.9 两种读回路径,以及有限归并 fan-in

当前 SpillMergeStream::current() 返回 const RowVectorPtr&。需要 RowVector* 的 gather 数组使用 current().get();访问列用 current()->childAt(...)。借用裸指针时仍要在 pop 可能推进到下一批之前完成复制,或显式保活原批次。见 current 与 currentIndex 契约。

读取方式 结构 典型用途
unordered reader BatchStream / UnorderedStreamReader HashJoin 重建表
ordered reader SpillMergeStream / TreeOfLosers Aggregation 合并同 key、OrderBy 合并有序 run

unordered 只需逐批把数据交回算子,不保证整体有序。ordered reader 依赖每条输入 stream 已有序,每次选择最小行,按需要维护相等 key 的信息。

当前 createOrderedReader 支持 numMaxMergeFiles。文件过多时,先选较小文件分轮归并,缩小最后同时参与归并的文件数量;0 不限制,1 被拒绝。这会增加额外读写,但控制同时打开的流及缓冲。这种多轮有序文件归并,不是 HashJoin 用更多 hash bits 做的递归分区。

图 3:HashJoin 读回后重建表;Aggregation 和 OrderBy 依赖有序流,必要时先做多轮文件归并。
图 9:HashJoin 读回后重建表;Aggregation 和 OrderBy 依赖有序流,必要时先做多轮文件归并。

文件:velox/exec/Spill.h:457

流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

class SpillPartition {
  SpillPartitionId id_;
  SpillFiles files_;      // 该 partition 的所有文件
  uint64_t size_;         // 总字节数

  // 无序读(HashJoin 用):逐文件顺序消费,不保序
  std::unique_ptr<UnorderedStreamReader<BatchStream>> createUnorderedReader(...);

  // 有序读(Aggregation/OrderBy 用):N 路败者树归并,保序
  std::unique_ptr<TreeOfLosers<SpillMergeStream>> createOrderedReader(...);

  // 均分文件列表为 N 个 shard(并行算子使用)
  std::vector<std::unique_ptr<SpillPartition>> split(int numShards);
};

// SpillPartitionSet 是按 SpillPartitionId 有序排列的 map
using SpillPartitionSet = std::map<SpillPartitionId, std::unique_ptr<SpillPartition>>;

createOrderedReader 的防 FD 耗尽机制:

若 files_.size() > numMaxMergeFiles(且 numMaxMergeFiles >= 2):
  1. 将文件分组,每组最多 numMaxMergeFiles 个
  2. 对每组文件创建 TreeOfLosers 做局部归并,写出中间文件
  3. 递归,直到文件数 <= numMaxMergeFiles
  4. 最终对剩余文件建 TreeOfLosers 返回

这避免了同时打开 O(所有文件) 个文件描述符导致的 OOM 或 FD 耗尽。

8.10 读流层:SpillMergeStream / BatchStream

文件:velox/exec/Spill.h:56

BatchStream — 无序批次流(HashJoin 用):

接口:nextBatch(batch) → bool
实现:FileSpillBatchStream        - 单文件
      ConcatFilesSpillBatchStream - 多文件顺序串联

SpillMergeStream — 有序合并流(Aggregation/OrderBy 用):

接口:hasData() / currentIndex() / current() / pop()
实现:FileSpillMergeStream        - 单文件
      ConcatFilesSpillMergeStream - 多文件顺序串联(文件间已有序)

关键方法 compare():按 sortingKeys_ 逐列比较当前行,供 TreeOfLosers 使用
关键方法 pop():消费当前行;若当前 batch 用完,调用 nextBatch() 加载下一批

TreeOfLosers<SpillMergeStream> 是 N 路归并排序的败者树实现,next() 每次返回所有流中最小的行所在的流。


9. 递归 Spill 流程图

HashBuild 的写出与恢复循环
图 10:HashBuild 的写出与恢复循环。已按当前实现修正标注,具体约束见相邻正文。
Fig. HashBuild 递归 spill 全流程:内存不足触发 spill,逐 spill partition 递归恢复直到 kFinish

10. 关键设计决策总结

10.1 限制、统计和清理

spillMemoryPool 为写出临时状态提供资源,避免所有临时分配都继续压在正在回收的 query/operator pool 上;它不是无限内存。SpillRun 指针、提取向量、压缩 buffer、读回 stream 和 merge state 都有峰值成本。

updateSpilledBytesAndCheckLimit 维护 query 级累计写出字节并检查限制。递归分区、重新写出和中间文件归并都可能让总 I/O 超过原始数据量。磁盘占用、累计 spilled bytes 和逻辑行大小不是同一个指标。

SpillStats 分别记录 rows/bytes/files/partitions、fill/sort/extract/serialization/write/read 等时间和计数。多个并行任务的累计时间可能超过 wall time,不要简单相加当成查询耗时。分析时要连同仲裁等待、Task pause、恢复次数和最大 spill level 一起看。

SpillReadFile 的生命周期契约 明确说明:析构不删除文件,调用方负责后续清理,例如 Task 统一清理生成目录。因此“reader 被释放”与“磁盘文件已删除”要分开验证。

10.2 有序 vs. 无序 Spill

算子 Spill 类型 原因
HashJoin Build 无序 恢复时重建 hash 表,顺序无关
HashAggregation Input 有序(按 group key) 支持流式归并聚合,避免全量恢复
HashAggregation Output 无序(hash 表 scan 顺序) 已是输出顺序,不需要重排
OrderBy Input 有序(按 sort key) 支持多路归并产生有序输出
OrderBy Output 有序(已排序传入) 同上

10.3 两阶段 Spill(Input/Output Spiller 分离)

Aggregation 和 OrderBy 都支持两阶段 spill:

  • Input 阶段可多次写出并清理内存状态:AggregationInputSpiller 按 grouping keys 做 hash 分区并排序,SortInputSpiller 则在单分区中形成有序 run。两者都支持后续归并,但分区规则不同。
  • Output 阶段处理已经完成聚合或排序、但尚未交给下游的剩余行。它沿相应输出迭代位置写出,避免重放已输出前缀;此路径的单分区和阶段约束应按具体子类实现理解,不能把输出写出等同于重新聚合全部输入。

两阶段互斥,通过 inputSpiller_/outputSpiller_ 非空判断。

10.4 内存预留机制(proactive spill)

触发链:ensureInputFits → maybeReserve → 触发内存仲裁 → reclaim → spill

关键参数:
  minSpillableReservationPct     确保始终有足够"可 spill"预留空间
  spillableReservationGrowthPct  每次扩容幅度,平衡内存利用率与 spill 频率

设计意图:通过 proactive reservation,在真正 OOM 之前触发 spill,避免在不可中断的关键路径上被动 OOM。

10.5 HashBuild 多线程 Spill 协调

所有 peer HashBuild 线程共享同一个 joinBridge
仲裁器触发 reclaim 时:
  1. 暂停所有 peer driver(task->pause())
  2. 检查所有 peer 均可回收
  3. 统一 spill 所有 peer 的 hash 表
  4. 清空所有 peer 的表 + 释放预留

这确保所有线程的 hash 表碎片被一次性清理,而不是只清理一个线程的。

10.6 递归 Spill(HashJoin 专属)

初次 spill:例如使用连续 3 个 hash 位 → 最多 8 个 spill partitions(具体起始位由配置决定)
第一次递归:使用下一段可用 hash 位 → 最多 8 个子 partitions
最多 maxSpillLevel 层(默认 3 层)

终止条件:
  1. 达到 maxSpillLevel
  2. 某 partition 能完全装入内存(不再 spill)
  3. partition 数据为空

10.7 probed flag Spill(Right/Full Outer Join)

对于 right join / full outer join,hash 表中每行有一个 probed 标记位,记录该行是否已与 probe 行匹配。Spill 时需要将此标记一同写入磁盘(作为额外的 bool 列),恢复时读回并设置回 RowContainer 的对应行。


11. C++ 编码品味:值得学习的细节

11.1 从当前实现学到的几个编程约束

写法 解决的问题 不能代替什么
SCOPE_EXIT / makeGuard 异常出口仍完成通知与异步收尾 不能修复错误所有权
exception_ptr 把并发任务错误交回调用线程 不是忽略写盘错误
stateCleared / finalized / sorted 明确阶段与一次性动作 不是线程同步原语
锁内取走状态、锁外 notify 避免通知回调扩大临界区 不是所有函数都无重活锁
TestValue::adjust 注入特定交错、回收与异常 不代表生产必走测试路径
release() 归还未使用的 reservation 不释放仍有引用的数据

特别是 Bridge::reclaim 当前会持锁执行 spill callback,而 Spiller 的异步写入返回前又要收齐结果。讨论“非阻塞”时必须说清是 Driver 等 future 让出 executor,还是某个同步 reclaim 调用仍占用当前线程。

这一章从代码本身出发,梳理 Velox spill 系统在 C++ 工程实践上值得学习的具体技法。每一条都有对应的代码位置。


11.2 SCOPE_EXIT:让清理逻辑紧贴触发逻辑

流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

// HashBuild::finishHashBuild
SCOPE_EXIT {
  // 无论函数从哪个 return 路径退出,都必须唤醒等待的 peer
  peers.clear();
  for (auto& promise : promises) {
    promise.setValue();
  }
};
// SortBuffer::getOutput
SCOPE_EXIT {
  pool_->release(); // 每次 getOutput 返回后都归还未用预留
};

SCOPE_EXIT(底层是 folly::ScopeGuard)比 try-catch 清晰得多:清理代码与触发它的上下文紧邻,阅读者一眼就能看到"这个函数结束时会做什么",而不需要去找 finally 块或析构函数。

对比写法:

// ❌ 反例:try-catch 把清理和逻辑分离
try {
  doWork();
  pool_->release(); // 容易漏掉 early return 路径
} catch (...) {
  pool_->release();
  throw;
}

// ✅ SCOPE_EXIT:写一次,所有路径都覆盖
SCOPE_EXIT { pool_->release(); };
doWork();

11.3 folly::makeGuard:异步任务的"排水阀"

流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

// SpillerBase::runSpill
auto sync = folly::makeGuard([&]() {
  for (auto& write : writes) {
    try {
      write->move();  // 强制消费所有 pending 任务
    } catch (const std::exception&) {
      // 清理路径不能 throw
    }
  }
});

// 主线程消费结果(可能 throw)
for (auto& write : writes) {
  results.push_back(write->move()); // 若某个任务失败,这里 throw
}
// 若 results 循环 throw,guard 析构时仍会 drain 剩余任务

这里的 makeGuard 扮演"排水阀"角色:无论主逻辑成功还是抛异常,都确保所有后台异步任务被 move()(消费)掉,避免后台任务持有的资源(文件句柄、内存)泄漏。

关键细节:guard 内部的 catch (const std::exception&) 是有意的——清理路径本身不能 throw,第一个错误已经在 results 循环里被捕获了。


11.4 跨线程异常传递:exception_ptr 模式

多线程 spill 时,后台线程的异常无法直接传播到主线程。Velox 使用 std::exception_ptr 将异常"装箱":

流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

// 后台线程(writeSpill)
try {
  // ... 做实际工作
  return std::make_unique<SpillStatus>(id, written, nullptr); // 成功
} catch (const std::exception&) {
  // 捕获任何异常,装箱为 exception_ptr
  return std::make_unique<SpillStatus>(id, 0, std::current_exception());
}

// 主线程(runSpill)
for (auto& result : results) {
  if (result->error != nullptr) {
    std::rethrow_exception(result->error); // 在主线程重新抛出
  }
  // ... 处理正常结果
}

结果对象 SpillStatus 的定义简洁体现了这个模式:

struct SpillStatus {
  SpillPartitionId partitionId;
  uint32_t rowsWritten;
  std::exception_ptr error; // nullptr 表示成功,否则携带异常
};

这使得多线程错误处理与单线程代码风格一致——主线程的调用者只需检查返回值或 catch 异常,不需要感知线程。


11.5 "1 + N" 并发策略:第一个任务不切线程

流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

for (const auto& [id, spillRun] : spillRuns_) {
  writes.push_back(
      memory::createAsyncMemoryReclaimTask<SpillStatus>(
          [partitionId = id, this]() { return writeSpill(partitionId); }));

  // 只有第 2 个及之后的任务才提交到后台 executor
  if ((writes.size() > 1) && executor_ != nullptr) {
    executor_->add([source = writes.back()]() { source->prepare(); });
  }
}

// 主线程顺序 move(),第一个任务在这里内联执行
for (auto& write : writes) {
  results.push_back(write->move()); // 若未 prepare(),move() 内联执行 lambda
}

设计精妙处:AsyncSource::move() 的语义是"如果还没执行就内联执行,已经执行就等待结果"。第一个任务从未被 prepare(),所以在 move() 调用时内联执行于主线程。

好处:

  • 只有一个任务时,不必为了这次分区写出额外提交 executor;文件系统、压缩或其他下层操作仍可能等待或调度,不能推导整次 spill 零线程切换。
  • 调用线程参与第一个任务,减少把全部任务提交后立即等待的开销;取得其他任务结果时仍可能等待,因此并不保证它全程都在计算。
  • 调用线程与 executor 中的工作可以重叠;实际并行度和收益仍受 executor 大小、分区分布、I/O、压缩以及内存峰值制约。

HashProbe::spillOutput 对多个 probe 算子也用了完全一样的模式。


11.6 CHECK vs DCHECK:按热度分配断言代价

流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

// 生产代码的不变式用 VELOX_CHECK(release 也执行)
VELOX_CHECK(!finalized_);                       // finishSpill 之后不能再写
VELOX_CHECK_EQ(numWritten, run.rows.size());    // 写出的行数必须对得上

// 热路径上的显然正确断言用 VELOX_DCHECK(只在 debug 执行)
VELOX_DCHECK_GE(partitionNum, 0);              // hash 分区号必然 >= 0(inner loop)
VELOX_DCHECK(isPartitionSpilled(id));          // appendToPartition 前必须已设置

原则:

  • VELOX_CHECK_* 用于外部可观察的不变式,违反时意味着调用方 bug 或状态异常,即使在 release 也要发现
  • VELOX_DCHECK 用于主要在调试构建检查的内部不变量;关闭检查不代表条件天然正确。选择 CHECK/DCHECK 应同时考虑错误后果、发布版需要的诊断和热路径成本,不能只根据“理论上不会错”决定。

在 fillSpillRuns 的内层循环(每行都执行一次)里用 VELOX_DCHECK 而非 VELOX_CHECK,避免 release 构建中每行多一次分支,对亿级行场景有实际性能影响。


11.7 FOLLY_LIKELY / FOLLY_UNLIKELY:把预期写进代码

流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

// createOrGetSpillRun:绝大多数情况 spillRun 已存在
if (FOLLY_UNLIKELY(!spillRuns_.contains(id))) {
  spillRuns_.emplace(id, SpillRun(*memory::spillMemoryPool()));
}

// getOutputWithSpill:几乎每行都需要复制,到达 batch 边界是例外
if (FOLLY_UNLIKELY(isEndOfBatch)) {
  gatherCopy(...);
}

// 最后一批输出大多非空
if (FOLLY_LIKELY(outputSize != 0)) {
  gatherCopy(...);
}

FOLLY_LIKELY/UNLIKELY 不仅是编译器分支预测提示,更是可执行的注释——它明确告诉读者"作者预期这条路径是 hot/cold 的"。这比写 // rarely happens 注释更可靠,因为注释会过期,而这行代码会随着逻辑一起维护。


11.8 TestValue::adjust:无侵入式测试注入

流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

// SpillerBase 构造函数
TestValue::adjust("facebook::velox::exec::SpillerBase", this);

// HashBuild::addInput 入口
TestValue::adjust("facebook::velox::exec::HashBuild::addInput", this);

// HashBuild::reclaim 入口
TestValue::adjust("facebook::velox::exec::HashBuild::reclaim", this);

TestValue::adjust 提供命名注入点,便于测试在指定时刻修改状态或抛出异常。是否编译掉及其调用成本取决于构建条件;应检查当前 TestValue 定义与构建选项,不能由“测试接口”直接推断所有生产构建均零开销。

配套的 TestScopedSpillInjection 更进一步:

// 测试代码:
TestScopedSpillInjection injection(20, ".*", 10);
// poolName 匹配 ".*" 时,testingTriggerSpill() 有 20% 概率返回 true,最多触发 10 次

// 生产代码:
if (testingTriggerSpill(pool_->name())) {
  spill(); // 测试强制触发 spill
}

测试注入能精确覆盖特定交错;其是否启用与代价必须结合构建配置判断。真正的并发正确性仍需要状态、锁和生命周期协议保证。


11.9 NanosecondTimer:RAII 计时的最小实现

流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

uint64_t execTimeNs{0};
{
  NanosecondTimer timer(&execTimeNs);  // 构造时记录起始时间
  // ... 被测代码 ...
} // 析构时写入 execTimeNs

spillStats_->spillFillTimeNanos.fetch_add(execTimeNs, std::memory_order_relaxed);

几个值得注意的细节:

① 先积累,再 fetch_add:不在块内直接 fetch_add,而是先存到局部变量,块结束后一次 atomic 操作。减少 atomic 操作次数,也使代码结构更清晰。

**② memory_order_relaxed**:spill 统计不需要与其他内存操作同步,用 relaxed 避免不必要的内存屏障。

③ 一个特殊处理:

// NOTE: Always set a non-zero sort time to avoid flakiness in tests which check sort time.
updateSpillSortTime(std::max<uint64_t>(1, sortTimeNs));

即使排序耗时 0 纳秒(几乎不可能但理论上存在),也写入 1,防止测试误判"排序未发生"。这是一个小而精的防御性设计。


11.10 SpillPartitionId:把所有接口设施配齐

SpillPartitionId 是一个需要在多种容器中使用的值类型,代码为它配齐了所有标准设施:

流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

// 1. 比较运算符:支持 std::map(有序容器)
bool operator<(const SpillPartitionId& other) const;
bool operator>(const SpillPartitionId& other) const {
  return other < *this; // 复用 < 实现 >,不重复逻辑
}
bool operator==(const SpillPartitionId& other) const = default; // C++20 default

// 2. std::hash 特化:支持 folly::F14FastMap/Set
namespace std {
template <>
struct hash<SpillPartitionId> {
  uint32_t operator()(const SpillPartitionId& id) const {
    return std::hash<uint32_t>()(id.encodedId()); // 直接哈希底层整数,O(1)
  }
};
}

// 3. fmt::formatter 特化:支持 fmt::format("{}", id)
template <>
struct fmt::formatter<SpillPartitionId> : formatter<std::string> {
  auto format(SpillPartitionId s, format_context& ctx) const {
    return formatter<std::string>::format(s.toString(), ctx);
  }
};

// 4. operator<< :支持 LOG(INFO) << id
inline std::ostream& operator<<(std::ostream& os, SpillPartitionId id) {
  return os << id.toString();
}

operator> 通过调用 operator< 实现而非独立写逻辑,是很好的习惯——减少了两个实现可能不一致的风险。


11.11 withWLock / withRLock:锁的最小化持有

流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

// SpillState::appendToPartition
partitionWriters_.withWLock([&](auto& lockedWriters) {
  // 只在创建 writer 时持有写锁
  if (!lockedWriters.contains(id)) {
    lockedWriters.emplace(id, std::make_unique<SpillWriter>(...));
  }
}); // 锁在此释放

// 实际写入在锁外进行(不同 partition 的写入完全并发)
return partitionWriter(id)->write(rows, ...);
// testingSpilledFilePaths(只读,用 rlock)
partitionWriters_.withRLock([&](const auto& partitionWriters) {
  for (const auto& [id, writer] : partitionWriters) {
    // ... 读取
  }
});

folly::Synchronized 的 withWLock/withRLock lambda 接口使锁的范围在代码结构上可见——lambda 的括号就是临界区的边界,不需要手动 lock_guard。

关键设计:写锁只保护"创建 writer"(初始化一次),不保护"写数据"。这是因为每个 partition 只有一个 writer,不同 partition 的写入天然不冲突,无需锁保护。最快的锁是不持有锁。


11.12 SpillRun::sorted:防止双重排序的状态位

流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

struct SpillRun {
  SpillRows rows;
  uint64_t numBytes{0};
  bool sorted{false}; // 是否已排序
};

void SpillerBase::ensureSorted(SpillRun& run) {
  if (run.sorted || !needSort()) { // 已排序则跳过
    return;
  }
  // ... 执行排序
  run.sorted = true; // 标记,防止重复排序
}

sorted 标志记录本轮已经按同一比较规则排序,主要避免重复工作。再次使用相同正确比较器排序不会破坏有序性,但不稳定排序可能改变相等键之间的相对顺序;SQL 未声明的稳定顺序也不能由此获得保证。状态位本身不是互斥或内存同步原语。


11.13 pool_->release():主动归还,而非等待回收

流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

// SortBuffer::noMoreInput()
// Releases the unused memory reservation after processing input.
pool_->release();

// SortBuffer::getOutput()
SCOPE_EXIT {
  pool_->release();
};

// HashBuild::ensureTableFits() 中扩容后
if (spiller_->spillTriggered()) {
  pool()->release(); // spill 已触发,扩容的预留用不上了,主动归还
}

pool_->release() 的语义是"我承诺不再需要当前预留的超额内存"。Velox 在每个"阶段完成"的节点系统性地调用它,而不是等待 GC 或析构。

release() 归还不再需要的预留,使父级 reservation 有机会下降;它不会释放仍存活的 buffer,也不会直接让其他 root 获得 capacity。跨 query 再分配还需要仲裁器 shrink/grow 的额度转移。应分别看 used、reservation 和 capacity 的变化。


11.14 succinctBytes():有信息量的日志

流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

LOG(WARNING) << "Failed to reserve " << succinctBytes(targetIncrementBytes)
             << " for memory pool " << pool_->name()
             << ", root pool: " << pool_->root()->name()
             << ", used: " << succinctBytes(pool_->usedBytes())
             << ", reservation: " << succinctBytes(pool_->reservedBytes())
             << ", root pool reservation: "
             << succinctBytes(pool_->root()->reservedBytes());

几个细节值得学习:

succinctBytes() 按可读单位表示字节数,减少人工换算;这是日志表达的改进,不对应一个可量化的固定倍数收益。

多处 spill 日志按 pool 名、root 名、used、reservation 和 root reservation 展示资源上下文,便于沿树定位压力。不同失败点仍有各自字段和格式;这是可参考的写法,不是所有日志完全一致的接口保证。

③ WARNING 而非 ERROR:预留失败不是错误,只是触发 spill 的信号。级别选择准确。


11.15 testing* 前缀:测试可观测性的统一约定

流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

// SpillState::testingSpilledFilePaths()
// SpillState::testingSpilledFileIds()
// SpillState::testingNonEmptySpilledPartitionIdSet()
// HashBuild::testingExceededMaxSpillLevelLimit()
// HashProbe::testingHasInputSpiller()
// HashProbe::testingExceededMaxSpillLevelLimit()

多处测试 accessor 采用 testing* 前缀,便于识别测试可观测性入口。前缀搜索能找到这一类方法,但不能据此穷尽所有测试接口、TestValue 钩子或 friend 访问。

  • 接口清晰度:任何 testing* 方法,调用者立刻知道"这不是业务接口"
  • 易于检索:rg "testing[A-Z]" 可定位遵循此前缀的符号;完整审核还需检查测试类、注入点和其他公开接口。
  • testing* accessor 让测试可观测性具有明确入口,便于评审暴露范围;局部仓库协作文件中的风格偏好不等于整个上游项目永远禁止 friend,不能把它当作接口设计的强制因果。

11.16 constexpr 局部常量:让魔数有名字

流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

void SpillerBase::fillSpillRuns(...) {
  constexpr int32_t kHashBatchSize = 4096; // 每次从 RowContainer 取多少行
  // ...
}

std::unique_ptr<SpillerBase::SpillStatus> SpillerBase::writeSpill(...) {
  constexpr int32_t kTargetBatchBytes = 1 << 18; // 256K
  constexpr int32_t kTargetBatchRows = 64;
  // ...
}

这些常量定义在函数内部而非文件级别,有两个好处:

① 最小作用域:kHashBatchSize 只在 fillSpillRuns 里有意义,定义在全局或类级别会误导读者以为它有更广泛的语义。

② 紧贴使用处:常量定义和使用之间距离最短,读者不需要跳转文件来理解这个数字的含义。

1 << 18 比 262144 更清楚地表达了"这是 2 的幂次"的意图,同时命名为 kTargetBatchBytes 说明了它是目标(上限)而非精确值。


11.17 编码品味小结

技法 核心价值 代表位置
SCOPE_EXIT 清理与触发紧邻,所有 return 路径都覆盖 SortBuffer::getOutput, HashBuild::finishHashBuild
folly::makeGuard 排水阀 异步任务无论成败都被消费,防止资源泄漏 SpillerBase::runSpill
exception_ptr 跨线程传递 后台线程异常装箱,主线程重抛,调用方无感知 SpillerBase::writeSpill / runSpill
"1 + N" 并发 首任务在调用线程执行,后续可并发;单 partition 可避免一次 executor 提交,仍有普通调用成本 SpillerBase::runSpill
DCHECK 分热路径 release 不付 inner loop 断言代价 fillSpillRuns inner loop
FOLLY_LIKELY/UNLIKELY 给编译器的分支概率提示;工作负载变化后可能失效,收益需测量 createOrGetSpillRun, getOutputWithSpill
TestValue::adjust 命名注入点;启用条件与成本依构建配置判断 所有关键函数入口
NanosecondTimer RAII 单行计时,早 return 安全 fillSpillRuns, ensureSorted
配齐值类型接口 <, >, ==, hash, formatter, << 一起提供 SpillPartitionId
withWLock 最小临界区 锁只保护初始化,写操作在锁外 SpillState::appendToPartition
sorted 防御位 防双重排序,用状态堵死隐患 SpillRun / ensureSorted
pool_->release() 主动归还 push 模型快速回转,而非等待析构 noMoreInput, getOutput
succinctBytes() 统一日志格式 固定信息层次,可读,可预期 所有 LOG(WARNING)
testing* 前缀约定 测试入口可辨;是否使用 friend 需按具体封装代价评估 所有 testing* 方法
constexpr 局部常量 最小作用域,紧贴使用处,有名字 kHashBatchSize, kTargetBatchRows

12. 从恢复需求推导 Spill 的分层

Velox spill 系统的代码量约 12,000 行,但它的设计并不复杂——真正值得品味的,是每一处决策背后的约束权衡和系统性思考。这一章从"为什么要这样设计"出发,梳理其核心设计哲学。


12.1 资源有限时扩大可执行工作集的范围

数据库引擎处理大数据时,内存不可能无限大。任何算子都需要回答这个问题:当数据超过可用内存时,怎么办?

Velox 通过可回收算子的 spill 路径扩大可执行数据规模,但不承诺开启 spill 后查询一定成功。不可回收区、不可 spill 的算子或执行模式、磁盘限额、I/O 错误、倾斜、递归层数和恢复时的内存需求仍可能导致失败。设计目标是提供可控的外部执行路径,而不是消除所有资源限制。

这一承诺对系统设计造成了深远影响:

  • 不能假设数据能放进内存——每一个数据结构都要有"放不下时的出路"
  • 不能只支持一次 spill——必须支持递归 spill(spill 后恢复,恢复时再次 OOM,再次 spill)
  • 回收请求可以在内存压力下到来,但真正进入算子 reclaim 必须满足 Task 安全停稳、算子可回收区及阶段约束;不能在任意正在修改内部结构的位置直接打断并 spill。

这三点决定了架构的基本形态。


12.2 分层:每层只管自己的事

整个 spill 系统被分成 5 层,每层接口极简:

算子层     → 决定「何时」和「触发哪种」spill
Spiller层  → 决定「如何」把 RowContainer 内容写到磁盘(行→列,分区,排序)
SpillState → 决定「写到哪个文件」
SpillFile  → 决定「如何」序列化/反序列化(PrestoPage 格式)
SpillConfig→ 决定「用什么参数」

这种分层带来的核心好处不是"解耦"这么虚的东西,而是每一层可以独立演进:

  • 更换序列化格式时,文件层可以集中部分读写改动,但 reader、writer、配置、类型支持和兼容性仍需一起检查。接口隐藏了 I/O 组织细节,不代表上层对任何格式替换都无感。
  • 要支持新算子 spill?只实现新的 Spiller 子类,底层文件 I/O 不变
  • 要调整内存触发策略?只改算子层的 ensureInputFits,文件写入逻辑不变

这在一个长期演进的系统里价值巨大。HashJoin、HashAggregation、OrderBy 三类算子的 spill 逻辑差异极大,但它们复用了同一套文件 I/O 和统计框架,维护成本大幅降低。


12.3 分区:把"全量恢复"变成"分批恢复"

Spill 最天真的设计是:内存不够时,把所有数据写到一个文件;恢复时,把整个文件读回来。

问题在于:如果数据总量是内存的 10 倍,一次性读回仍然 OOM。

Velox 的解法是按 hash 分区:

数据 → hash(key) → partition 0, 1, 2, ..., 7(8 个分区)

恢复时:
  一次只加载 partition 0 → 处理完 → 加载 partition 1 → ...
  均匀分布时单分区约 1/8 的行;倾斜、工作缓冲和重建开销会改变内存峰值

恢复复用了 processSpillInput → addInput 的行处理能力,但有专门的入口协议:setupSpillInput 清理旧 table/spiller/reader,重置 key 与 dependent channels,重建表并按恢复分区计算下一层 hash 位;processSpillInput 逐批输入,处理 yield / 状态变化,结束后调用 noMoreInputInternal。复用的是批次处理,分区身份、递归深度和轮次状态不能省略。

SpillPartitionId 的 bit 编码设计支持这种分批恢复天然地扩展到多层递归:

Level 0:bit[31:29] = 0, 分区号存在 bit[2:0]  → 8 个 partition
Level 1:bit[31:29] = 1, 分区号存在 bit[5:3]  → 每个 Level 0 partition 又有 8 个子 partition
Level 2:...

递归受到配置层数、可用 hash 位范围和 SpillPartitionId 表示能力的共同限制。当前 ID 每层最多 3 个分区位,表示 level 0..3;不能由某个配置项的“不限制”值推导整个系统无限递归。


12.4 有序与无序的分叉:根据"读回时的需求"决定

Spill 时要不要排序,完全由读回时的需求决定,而不是写入时的方便:

算子 读回时的需求 要不要排序
HashJoin Build 重建 hash 表(插入顺序无关) 不排序(节省 CPU)
HashAggregation 流式归并聚合(需要相同 key 相邻) 排序(按 group key)
OrderBy 多路归并产生有序输出(需要各 run 内有序) 排序(按 sort key)

这个设计的价值在于:对于不需要排序的场景,完全不付排序的代价。HashJoin 是最常见的 spill 场景,它不排序节省了大量 CPU 时间。

而对于需要排序的场景(Aggregation、OrderBy),排序在 spill 时一次性完成,读回时的 N 路归并就是线性扫描,极其高效。如果不在 spill 时排序,读回时就必须全量加载再排序,峰值内存反而更高。


12.5 两阶段 Spill:尊重算子的生命周期

Aggregation 和 OrderBy 各有两个 Spiller(input/output),这不是过度设计,而是对算子生命周期的精确建模:

Input 阶段:数据持续流入 → 多次触发 spill → hash 表/排序缓冲反复清空重填
Output 阶段:数据已全部处理 → 只可能触发一次 spill → 之后只有读出

两个阶段的需求截然不同:

Input 阶段分别使用 AggregationInputSpiller 和 SortInputSpiller:前者按 group key 分区并生成有序数据,后者为单分区外部排序生成 run。它们复用基础写出机制,但不能套用完全相同的分区和恢复策略。

  • 聚合输入可按 grouping keys 分区;OrderBy 输入保留单分区并用多路有序归并完成全局排序。两者的目的都是控制外部执行的工作集,采用的组织方式不同。
  • 需要排序(支持流式归并)
  • 可以多次触发(每次清空 hash 表继续处理输入)

Output 阶段 spill(AggregationOutputSpiller / SortOutputSpiller):

  • 无需分区(只触发一次,全量写出)
  • 无需排序(数据已经就绪,直接写入已有顺序)
  • 只触发一次(之后只读,无需再写)

用同一个 Spiller 处理两个阶段会导致代码复杂化,且无法针对性优化。分开设计使每个 Spiller 的逻辑极度简单。


12.6 SpillRun 切分:在多个约束之间找到平衡点

maxSpillRunRows 参数控制每次 fill-run 的行数,这个设计同时解决了四个互不相关的问题:

约束 1:指针数组内存(SpillRun.rows 每行 8 字节,1 亿行 = 800MB)
约束 2:排序辅助内存(大 run 排序 cache miss 严重)
约束 3:PrestoPage 2GB 序列化限制
约束 4:指针与中间批次峰值;源 RowContainer 的释放点由算子决定,不能将清理 run 指针视为已经释放源行

一个参数,四个收益,这是典型的"一石多鸟"设计。

值得注意的是,这个设计并没有在代码里写四个 // TODO: 原因X,而是通过一个统一的机制自然地满足了所有约束。这种"机制而非策略"的思维是优雅设计的标志。


12.7 批量行列转换:峰值内存与 I/O 效率的双重优化

extractSpillVector 使用最多 64 行的行数目标与 256 KiB 的字节目标组织批次;字节限制不是绝对硬上限。实现累计行大小,在达到目标时仍保留最后纳入的一行,单个大值也可能超过目标。它控制常见批次的物化峰值,后续序列化 / 压缩缓冲还要单独记账。

峰值内存:

不切分时:SpillRun 10 万行 → 物化为一个 10 万行 RowVector → 可能数 GB
切分后:  每次 64 行 RowVector → 几百 KB → prepareForReuse 复用同一块内存
          物化峰值受单批影响,但总峰值还包括 run 指针、源行、序列化缓冲及并行任务

I/O 效率:

SpillWriter 有 writeBufferSize 写缓冲(默认几 MB)
每次 appendToPartition 追加到缓冲,满了才刷盘
64 行 ≈ 几十 KB → 多次 append 后缓冲满 → 一次大 write()
这比每 10 万行一次大 write() 的 I/O 效率相当,但内存峰值小得多

64 行目标使固定宽度列的单批数据较小,有利于控制物化峰值;真实工作集同时包含多列和变长值。源码给出了参数与提取行为,但不能据此宣称“最优 cache 利用率”或固定的 L1 命中保证。


12.8 主动预留 vs. 被动仲裁:两道防线

Velox 的内存触发有两条路径,共同构成防 OOM 的双重保险:

防线 1 — 主动预留(proactive):
  ensureInputFits → maybeReserve(targetIncrementBytes)
    ├─ 成功:继续处理(预留了足够空间)
    └─ 失败:[内存仲裁在此可能已触发 reclaim]

防线 2 — 被动仲裁(reactive):
  内存仲裁器观测到系统内存紧张 → 调用 reclaim()
    └─ 强制 spill 当前算子

两者的分工:

  • 防线 1 发生在算子自己的处理循环中,知道"我即将需要多少内存",可以精确预留
  • 仲裁层可以在其他参与者的分配压力下请求回收,但回收执行仍需要 Task pause、算子可回收区和具体 spill 状态机配合。请求到达的时间与算子实际可回收的时间是两回事。

ReclaimableSectionGuard 是连接两者的桥梁——算子在调用 maybeReserve 前设置它,告诉仲裁器"现在可以来抢我的内存"。这样即使 maybeReserve 本身触发了全局内存整理,算子也能安全地被 spill。


12.9 NoRowContainerSpiller:让接口适应现实,而非让现实适应接口

HashProbe 侧 spill 时,行数据在 RowVectorPtr 里,根本不经过 RowContainer。如果强制所有 spill 都走 RowContainer → extractSpill → 磁盘 的路径,就需要先把 RowVector 插入 RowContainer 再提取出来,白白多一次转换。

NoRowContainerSpiller 的存在体现了一个原则:抽象应该覆盖真实场景,而不是把真实场景削足适履。它直接接受 RowVectorPtr,跳过整个 fill/extract 流程,只保留写文件这一核心功能。

同样的思路也体现在 SortOutputSpiller 上:sortedRows_ 已经是排好序的行指针数组,如果再经过 fillSpillRuns → 按 hash 分区 → sort 这个流程,不仅多余还会破坏已有的顺序。它重写接口,接受外部已排序的 SpillRows,直接写入。


12.10 HashJoinBridge:异步协调的最小化接口

Build 和 Probe 在不同线程运行,需要协调三件事:

  1. Build 完成 → 通知 Probe 开始
  2. Probe 完成 → 通知 Build 恢复 spill partition
  3. Build/Probe 任意一方 spill → 另一方配合

HashJoinBridge 用表发布、probe 等待和 spill 输入领取等接口重复表达每轮交接;取消、peer barrier、共享表缓存和回收还有其他入口。下面三个接口说明常规 spill 轮次,不是类的全部同步能力。

流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

setHashTable(table, spillPartitionSet) // Build → Probe:我建好了,这些 partition 被 spill 了
probeFinished()                        // Probe → Build:我探测完了,请处理下一个 spill partition
appendSpilledHashTablePartitions(set)  // Probe reclaim → Build:我 spill 了你的 hash 表

在表示与配置允许的递归层级中,每轮都可以复用这组事件接口。当前 SpillPartitionId 支持 level 0..3、每层至多 3 个分区位;恢复还受可用 hash 位与配置上限限制。接口复用减少了按层另造同步结构的需要,但不会消除这些边界。

统一 bridge 协议使每轮恢复复用相同状态和事件接口,减少按层手工维护同步对象的复杂度;复杂度是否指数增长取决于另一种实现的结构,不能直接给出这样的必然结论。


12.11 DictionaryVector wrap:零拷贝分区

Probe 侧 spillInput 需要将一个 batch 的行按 partition 分发到多个 spill 文件。最直接的做法是按 partition 拷贝行数据,但这意味着每行数据被读一次、写一次——对于宽行(含字符串列)代价极高。

Velox 的解法是 DictionaryVector:

原始 input:[row0, row1, row2, row3, row4]
partition 0 的 indices buffer:[1, 3]
partition 1 的 indices buffer:[0, 2, 4]

wrap(2, partition0_buffer, input) → DictionaryVector{indices=[1,3], base=input}
  ↓ 序列化时才真正读取 row1, row3 的数据

DictionaryVector 分区保留 base vector 并创建索引,避免在这一分区步骤复制全部 payload。后续 lazy 加载、提取、序列化、压缩和文件缓冲仍可能发生读取或复制,因此它不是整条 spill 路径的“一读一写”保证。

这是"懒求值"(lazy evaluation)思想在 spill 路径中的具体应用。


12.12 小结:设计模式归纳

12.13 代码阅读与验证边界

建议按 SpillerBase::spill → runSpill → extractSpillVector → SpillState → ordered reader 阅读基础设施,再分头看 HashBuild reclaim、GroupingSet spill、SortBuffer spill。

本次核对了这些实现及相关测试入口,没有运行完整 Velox C++ 测试或性能实验。文章中的阈值是源码参数,流程图是依赖与状态示意;它们不能替代指定配置下的执行统计。

回顾整个 spill 系统,可以归纳出几个贯穿始终的设计模式:

① 约束驱动设计(Constraint-Driven Design) 每个设计决策都能追溯到一个具体约束:SpillRun 切分 → PrestoPage 2GB 限制;64 行批次 → L1 cache;DictionaryVector → 零拷贝。没有无目的的"为了灵活性而灵活性"。

② 机制 vs. 策略分离(Mechanism vs. Policy Separation) SpillerBase 提供机制(fill/sort/write 的流程),各子类通过 needSort()、重写 extractSpill() 等钩子注入策略。算子层决定"何时 spill",Spiller 层决定"如何 spill",文件层决定"如何存储"。

③ 路径复用(Path Reuse) HashBuild 恢复 spill partition 时走的是和正常处理输入完全相同的 addInput 路径。HashProbe restore 轮也走同一套 addInput → spillInput → probe 路径。不为"恢复"写特殊逻辑,而是让恢复成为正常路径的自然一部分。这减少了测试覆盖的难度,也降低了特殊路径引入 bug 的风险。

④ 懒求值(Lazy Evaluation) 数据尽可能晚才物化:行指针不拷贝数据(SpillRun)→ DictionaryVector 不拷贝行(probe 分区)→ lazy 列在 spill 前才 loadedVector()。每一层都只在真正需要时才付出代价。

⑤ 尊重生命周期(Lifecycle Awareness) Input Spiller 和 Output Spiller 的分离不是技术必要性,而是对算子语义的尊重:input 阶段和 output 阶段的资源需求、触发频率、分区需求都不同,强行合并会产生大量条件分支,拆开则各自极简。


这套设计的优雅,不在于它的聪明,而在于它的诚实——每个决策都清楚地知道自己在解决什么约束,不做多余的事,也不留下欠债。


13. 以恢复工作集检验 Spill 的设计取舍

Spill 的评价单位应是一次完整的外部执行过程。只看写盘吞吐,会遗漏状态物化的峰值内存、后续恢复需要的组织方式、重复序列化和跨 Driver 的等待。写得快但恢复时重新加载全部数据,并不能解决原来的容量问题。

选择 设计出发点 需要同时计算的成本
有序 run 或无序分区 由恢复时需要 merge 还是重建索引来决定 排序 CPU、分区数量、文件组织和读回工作集
Input / Output Spiller 分开 输入仍在增长与结果已经开始输出具有不同状态约束 两套阶段接口需要一致的游标、清理和异常处理
小批量提取与写出 控制行列转换和序列化产生的临时峰值 批次太小可能增加调用和 I/O 开销,参数需要结合行宽判断
DictionaryVector 包装分区 通过索引复用输入 payload,减少分区阶段的复制 必须保留底层向量生命周期;后续序列化仍会读取和处理值
并行 spill 与异步写出 重叠独立工作,缩短回收执行时间 调度、内存峰值、错误传播与等待所有任务结束的成本

有价值的代码品味是把状态转换和资源责任写进结构中。RAII 清理、捕获跨线程异常、等待已提交任务、明确有序状态、主动归还 reservation,都在维护“失败后仍能安全结束”这一要求。它们不是可以脱离线程与对象生命周期随意复制的语法技巧。

分层也有边界:文件层可以隐藏部分 I/O 细节,但序列化格式与配置仍可能影响读写两侧;统一写出流程可以复用机制,却不能替算子决定排序、匹配标记和输出进度。抽象的价值在于让共性复用,同时使差异有明确的位置,而不是声称替换一个组件就绝不会影响其他层。

源码核对(2026-09-20):本轮按 Velox 1d1b76567870 核对关键接口、控制流、默认值与边界条件。当前源码摘录附固定版本链接;流程伪代码用于说明分支,不是可直接编译的程序。未对全文示例做独立编译或性能复测。涉及宿主集成与历史实验的数据,按各节标注的来源理解。

对照贯穿例子,Spill 的验收问题可以具体到:同 key 是否在同分区、写出的是否为可合并中间态、每条输入是否只贡献一次、恢复时是否只保留受控工作集、异常后是否仍有任务访问已释放容器。文件写成功只满足其中一项。有限 fan-in 能限制读回阶段同时打开的流和缓冲,却增加中间读写;一个极大的聚合状态、耗尽的磁盘或不可回收阶段仍可能失败,因此 spill 扩大了可处理规模,但不保证任意查询永不 OOM。