1. 1. 1. 有序数据合并,为什么仍然会耗尽内存
    1. 1.1. 从小时分区合并到天分区
    2. 1.2. source 的有序契约比文件有序更强
    3. 1.3. 24 路输入、每轮 8 路的完整路径
  2. 2. 2. 算子、队列、游标与文件流各做什么
  3. 3. 3. Pipeline lazy start:先创建连接,再允许读取
    1. 3.1. 为什么启动门放在 CallbackSink
    2. 3.2. 下一组 source 什么时候启动
  4. 4. 4. 一份队列协议,连接启动、数据和结束
    1. 4.1. LocalMergeSource 保存哪些状态
    2. 4.2. 为什么 promise 在锁外通知
  5. 5. 5. 归并内核:比较、游标、批次复制
    1. 5.1. TreeOfLosers 如何决定下一行
    2. 5.2. 不是每次选中一行就立即复制整行
    3. 5.3. 行宽估计怎样影响输出批次
  6. 6. 6. 多轮控制:从当前 source 组切换到最终归并
    1. 6.1. 字段组合表达执行进度
    2. 6.2. 为什么最后一轮也落盘
    3. 6.3. 阻塞 future 怎样回到 Driver
  7. 7. 7. 文件组读回:有序拼接与真正的资源边界
    1. 7.1. 一条 stream 不等于只打开一个文件
  8. 8. 8. Async stream:怎样调度,怎样控制领先量
    1. 8.1. streamIdx 是几组对象之间的对应关系
    2. 8.2. 快路径直接读,满队列后挂 continuation
    3. 8.3. 有界队列约束的是什么
  9. 9. 9. 异步线程的异常怎样回到查询执行线程
    1. 9.1. 错误信息与流结束必须分别表达
    2. 9.2. 同步读取失败与 continuation 失败
    3. 9.3. 一条流出错以后,其余流如何停止
    4. 9.4. 为什么在 atEnd 处统一重抛
  10. 10. 10. 对象所有权、回调寿命与提前取消
    1. 10.1. 谁负责让对象继续存在
    2. 10.2. weak_ptr 不是所有回调都安全的证明
    3. 10.3. 析构和 close 实际做了什么
  11. 11. 11. 启用条件、配置与运行观测
    1. 11.1. 配置 F 之前,先满足 spill 的接入条件
    2. 11.2. 怎样读两个阶段的统计
  12. 12. 12. 源码中的测试验证了哪些约束
  13. 13. 13. 设计取舍:用阶段边界管理工作集
    1. 13.1. 内存模型需要同时计算两阶段
    2. 13.2. 有序 run 是简化系统的关键中间表示
    3. 13.3. 换来的稳定性,付出的延迟与 I/O
    4. 13.4. 错误处理的价值在于可解释的结束协议
  14. 14. 14. 参考文章与源码入口
Macduan Notes

Velox Multi-Round LocalMerge:Lazy Start、异步读回与异常收敛

LocalMerge 是 Velox 在一个 Task 内合并多条有序输入流的算子。上游可以并行扫描和解码,LocalMerge 用一个 Driver 按排序键选出下一行,向下游提供一条有序的 RowVector 流。输入已经有序,所以这里的核心工作是归并,而不是重新排序全部数据。

Multi-round LocalMerge 在这个模型上加入了分组启动、归并落盘、异步读回:一次只启动一部分上游 source,把这一组的归并结果写成有序 spill 文件组;各组完成后,再归并这些文件组。它要控制的是同时活跃的 reader 和输入工作集,代价是一次完整的中间写入与读回,以及更晚的最终首批输出。

这个特性里有三套相互配合的协议:lazy start 决定谁现在可以开始生产;有界队列决定已经开始的 producer 可以领先 consumer 多远;异常与 EOF 协议决定异步工作怎样停止,以及失败怎样回到执行查询的 Driver。先理解它们各自的职责,再看具体类和回调,代码中的 future、指针和状态字段才有明确的位置。

本文围绕我参与实现的特性展开,背景来源于 Multi-Round Lazy Start Merge(官方文章发表于 2025-11-09,作者 Meng Duan、Xiaoxuan Meng、Pedro Pedreira)。本文发布时间设为 2026-01-18;实现细节于 2026-09-20 按 Velox 1d1b76567870 重新核对。下文描述的是这个固定快照,包含后续演进,不代表所有代码在 1 月都已相同。源码节选保留原实现;流程化代码明确标为伪代码。

1. 有序数据合并,为什么仍然会耗尽内存

从小时分区合并到天分区

推荐模型训练数据常按主键分桶,并在桶内按键排序。合并 24 个小时分区时,希望同一个主键的所有行在天分区中连续存放。如果每个参与归并的输入都已按相同规则排序,问题就是多路有序归并:维护每路当前的最小行,从中不断选出全局最小行。LocalMerge 保留所有输入行;它不会因为 key 相同就做去重、聚合或覆盖。

官方文章讨论的是 PB 级训练数据准备场景。Spark 可以利用已有的 bucketing 和 ordering 避免一部分 shuffle,但文中所述的执行路径仍把依次读入的数据交给 sorter,产生中间 spill。原文观察到,那个场景下 row-based spill 相对数据湖中的列式训练数据约有 4 倍空间放大。这是原文特定工作负载的观察,不能当作任意 Spark 作业的固定倍率,也不是本文重新测得的结果。

换成有序归并后,不需要把所有输入行重新装进 sorter;但需要同时看到每条活跃输入流的头部。宽表的 reader 还可能持有列元数据、字典、页缓冲、解压空间和预取结果。上千列乘上多条输入流,即使队列只缓存少量 RowVector,读取侧的内存也可能很大。不同格式和 reader 的按需加载策略不同,不能把它简化成“任何 reader 永远为每列完整保留一页”的常量公式。

source 的有序契约比文件有序更强

一个 MergeSource 表示一条完整的生产流,通常连接一份上游 Driver 和一个队列。它必须保证跨批次的整体有序性,排序键、升降序和 NULL 排序也必须与 LocalMerge 一致。两个文件各自有序,并不意味着把它们随意串在同一条 Driver 上后仍然有序;宿主规划 split、文件和 Driver 时必须保证这个契约。LocalMerge 不会替宿主修复错误的文件拼接顺序。

对象这里的含义容易混淆的边界
source一条整体有序的生产流不是必然等于一个文件、一个 split 或一个线程
pipeline / Driver算子链的结构 / 一份可调度的执行进度Driver 可以换 executor 线程,也可以在 future 上暂停
round / run一次处理一组原始 sources / 这组的有序输出一轮可生成多个 spill 文件,仍然只是一条有序 run
spill file group保存同一条 run 的一个或多个文件最终归并按组计路数,组内按写出次序拼接
async stream一条 BatchStream、队列和回调组成的逻辑流不是一个名为 AsyncStream 的通用 C++ 类,也不是专属线程

24 路输入、每轮 8 路的完整路径

令原始 source 数为 N,每轮上限为 F。N = 24、F = 8 时,先启动 0–7,归并并写出文件组 G0;这一组结束后再启动 8–15,最后启动 16–23。假设三组都非空,最终启动 3 条 spill stream 读回 G0、G1、G2,归并出最终结果。某一组完全没有行就不会生成文件组。

第一阶段的三轮依次执行;第二阶段才并行读回三个有序文件组并向下游输出。
图 1:第一阶段的三轮依次执行;第二阶段才并行读回三个有序文件组并向下游输出。 打开 SVG 原图

这是多个局部归并轮次,加一次最终归并。当前实现不会反复递归,把第二阶段也一直限制成 F 路。N ≤ F 时走普通流式 LocalMerge,无需为了“多轮模式”强制落盘。N > F 时所有原始输入组的结果都要 spill,最后一组同样如此。

2. 算子、队列、游标与文件流各做什么

规划层负责建立跨 pipeline 的连接,Merge 负责阶段与轮次,SourceMerger 负责有序归并,MergeSource 负责同步。SpillMerger 把文件读回转成 MergeSource 所提供的相同接口,因此最终归并可以继续使用 SourceMerger。

模块主要类或函数职责
计划与执行连接LocalMergeNode / LocalPlanner / Task定义排序规则,建立 producer pipeline,登记并共享 MergeSource
生产入口CallbackSink启动门检查、批次入队、生产结束通知
跨线程同步MergeSource / LocalMergeSourcestarted、队列、EOF、producer / consumer promises
算子控制LocalMerge → Merge → SourceOperator选择 source 组、调度 spill、处理 Driver 阻塞和结束
归并内核SourceMerger / SourceStream / TreeOfLosers比较排序键、推进游标、组装输出向量
中间写出MergeSpiller接收已经有序的 RowVector,写出同一 run 的文件
最终读回SpillMerger / BatchStream / SpillReadFile并行读取不同文件组,经队列交给单线程归并
归并侧的主要类关系。TreeOfLosers 拥有 SourceStream;SourceStream 只借用 MergeSource 指针。
图 2:归并侧的主要类关系。TreeOfLosers 拥有 SourceStream;SourceStream 只借用 MergeSource 指针。 打开 SVG 原图

LocalMerge 在下游 pipeline 中是 SourceOperator:它通过跨 pipeline 的共享队列获取输入,而不是由同一条 Driver 中的上一个 Operator 调用它的 addInput。上游 pipeline 尾端的 CallbackSink 才是实际接收扫描输出并入队的位置。Task 按 split group 和 plan node id 保存这些共享 sources。

原文的 24 路扫描画成“两条 pipeline,扫描 pipeline 创建 24 份 Driver,归并 pipeline 创建 1 份 Driver”,适合解释该工作负载。当前 LocalPlanner 还需要按具体计划结构切分不同 source 的 pipeline,不能把这个示例当作任意 LocalMerge 计划都恰好只有两个 DriverFactory 的规定。归并端要求单线程,则同时在 LocalMergeNode::requiresSingleThread 和 LocalMerge 构造检查中体现。

源码:LocalMergeNodemakeMergeSinkSupplierTask::addLocalMergeSourceMerge 与 SourceMerger 声明

3. Pipeline lazy start:先创建连接,再允许读取

为什么启动门放在 CallbackSink

如果上游已经把所有 reader 打开并开始读取,下游即使每次只归并 8 路,也未必能减少扫描侧的峰值内存。多轮归并需要更早的控制点:未被选中的 producer 暂时不要向源头取数。

LocalPlanner::makeMergeSinkSupplier 为每份 CallbackSink 创建一个 LocalMergeSource,并安装两个捕获该 source 的回调。consumerCb 调用 enqueue;startCb 调用 started。producer 知道怎样发送数据,但是否允许启动,由下游调用 source->start() 决定。

Driver 正常推进时先初始化算子,然后从 sink 向 source 检查可执行状态。于是 CallbackSink::isBlocked 能在上游 getOutput 前拦住这条 pipeline。以当前 TableScan 为例,initialize 并不完成全部文件读取工作,dataSource 的创建和 addSplit 发生在后续取 split 的路径中。这个门控能推迟昂贵的读路径,但不是承诺所有 Operator 的 initialize 都无成本,更不是延迟创建所有 Driver 对象。

源码节选:启动检查只做一次,恢复依靠 future 通知 · BlockingReason CallbackSink::isBlocked

BlockingReason CallbackSink::isBlocked(ContinueFuture* future) {
  if (startedCb_ != nullptr) {
    blockingReason_ = startedCb_(&future_);
    startedCb_ = nullptr;
  }
  if (blockingReason_ != BlockingReason::kNotBlocked) {
    *future = std::move(future_);
    blockingReason_ = BlockingReason::kNotBlocked;
    return BlockingReason::kWaitForConsumer;
  }
  return BlockingReason::kNotBlocked;
}

startedCb_ 首次调用后被清空。若 source 尚未启动,LocalMergeSource 注册 producer promise,把 future 交给 CallbackSink;Driver 保存阻塞状态并退出当前执行。它不会占着 executor worker 原地等待,也不需要以后每次轮询 startedCb_。下游 start() 把 started_ 置为 true 并兑现 promise,正常的 Driver continuation 让它再次变成可运行状态。

下游 LocalMerge 打开启动门,上游 Driver 才继续向 TableScan 取数。时序竖线仅表示参与者,不代表专属线程。
图 3:下游 LocalMerge 打开启动门,上游 Driver 才继续向 TableScan 取数。时序竖线仅表示参与者,不代表专属线程。 打开 SVG 原图

下一组 source 什么时候启动

Merge::maybeStartNextMergeSourceGroup 仅在当前 sourceMerger_ 为空、并且还存在未启动 source 时工作。它取 sources_ 中下一段连续索引,建立 SourceStream 和 SourceMerger,再对这组调用 start(),最后推进 numStartedSources_。启动是显式的一次性操作;LocalMergeSource::start 检查不能重复启动。

一组到达 EOF 后,Merge 结束该组的 spill 并丢弃本轮 SourceMerger;下一次执行推进时才能打开下一组。上游已经结束的 Driver 随正常执行路径关闭算子,TableScan::close 会取消 dataSource 的工作。但“本组 EOF”并不是操作系统线程 join,也不保证所有分配立即归零;具体 reader、缓冲和回调的最终释放仍遵循各自所有权。

这里使用的是普通 Driver 阻塞协议,与内存仲裁里 Task 的 pause / suspended 协调不同。把 lazy start 解释成“把整个 Task 暂停”,会看错控制范围:等待的是尚未获准开始的 producer,当前组和下游依然在执行。

源码:Driver::runInternalTableScan::initializeTableScan::close组启动控制

4. 一份队列协议,连接启动、数据和结束

LocalMergeSource 保存哪些状态

LocalMergeSource 的内部队列由 folly::Synchronized 保护,批次存于 boost::circular_buffer<RowVectorPtr>。started_ 表示生产是否获准启动,atEnd_ 表示生产端声明不会再提供数据;producerPromises_ 和 consumerPromises_ 分别承接两个方向的等待。容量以批次数计量,默认是 2,而不是按字节计量。

操作与条件状态变化调用者怎样继续
started:尚未启动保存 producer promise返回 kWaitForConsumer,等待 start
startstarted_ = true,通知 producer已等待的 Driver 可以恢复
enqueue:普通批次先 push,再检查是否变满;通知 consumer变满则返回 future,当前批次已被接受
next:队列非空pop 一批并通知 producer返回这一批,释放一个队列槽
next:空且没有 EOF保存 consumer promise返回 kWaitForProducer
enqueue(nullptr)atEnd_ = true,通知 consumer不占槽,不因满队列而阻塞
next:空且已经 EOF返回空数据且不阻塞consumer 将这条流视为结束

源码节选:批次先入队,满队列只阻止继续生产 · BlockingReason enqueue

BlockingReason enqueue(
        RowVectorPtr input,
        ContinueFuture* future,
        ScopedPromiseNotification& notification) {
      if (!input) {
        atEnd_ = true;
        notifyConsumers(notification);
        return BlockingReason::kNotBlocked;
      }
      VELOX_CHECK(started_);
      VELOX_CHECK(!data_.full(), "LocalMergeSourceQueue is full");

      for (auto& child : input->children()) {
        child->loadedVector();
      }

      data_.push_back(input);
      notifyConsumers(notification);

      if (data_.full()) {
        producerPromises_.emplace_back("LocalMergeSourceQueue::enqueue");
        *future = producerPromises_.back().getSemiFuture();
        return BlockingReason::kWaitForConsumer;
      }
      return BlockingReason::kNotBlocked;
    }

enqueue 返回 kWaitForConsumer 时,调用者不能重试同一个 batch。这个 batch 已经 push 进队列;future 表示的是下一次生产需要等待空位。CallbackSink 保存这份 future 交给 Driver;SpillMerger 则在它上面挂 continuation。两条路径使用同一套背压规则。

queueSize = 2 的往返。EOF 不消耗容量;消费者先排空已有批次,然后才观察结束。
图 4:queueSize = 2 的往返。EOF 不消耗容量;消费者先排空已有批次,然后才观察结束。 打开 SVG 原图

为什么 promise 在锁外通知

每次改变队列状态时,会把应当唤醒的 promises 移入 ScopedPromiseNotification。queue_.withWLock 结束后,notification 析构再兑现 promise。这样避免在持有队列锁时触发后续调度或回调重入,缩小临界区,也降低重新进入同一队列导致锁交互的风险。状态必须先更新,再通知等待者;通知本身不替代受锁保护的状态。

空队列与 EOF 必须分开:一个 source 暂时没有头部行,下游无法判断它的下一行是否小于其他 source 的当前头部。比如当前可见两路是 10 和 20,第三路上一批最后一行是 1,那么下一批可能从 3 开始;直接跳过这路输出 10 就破坏了排序。只有明确结束,才可以不再考虑这路。

当前实现还支持 drain 状态供 MixedUnion 等路径使用;LocalMerge 的 consumerCb 明确拒绝 drained 输入。不能把 drain 当成这里正常 EOF 的替代品。另外,LocalMergeSource::close 当前是空实现,不能据接口名称想象出“close 自动通知全部等待者”的路径。

源码:LocalMergeSource 及其内部队列

5. 归并内核:比较、游标、批次复制

TreeOfLosers 如何决定下一行

SourceStream 包装一条 MergeSource,保存当前 RowVector、当前行位置、排序列指针和待输出位置。SourceMerger 把所有 SourceStream 的 unique_ptr 交给 TreeOfLosers;它自己另存一组裸指针,用于轮询阻塞、估计行宽和复制输出,这些指针不再拥有对象。

败者树把两路比较中的败者留在内部节点,胜者向上推进。选出全局最小行后,只有获胜 stream 的头部发生变化;下一次选择只需沿这条路径重新比较,不必每次重新扫描所有输入头部。构建代价与路数同阶,稳定推进每行的比较次数通常为 O(log k),比较本身的代价还取决于键数量、类型和共同前缀长度。当前 SourceMerger 调用的是 TreeOfLosers::next(),不是供聚合等场景使用的 nextWithEquals()。

SourceStream::operator< 逐个排序键调用 vector 的 compare,并传入 ascending、nullsFirst 等 CompareFlags。第一个非零比较结果决定顺序;所有 key 相等则返回 false。这满足归并的排序需求,但没有承诺 SQL 未声明的相同 key 稳定顺序。需要确定性的 tie-break 顺序,就应当把 tie-break 字段加入排序键。

比较只决定哪一行属于哪个输出位置;payload 复制按列完成,并尊重输入批次的生命周期。
图 5:比较只决定哪一行属于哪个输出位置;payload 复制按列完成,并尊重输入批次的生命周期。 打开 SVG 原图

不是每次选中一行就立即复制整行

每次 winner 产生后,setOutputRow 记录它在输出批次中的位置。对同一 source 而言,被消耗的输入行始终按顺序推进,但对应输出位置可能间隔着其他 source 的行。copyToOutput 由这些选择位置生成 source row indices,再逐列调用 child vector 的 copy。这样把排序控制与列向量复制分离。

有两个不能拖延的复制时机。其一,当前 source 批次即将耗尽,必须在 pop 触发 fetchMoreData 替换 data_ 前复制所有已选行。其二,输出批次达到行数上限,要把所有 source 尚未复制的选择一起写出。SourceStream::pop 中“没有剩余 selection”的检查,正是在保护前一种生命周期边界。

如果输出只攒了一部分行,而某路取下一批时阻塞,getOutput 可以返回 nullptr,同时保留 output_ 和 outputRows_。空返回只是暂时不能继续形成完整结果,不等于已经丢弃部分输出或算子完成。外层仍需配合 isBlocked 和 isFinished 判断执行状态。

createOutputVector 优先基于第一个有效输入调用 createEmptyLike<RowVector>,使子列编码能支持 FlatMapVector 等输入;没有模板时才按 RowType 创建。这里仍然会复制所选 payload,不能把“参考输入编码建输出”写成整个归并是零拷贝。

行宽估计怎样影响输出批次

源码节选:每个 SourceMerger 首次决定输出行数 · void SourceMerger::setOutputBatchSize

void SourceMerger::setOutputBatchSize() {
  if (outputBatchRows_ != 0) {
    return;
  }
  size_t numEstimations{0};
  int64_t estimateRowSizeSum{0};
  for (auto* stream : streams_) {
    const auto estimateRowSize = stream->estimateRowSize();
    if (estimateRowSize.has_value()) {
      ++numEstimations;
      estimateRowSizeSum += estimateRowSize.value();
    }
  }

  if (numEstimations == 0) {
    outputBatchRows_ = maxOutputBatchRows_;
    return;
  }

  const auto estimateRowSize =
      std::max<vector_size_t>(1, estimateRowSizeSum / numEstimations);
  outputBatchRows_ = std::min<vector_size_t>(
      std::max<vector_size_t>(1, maxOutputBatchBytes_ / estimateRowSize),
      maxOutputBatchRows_);
}

各有效 stream 用当前批次的 estimateFlatSize()/size 估计每行大小,再对这些估计取按 stream 等权的算术平均,不是把全部批次字节与行数加总后的加权平均。最终行数为 min(R, max(1, B / max(1, 平均行宽))),其中 B 是 preferred_output_batch_bytes,R 是构造时传入的 maxOutputBatchRows。没有有效估计时直接用 R。

这个决定只在 outputBatchRows_ 为 0 时做一次。因此是每个 SourceMerger 初始批次驱动的估计:新一轮和最终读回会分别创建新的 SourceMerger,但并非每个输出 batch 都重新自适应。变长字段、编码变化和后续行宽偏斜都可能让实际分配超过估计字节目标;B 不是硬内存配额。

还有一个参数命名容易误导:当前 Merge 构造使用无参数 outputBatchRows() 初始化 maxOutputBatchRows_,这条分支返回 preferred_output_batch_rows,默认 1024。这里没有传 averageRowSize,也就不会走 Operator::outputBatchRows 中依据 max_output_batch_rows 计算上限的另一条分支。换句话说,当前 LocalMerge 的 R 默认是 1024,而不是看到字段名就推测出的全局 10000 行上限。

源码:TreeOfLosersSourceStream 比较copyToOutputcreateOutputVectoroutputBatchRows

6. 多轮控制:从当前 source 组切换到最终归并

字段组合表达执行进度

字段含义关键变化
sources_ / numStartedSources_全部原始输入 / 已启动前缀的长度sources_ 保留原始连接;计数按组推进
sourceMerger_当前原始输入组的归并状态一组一个实例,组结束后清空
mergeOutputSpiller_当前组的输出 writer首次出现输出才创建,finishSpill 后清空
spillFileGroups_已经完成的非空有序 runs每结束一个非空组追加一项
spillMerger_第二阶段的整体读回管理器所有原始组完成后创建一次
sourceBlockingFutures_ / finished_待处理的阻塞通知 / 算子完成标记等待与结束分别表达
Merge 的解释性状态图。是否需要 spill 由总输入规模决定,而不是每一轮临时决定。
图 6:Merge 的解释性状态图。是否需要 spill 由总输入规模决定,而不是每一轮临时决定。 打开 SVG 原图

流程伪代码:展示阶段切换;省略统计、检查和移动语义

// Control-flow pseudocode, not a compilable replacement for Merge.
if (spillMerger_) {
  return readFinalMerge();
}

auto batch = sourceMerger_->getOutput(futures, atEnd);
if (maxNumMergeSources_ < sources_.size()) {
  spillIfNonEmpty(batch);  // also applies to the last source group
  batch = nullptr;
}

if (atEnd) {
  finishCurrentGroup();
  if (numStartedSources_ == sources_.size()) {
    if (numSpilledRows_ > 0) {
      setupSpillMerger();
    } else {
      finished_ = true;
    }
  }
}
return batch;

为什么最后一轮也落盘

needSpill() 的实现是 maxNumMergeSources_ < sources_.size()。它在整个原始输入阶段保持同一判断。假设前两组已经落盘,最后一组只有几条 source,也不能直接把最后一组的结果作为最终输出:那些已经落盘的组可能还有更小的 key。当前实现选择统一的阶段边界,让最后一组也形成完整 run,然后再对所有 runs 做最终归并。

MergeSpiller 接收已经排好序的 RowVector。它继承 NoRowContainerSpiller,不需要先进入 RowContainer 再做一次排序。当前路径用 SpillPartitionId{0} 写同一个 partition;一个组结束时调用 finishSpill,取出该 partition 的文件列表,追加到 spillFileGroups_。这里的 partition 是 spill 管理单元,不是在排序键上重新 hash 分桶。

如果某一组全空,就不会创建 mergeOutputSpiller_,也就不会产生空文件组。正常多轮路径里每个输入行写入一次中间 spill;一个文件组可能因文件大小控制拆成多个文件,所以“行数只 spill 一次”和“每组恰好一个文件”是两件不同的事。

阻塞 future 怎样回到 Driver

SourceMerger::isBlocked 在当前 future 列表为空时检查各个 stream,收集取数所需的 futures。Merge::isBlocked 每次从列表末尾拿出一个 future,返回 kWaitForProducer;恢复执行后继续处理剩余等待。这不是一个显式 collectAll,也不能理解成“任意一路变可读就一定能输出”。为了选出下一行,归并仍必须知道所有未结束流的有效头部。

第二阶段的 futures 由 SpillMerger::getOutput 经同一个 sourceBlockingFutures_ 容器交回外层。上层 Driver 并不需要知道这次缺数据是因为上游扫描尚未产出,还是 spill worker 尚未把读回批次放入队列。

源码:getOutputFromSourcefinishMergeSourceGroupMerge::isBlockedMergeSpiller

7. 文件组读回:有序拼接与真正的资源边界

一轮的归并输出本身全局有序。写满一个 spill 文件再写下一个,只是把同一条有序 run 分段,因此文件组内部可以直接按原次序拼接;不同组的 key 范围可能重叠,才需要在组之间做有序归并。

文件流模块的类关系及数据方向。BatchStream 的同步接口与 SpillMerger 的异步调度分处两层。
图 7:文件流模块的类关系及数据方向。BatchStream 的同步接口与 SpillMerger 的异步调度分处两层。 打开 SVG 原图

SpillMerger 为每组建立一个 ConcatFilesSpillBatchStream,它实现 BatchStream::nextBatch(RowVectorPtr&)。内部 fileIndex_ 指向当前文件;成功读到一个 batch 就返回 true;当前文件结束后 reset 对应 SpillReadFile,再前进到下一个文件;整个组读完后清空文件列表,设置 atEnd_ 并返回 false。调用者必须顺序推进同一条 stream,不能并发调用这个非线程安全对象。

源码节选:顺序拼接一个组内的文件 · bool ConcatFilesSpillBatchStream::nextBatch

bool ConcatFilesSpillBatchStream::nextBatch(RowVectorPtr& batch) {
  TestValue::adjust(
      "facebook::velox::exec::ConcatFilesSpillBatchStream::nextBatch", nullptr);
  VELOX_CHECK_NULL(batch);
  VELOX_CHECK(!atEnd_);
  for (; fileIndex_ < spillFiles_.size(); ++fileIndex_) {
    VELOX_CHECK_NOT_NULL(spillFiles_[fileIndex_]);
    if (spillFiles_[fileIndex_]->nextBatch(batch)) {
      VELOX_CHECK_NOT_NULL(batch);
      return true;
    }
    spillFiles_[fileIndex_].reset();
  }
  spillFiles_.clear();
  atEnd_ = true;
  return false;
}

SpillReadFile 继承 SerializedPageFileReader,当前使用 Presto vector serde,并配置 timestamp、compression 和 nullsFirst 相关读取选项。反序列化重新产生 RowVector,不需要把中间文件重新交给 Parquet / ORC reader。这将复杂的原始列式 reader 工作集换成 spill 文件读取工作集,但仍有 I/O buffer、反序列化和输出向量成本。

一条 stream 不等于只打开一个文件

这里有一处不能仅凭 ConcatFiles 的名字推断:Merge::setupSpillMerger 会遍历所有组的所有文件,提前创建 SpillReadFile。其基类构造函数立即打开文件并创建 FileInputStream;当前 FileInputStream 构造会分配缓冲并执行首段读取,支持异步预读时还会分配第二个缓冲。

因此,“每组按序调用 nextBatch”并不等于“每组只初始化当前一个文件”。最终归并有 G 路,但若共有 K 个 spill 文件,初始化期间的文件句柄、首读和读缓冲成本与 K 也有关;不能只按 G 乘以一个 buffer 估计内存。文件耗尽后 reset reader 会逐渐回收这些资源。把 reader 创建延迟到实际切换文件,是可以讨论的优化方向,当前快照并不是这样实现的。

另一个边界是初始化错误:打开文件、分配 buffer 或首读就失败时,异常可能在 setupSpillMerger 的 Driver 调用栈上直接抛出,尚未进入异步 worker。后面讨论的异常汇总协议,针对的是已经启动的 spill stream 读回工作,不能声称所有文件错误都必定先经过 setError。

代码中还存在 ConcatFilesSpillMergeStream,它是另一种 SpillMergeStream 游标;本特性的 async 读回路径使用 ConcatFilesSpillBatchStream,经队列后再使用 SourceStream。也不要把它与通用 AsyncSource<T> 混为一谈,这条路径没有通过那个类实现异步化。

源码:setupSpillMergerConcatFilesSpillBatchStreamSpillReadFile 构造SerializedPageFileReaderFileInputStream

8. Async stream:怎样调度,怎样控制领先量

streamIdx 是几组对象之间的对应关系

SpillMerger 保存长度相同的 batchStreams_ 和 sources_:batchStreams_[i] 负责第 i 个文件组,sources_[i] 是它输出到的 LocalMergeSource。内部 SourceMerger 消费这些 sources,读文件回调通过 streamIdx 找到对应对象。sources 在创建后全部 start,因为第二阶段要看到每个文件组的头部,才能开始最终归并。

start() 检查 spill executor 非空,然后为每条 stream 提交一份初始工作。executor 的 worker 数量和 stream 数量没有一一对应关系:例如有 12 条 stream、4 个 worker,待运行回调共享这 4 个线程;同一条 stream 恢复后也可以由另一个 worker 执行。

异步调度解耦了多路读取和单线程归并;每路仍然串行调用自己的 BatchStream。
图 8:异步调度解耦了多路读取和单线程归并;每路仍然串行调用自己的 BatchStream。 打开 SVG 原图

快路径直接读,满队列后挂 continuation

源码节选:完整的单 stream 读回回调,包含异常路径 · void SpillMerger::readFromSpillFileStream

void SpillMerger::readFromSpillFileStream(
    const std::weak_ptr<SpillMerger>& mergeHolder,
    size_t streamIdx) {
  TestValue::adjust(
      "facebook::velox::exec::SpillMerger::readFromSpillFileStream", nullptr);
  const auto merger = mergeHolder.lock();
  if (merger == nullptr) {
    LOG(ERROR) << "SpillMerger is destroyed, abandon reading from batch stream";
    return;
  }

  try {
    if (hasError()) {
      finishSource(streamIdx);
      return;
    }

    RowVectorPtr vector;
    if (!batchStreams_[streamIdx]->nextBatch(vector)) {
      VELOX_CHECK_NULL(vector);
      finishSource(streamIdx);
      return;
    }

    ContinueFuture future{ContinueFuture::makeEmpty()};
    const auto blockingReason =
        sources_[streamIdx]->enqueue(std::move(vector), &future);
    if (blockingReason == BlockingReason::kNotBlocked) {
      VELOX_CHECK(!future.valid());
      readFromSpillFileStream(mergeHolder, streamIdx);
    } else {
      VELOX_CHECK(future.valid());
      std::move(future)
          .via(executor_)
          .thenValue([this, mergeHolder, streamIdx](auto&&) {
            readFromSpillFileStream(mergeHolder, streamIdx);
          })
          .thenError(
              folly::tag_t<std::exception>{},
              [this, mergeHolder, streamIdx](const std::exception& e) {
                const auto merger = mergeHolder.lock();
                if (merger != nullptr) {
                  LOG(ERROR) << "Stop the " << streamIdx
                             << " th source on error: " << e.what();
                  setError(std::make_exception_ptr(e));
                  finishSource(streamIdx);
                }
              });
    }
  } catch (const std::exception& e) {
    LOG(ERROR) << "The " << streamIdx
               << " spill stream failed with error: " << e.what();
    setError(std::current_exception());
    finishSource(streamIdx);
  }
}

回调成功读出一个 batch 后调用 enqueue。若队列没有满,当前实现直接递归调用 readFromSpillFileStream,继续在当前线程读取;不会每个 batch 都重新 executor->add。若队列变满,返回的 future 上挂 via(executor_).thenValue(...),当前回调便可以退出;consumer 取走数据并通知 producer 后,再调度下一次读取。

所以 async 指的是不同 stream 可以并行,以及等待队列空位时让出 worker。BatchStream::nextBatch 本身是同步 API:读文件、等待底层预取完成、反序列化都会占用执行该回调的 worker。底层 FileInputStream 可以另有 read-ahead,这与 SpillMerger 这一层的调度并不冲突,也不是同一件事。

每条 stream 在“当前读取”与“等待背压 future”之间转换,逻辑上只有一个推进链,所以不需要把非线程安全的 BatchStream 变成共享并发容器。队列锁保护跨线程传递;stream 的顺序性保护 fileIndex_ 和 reader 游标。两个层次各管一段状态,职责比较清楚。

有界队列约束的是什么

q 个队列槽只限制已入队但尚未取走的批次。SourceStream 还持有从队列取走的当前批次,producer 读取时可能另有正在构造的向量,SourceMerger 还持有输出批次。因此内存估算不能只算 q × batchBytes,更不能由 q = 2 推导出每路最多只会存在两份 RowVector。

增大 q 可以吸收读延迟抖动,让扫描或解码领先一点;但 consumer 是单线程归并,增加缓存不会消除比较与输出复制瓶颈。对宽表来说,单个 batch 可能很大,增加一格队列的成本会乘以所有活跃 stream。

同步递归快路径减少了调度开销,但递归深度不严格受 q 限制:consumer 足够快时,producer 可以持续看到“队列未满”。若继续演进,可以把快路径改成显式循环,并设置合理的让出策略。这里是对实现取舍的分析,不是已经测得栈溢出,也不是当前已有的 fairness 保证。

源码:初始调度与后续 continuation

9. 异步线程的异常怎样回到查询执行线程

错误信息与流结束必须分别表达

worker 里直接 throw,不会自动在正在运行 LocalMerge 的 Driver 栈上出现。另一方面,只保存 exception_ptr 而不结束对应 source,也会使 consumer 继续等待一个永远不会再生产数据的流。当前实现把两件事配对:保存错误,再给这条流发送 EOF

SpillMerger 用 mutex_ 保护共享 exception_。第一个成功写入这个槽的异常被保留,后续错误不覆盖。这里的“第一个”是锁与写入顺序,不是跨线程的绝对发生时间顺序。各 worker 仍可能记录自己的错误日志,但最终重抛的是槽中保留的那一个。

源码节选:保留第一个写入的异常 · void SpillMerger::setError

void SpillMerger::setError(const std::exception_ptr& exception) {
  std::lock_guard l(mutex_);
  if (exception_ != nullptr) {
    return;
  }
  exception_ = exception;
}

同步读取失败与 continuation 失败

位置当前捕获方式语义边界
nextBatch / enqueue 等 try 内同步调用catch (const std::exception&);setError(std::current_exception())保留当前被捕获异常的动态类型和载荷,然后 finishSource
future continuation 链上的匹配错误thenError(folly::tag_t<std::exception>{});make_exception_ptr(e)e 的静态类型是 std::exception,复制时可能切片,不能保证派生异常信息完整保留
其他 stream 再次进入回调hasError() 后直接 finishSource协作停止后续读取,不覆盖共享错误
reader 创建、打开文件或初始提交阶段失败在启动调用栈上的同步异常不必经过 worker 的共享异常槽
非 std::exception 的 throw不是这两条 catch/tag 的覆盖范围不能将这里描述成兜住任意 C++ 异常的万能边界

std::current_exception() 和 std::make_exception_ptr(e) 在这里并不等价。前者保存正在处理的异常对象;后者根据表达式静态类型创建异常副本。当前 thenError 参数是 const std::exception&,因此可能丢失派生类型和自定义消息。正文保留这个实现事实,不把两条错误路径统一写成“完整保留原始异常”。若要增强错误信息,可以考虑保持 exception_wrapper 或 exception_ptr 的原始动态异常信息,但那属于后续代码改进。

一条流出错以后,其余流如何停止

错误先记录到共享槽,再通过每个 source 的 EOF 收敛;归并端到达结束时统一重抛。箭头中的 EOF 实际经各自队列传递。
图 9:错误先记录到共享槽,再通过每个 source 的 EOF 收敛;归并端到达结束时统一重抛。箭头中的 EOF 实际经各自队列传递。 打开 SVG 原图

失败的回调先 setError,再 finishSource(streamIdx)。finishSource 调用 enqueue(nullptr),并检查没有产生有效等待 future。EOF 不占队列容量,因而不需要先为一个结束标记等待空槽。其他 worker 下次进入读回回调时先检查 hasError,发现已有错误就不再读取下一批,同样发送自己那路的 EOF。

如果另一个 worker 已在 nextBatch 内执行,它不能被这个布尔检查立即打断;如果它因满队列而挂起,需要 consumer 继续取数据,触发 continuation,才会再次看到共享错误。因此当前协议是 cooperative 收敛。慢 I/O、调度停滞或尚未恢复的 continuation 都会影响错误传播延迟,不能理解成某路一失败所有 worker 就立刻消失。

各队列已有的 RowVector 不会因为 atEnd_ 置位立即丢掉。consumer 先排空它们,SourceMerger 才能看到全部 source 的终态。这也意味着错误发生前后可能已有一个输出前缀交给下游;LocalMerge 不提供回滚,宿主必须把整个查询的最终失败与结果提交语义衔接起来。

为什么在 atEnd 处统一重抛

SpillMerger::getOutput 先处理 SourceMerger 的阻塞,再取归并输出;只有 atEnd 为 true 时调用 checkError。若存在错误,checkError 先 reset SourceMerger、清空 batchStreams_ 和 sources_,再 std::rethrow_exception。这样把读回错误最终交回运行 Merge 的 Driver,由查询执行层处理失败,而不是让读回 worker 各自独立地决定算子结束。

源码节选:归并结束处清理并重抛 · void SpillMerger::checkError

void SpillMerger::checkError() {
  if (hasError()) {
    sourceMerger_.reset();
    batchStreams_.clear();
    sources_.clear();
    std::rethrow_exception(exception_);
  }
}

EOF 在这里提供“每路不会继续正常生产”的协议边界,但它不等价于 executor->join():一个回调可能已经发完 EOF,仍在退出函数;future 链和局部 shared_ptr 也可能尚未析构。理解这一点,才能把正常错误收敛与下一节的对象生命周期边界区分开。

源码:SpillMerger::getOutput 与 atEnd 检查

10. 对象所有权、回调寿命与提前取消

谁负责让对象继续存在

关系持有方式它解决的事
Merge → SpillMergershared_ptr让最终归并管理器可以在执行回调期间保活
SpillMerger → MemoryPool / SpillStatsshared_ptr异步工作仍可依赖池对象和统计对象的生命周期
SpillMerger → BatchStreamvector<unique_ptr<…>>集中管理每个文件组的读取状态
SpillMerger → MergeSourcevector<shared_ptr<…>>持有各路队列,让 consumer 与 producer 共享连接
SpillMerger → SourceMergerunique_ptr唯一拥有最终归并状态;其 cursor 借用 sources
SpillMerger → executorfolly::Executor*非拥有引用,executor 必须由外部正确管理寿命
读回回调 → SpillMergerweak holder,成功 lock 后局部 strong holder避免后续等待链无限保活,并覆盖成功 lock 之后的这次调用

readFromSpillFileStream 在进入主体前锁定 weak_ptr,得到局部 shared_ptr merger。只要这次调用已经成功获得强引用,函数执行期间 SpillMerger 就不会因外部最后一个引用释放而析构。pool_ 也用 shared_ptr 保存,避免把一个裸池指针交给异步读回后完全失去对象寿命支撑。

weak_ptr 不是所有回调都安全的证明

源码节选:初始提交实际捕获 this,并在执行时创建 weak holder · void SpillMerger::scheduleAsyncSpillFileStreamReads

void SpillMerger::scheduleAsyncSpillFileStreamReads() {
  VELOX_CHECK_EQ(batchStreams_.size(), sources_.size());
  for (auto i = 0; i < batchStreams_.size(); ++i) {
    executor_->add([&, streamIdx = i]() {
      readFromSpillFileStream(std::weak_ptr(shared_from_this()), streamIdx);
    });
  }
}

必须对照代码看完整捕获链。初始 executor 回调使用 [&, streamIdx = i],调用成员方法意味着它持有的是隐式捕获的原始 this;直到回调开始执行,才通过 shared_from_this() 建立 weak_ptr。不能把它画成“提交前已安全捕获 weak_ptr”。后续 thenValue / thenError 也显式捕获 this;弱引用保护的关键位置,是成功 lock 后形成的局部强引用。

因此,weak holder 体现了避免长期强持有的设计意图,但单凭它不能证明所有“已提交但尚未执行”的初始回调和提前销毁交错都安全。外部 executor 的寿命、operator 关闭顺序、Task 取消和未执行回调,需要一起审视。这里指出的是当前实现的证明边界,不是宣称已经复现了一个取消崩溃。

析构和 close 实际做了什么

SpillMerger 析构按顺序 reset sourceMerger_、clear batchStreams_、clear sources_,没有 join executor。Merge::close 记录统计、调用原始 sources 的 close,再调用 Operator::close;它也不等于立即清空全部 spill 成员。前面已经看到,当前 LocalMergeSource::close 是空方法,不能靠调用它推断所有 producer promises 都已经被主动兑现。

底层资源还有自己的收尾:当前 FileInputStream 析构在存在 read-ahead future 时会等待它完成,并处理预取异常。因而“SpillMerger 没有 join”不意味着销毁任何 reader 都完全不可能等待;只是它没有一个统一的、显式收拢所有 stream 回调的 executor join 屏障。

如果把这部分继续设计得更明确,可考虑提交时就创建合适的 weak / strong handle、使用显式活动回调计数或完成 future、把关闭队列与唤醒等待者作为可证明的协议,并规定 executor 必须活到哪个阶段。这些是可以比较的改进方向,本文不把它们写成当前已经存在的机制。

源码:SpillMerger 所有权字段析构顺序Merge::close底层预取收尾

11. 启用条件、配置与运行观测

配置 F 之前,先满足 spill 的接入条件

LocalMergeNode::canSpill 检查排序键非空和 local_merge_spill_enabled。DriverCtx::makeSpillConfig 还检查全局 spill_enabled,以及是否有 spill 目录或创建目录回调。LocalMerge 构造最终只有在得到 spillConfig 且其中 executor 非空时,才读取并应用 local_merge_max_num_merge_sources。

也就是说,缺少 spill executor 时,普通 LocalMerge 构造路径会保留默认的无限 fan-in,回到普通归并;不是仅设置一个 F 就必然开始多轮,也不是这个入口必然直接报错。另一方面,直接创建并 start 一个 SpillMerger 而没有 executor,则会触发检查失败。这两个入口的行为要区分。

配置或接入项当前默认值在本特性中的作用
spill_enabledfalse全局允许 spill;还需要宿主提供 spill 目录
local_merge_spill_enabledfalse允许有排序键的 LocalMerge 创建 spill 配置
local_merge_max_num_merge_sourcesuint32_t 最大值第一阶段每组的 F;要求大于 0,1 是合法值
local_merge_source_queue_size2每路可排队的 RowVector 批次数,两个阶段都会使用
preferred_output_batch_rows1024当前 Merge 构造传入的输出行数上限 R
preferred_output_batch_bytes10 MiB结合首批估计行宽缩小输出行数,不是硬内存限额
max_output_batch_rows10000其他 outputBatchRows 分支使用;当前 Merge 无参数调用不读取此上限
max_spill_bytes100 GiB;0 表示不设此限制查询级累计 spill 限制,不是 LocalMerge 专属预算
spill executor / spill directory宿主提供异步读回调度资源 / 中间文件存储位置

C++ 配置示例:键值传给宿主使用的 QueryConfig 构造路径;不是完整可运行查询

// Query configuration example; executor and spill directory are host setup.
std::unordered_map<std::string, std::string> config{
    {"spill_enabled", "true"},
    {"local_merge_spill_enabled", "true"},
    {"local_merge_max_num_merge_sources", "8"},
    {"local_merge_source_queue_size", "2"},
    {"preferred_output_batch_rows", "1024"},
    {"preferred_output_batch_bytes", "10485760"},
};

F = 1 虽然合法,却会把每条原始 source 单独写为一个 run,最终仍可能同时面对 N 条 spill stream。F ≥ N 则不发生中间 spill。因此 F 是两阶段工作集之间的权衡参数,不是越小就一定越省总内存。配置生效后应通过实际 spilledRows、阶段统计和计划中的 source 数验证行为,不能只看配置已经写入。

怎样读两个阶段的统计

streamingSourceReadWallNanos 记录原始输入阶段的墙钟跨度;发生多轮时,这段还包括归并输出写 spill 的时间。spilledSourceReadWallNanos 记录读回 spill 并产生最终输出的阶段跨度。它们不是逐个 worker CPU 时间之和,也不是纯磁盘 I/O 时间,里面可能包含等待、调度和归并开销。

分析时应一起观察输入行数、spilledRows / spilledBytes / spilledFiles、读写 I/O 与序列化统计、峰值内存和 executor 使用情况。阶段一很长,可能是扫描或写出慢;阶段二很长,可能是 spill 读回、解码或单线程比较复制慢。先用阶段指标缩小范围,再结合对应资源判断,不能从一个 wall-time 名字直接归因。

源码:LocalMerge 配置makeSpillConfigLocalMerge 构造阶段统计

12. 源码中的测试验证了哪些约束

这次文章核对阅读了 MergeTest 和 MergerTest 中的相关用例,下面说明现有测试的验证目标。本文没有重新编译或运行 Velox C++ 测试,也没有执行文中的配置示例;代码块中的原样节选由固定快照提取,伪代码用于解释控制流。博客构建、图形渲染和页面检查单独进行,不能替代执行引擎测试。

测试入口验证点解读边界
MergeTest.localMergeSpillBasic / localMergeSpill有序结果、多种 fan-in、spill 行数与文件统计特定数据大小下每组一个文件,不能外推所有组永远只有一个文件
localMergeSpill 的辅助测试逻辑F ≥ N、空输入、缺少 executor 等分支不满足多轮条件时可退回普通归并,spill 计数应为零
MergeTest.localMergeSpillPartialEmpty部分 source 为空的多轮行为空输入不能破坏结果或强制产生空文件组
MergeTest.localMergeSpillWithException在不同 nextBatch 调用位置注入异常验证相应读取异常能使查询失败并带回消息
MergeTest.localMergeSmallBatch批末取数阻塞,保留部分输出状态返回 nullptr 不应丢失已经选择的输出行
MergeTest.localMergeStart阻塞 start 时,Values::getOutput 尚未执行直接证明启动门早于上游读取
MergeTest.localMergeAbort归并端抛错,同时异步 reader 延迟覆盖一种提前失败交错;测试最后显式 join spill executor
MergerTest.sourceMerger / sourceMergerWithEmptySources / spillMerger归并内核及多组数、批大小、队列容量组合分层验证数据顺序与结束行为
MergerTest.spillMergerException读回流故障注入不能据此证明所有 continuation 错误均保留派生异常类型

这些测试中若干故障注入和启动时序用例处于 DEBUG 条件下。它们是理解协议的很好入口,但不能将“有 abort 测试”扩展成所有 raw this 初始回调、executor 销毁、取消交错都已经得到形式化保证。读实现时需要把测试覆盖到的具体时序与未覆盖的边界分开。

测试源码:MergeTestMergerTest

13. 设计取舍:用阶段边界管理工作集

内存模型需要同时计算两阶段

把单个原始 source 的 reader 工作集记为 R,把平均输入批次占用记为 V,把 queueSize 记为 q。第一阶段活跃 source 数至多 F,其主要数据工作集可以粗略写成 F × (R + (q + 1) × V),再加上 producer 正在生成的向量、输出批次、spill 写缓冲、已创建对象和其他 Task 状态。这里的 +1 是 consumer 从队列取走后仍保留的当前批次。

第二阶段令非空组数为 G,G ≤ ceil(N / F),总 spill 文件数为 K。主要项变成 G × (q + 1) × Vspill,加上 K 个 reader 的构造与缓冲成本、活跃反序列化工作集、输出和其他状态。底层预取可能使某文件有两个 buffer,bufferSize 还会受文件大小约束,所以这依然是分析框架,不是可直接承诺的精确峰值公式。

如果原始 reader 很重、spill reader 相对轻,分轮通常能显著降低第一阶段同时存在的昂贵读状态;但减小 F 会增大 G,还可能增大中间文件数与调度开销。最终阶段不再受 F 限制,提前创建 K 个 reader 又增加了一项不能忽略的成本。选择 F 时必须同时估计两个阶段,结合宽度、行宽偏斜、文件大小和 executor 能力测量。

有序 run 是简化系统的关键中间表示

第一阶段将复杂的扫描生命周期压缩成“这一组已完成的一条有序 run”。到第二阶段,只需知道如何顺序读每组并比较组间头部,SourceMerger 不必理解原始文件格式,也不需要了解哪些 Driver 尚未开始。文件组保留排序性质,队列统一同步接口,这两个边界共同降低了模块之间的耦合。

把归并保持单线程,也明确了所有权:败者树、cursor 和输出 selection 在一个执行上下文中推进,异步 worker 只负责向受锁保护的队列生产批次。这比让多个线程共同修改 winner 树与输出位置容易推理。代价是最终的比较与 payload 复制仍受单线程吞吐限制,增加读取线程只能隐藏一部分读延迟,不能无限提升整体吞吐。

换来的稳定性,付出的延迟与 I/O

N > F 的路径会把全部输入行写入中间 spill,再读回来;与可直接流式输出的普通 LocalMerge 相比,它增加了序列化、磁盘写读和额外归并。这种选择适合原始输入已经有序、reader 工作集偏大、吞吐和完成稳定性比低首行延迟更重要的批处理。它不能被宣传成任何数据规模上都会更快的优化。

第一阶段所有组完成后才产生最终输出,这个阻塞边界也限制了下游 LIMIT 可能带来的早停收益。一个可讨论的替代是让最后一组实时输出与旧 spill groups 一起归并,但那会让原始 reader 与 spill reader 同时活跃,重新增加内存与生命周期组合。当前实现选择统一写完再读回,控制流更规整,成本也更明确。

错误处理的价值在于可解释的结束协议

这套异步设计值得保留的部分,是不让异常只停留在某条线程的局部日志中:共享槽保存失败信息,EOF 推进所有流的可观察结束,最终在消费端重抛。代价是错误传播需要等待协作收敛,无法天然提供快速取消、事务回滚或统一 join。讨论代码品味时,应同时看到清楚的模块边界和仍需补强的生命周期细节。

未来可以沿几个明确方向演进:根据可用内存和实际 reader 工作集动态选择 F;在最终组数很大时引入额外归并层级;延迟创建组内后续文件 reader;保持完整动态异常信息;以显式完成协议覆盖取消和初始排队回调。官方原文把动态增大 fan-in 列为 Future Work;本文核对的快照仍使用配置中的固定 F,不能把这项愿景写成已落地行为。

Multi-round LocalMerge 的核心,是承认“有序归并只需要少量当前行”并不意味着“同时打开任意多条输入都便宜”。lazy start 把读取工作集控制前移,spill run 把阶段之间的状态固化,async stream 在有界缓存内重叠读取与计算。它们共同形成了这条执行路径;任何一层的结束、背压或所有权语义不清楚,都会影响另一层能否正确推进。

14. 参考文章与源码入口

源码节选来自 Apache-2.0 许可的 Velox 项目,原始版权和许可见各链接文件头及 项目 LICENSE。本文的流程图、类图和时序图按当前实现重新绘制。