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

Velox Multi-Round LocalMerge

把 24 个小时分区合并为一个天分区时,如果训练数据已经按主键分桶并排序,就应当直接合并这些有序输入,避免把同一批数据再完整排序一遍。官方文章讨论的推荐训练数据可达到 PB 级:即使 Spark 利用已有分桶避免了部分 shuffle,顺序读取各分区后再进入 sorter 的路径仍会产生大量中间 spill;原文还观察到该工作负载中的 row-based spill 相对湖中列式数据约有 4 倍空间放大。

Velox 的 LocalMerge 可以直接归并有序流,但“同时开始读 24 路”又带来另一个问题:宽表的 Parquet reader 会持有元数据、字典、页和解压缓冲,再加上各路当前批次,扫描端就可能占用很大的内存。做 Multi-Round LocalMerge,是为了把一次同时活跃的昂贵 reader 数量限制下来:分组启动扫描,每组归并成有序 spill run,全部完成后再异步读回这些较轻的中间流。下面先说明扫描、队列与归并的连接,再展开轮次控制和异步协议。

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 月都已相同。源码节选保留原实现;流程化代码明确标为伪代码。

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

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

1.1 从小时分区合并到天分区

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

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

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

1.2 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++ 类,也不是专属线程

1.3 路输入、每轮 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。

2.1 扫描 pipeline 与归并 pipeline 怎样相接

以“一个有序 Parquet 文件对应一条生产流”为例,每份 producer Driver 内有 TableScan → CallbackSink:TableScan 是读取端,CallbackSink 是同一 pipeline 的尾端。沿逻辑数据流看,CallbackSink 在 TableScan 下游、LocalMerge 上游;但 CallbackSink 与 LocalMerge 分处不同 pipeline,不通过同一 Driver 内的 addInput 直接相连,而是经 MergeSource 队列交接。

多份扫描 Driver 各有 TableScan、CallbackSink 和对应的本地队列;单个 LocalMerge Driver 合并这些生产流。
图 2:多份扫描 Driver 各有 TableScan、CallbackSink 和对应的本地队列;单个 LocalMerge Driver 合并这些生产流。 打开 SVG 原图

这里要区分 pipeline 定义与 Driver 实例:同一种 TableScan → CallbackSink 算子链可以实例化多份 Driver;计划中若有多个独立 source 子树,也会产生不同的 producer pipeline / DriverFactory。LocalMerge 所在的消费 pipeline 用一份 Driver 推进归并,后面还可以接结果输出等算子。因此实际可以看到多条扫描 pipeline,但不能把“24 个文件”硬编码成“必须有 24 个不同的 pipeline 定义”。

图中每条流对应一个文件只是便于理解。一个 source 也可以读取一个已经保证跨文件整体有序的序列;反之,任意几个各自有序的文件不能随意串起来就当作一条有序 source。真正的连接单位是一份 producer Driver 提供的一条完整有序流,不是文件名本身。

模块主要类或函数职责
计划与执行连接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 指针。
图 3:归并侧的主要类关系。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 构造检查中体现。

源码:LocalMergeNode,makeMergeSinkSupplier,Task::addLocalMergeSource,Merge 与 SourceMerger 声明。

2.2 沿一条输入流看启动、背压与结束的往返

在前面的 N=24、F=8 例子中,先追踪第一组的 source 0。规划阶段已经安装 TableScan → CallbackSink 与 LocalMergeSource 的连接,但未被 start 的 source 通过 startCB 让 producer Driver 等待。LocalMerge 选中第一组并 start 相应 sources,promise 完成后 Driver 重新入队,扫描出 batch 后由 consumerCB 入队。

若队列满,consumerCB 给出的 future 让 producer Driver off-thread;SourceMerger 取走批次、释放槽位后兑现 promise,producer 再次调度。当前组所有 source 结束后,归并结果形成一个有序文件组,才推进下一组。三个组都完成后,SpillMerger 为每组建立 BatchStream 与新队列,worker 读回文件,归并 Driver 从这些队列取得最终有序结果。

同一条链中的等待有两个调度主体:第一阶段是 Driver 的 BlockingState::setResume;最终阶段的 producer 是 SpillMerger 的 executor continuation。两者都通过队列 future 表达背压,但不能把文件 worker 当成 Task Driver,也不能把 BatchStream::nextBatch 当成返回 future 的异步 API。读回错误先记录再发送 EOF,消费端到达结束检查后重抛。

代码入口:回调安装、LocalMergeSource 队列、异步读回。

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

3.1 startCb、consumerCb 与 MergeSource 在什么时候安装

CallbackSink 是可接入不同消费者的通用 sink,它不硬编码 LocalMerge 指针。这个连接在LocalPlanner 为 producer pipeline 生成 sink supplier、Task 根据 DriverFactory 创建具体 Driver 和算子时建立。makeMergeSinkSupplier 本身先返回一个工厂闭包;执行这份工厂、真正构造某份 CallbackSink 时,才创建并注册属于它的 MergeSource。不是每次 addInput 都临时创建队列,也不是开始 spill 后再补装连接。

名称这条 LocalMerge 路径安装的函数什么时候调用
startCb(成员名 startedCb_)捕获本路 mergeSource,调用 mergeSource->started(future)CallbackSink 第一次 isBlocked;检查下游是否已允许开始读取
consumerCb(成员名 consumeCb_)捕获同一路 mergeSource,调用 enqueue(input, future)每次 addInput 发送批次;close 时传 nullptr 通知 EOF
mergeSourceTask::addLocalMergeSource 创建的共享对象在 CallbackSink 工厂执行时创建并登记,供两边运行期间共同使用

startCb 的名字容易让人误解:它不是让 producer 自己打开开关,而是询问 started 状态并取得等待 future。真正的启动动作是消费端 LocalMerge 在选中本组后调用 source->start()。consumerCb 才负责数据传递及生产端背压:普通批次已入队后如果变满,返回 kWaitForConsumer,让 CallbackSink 记录 future,稍后由 Driver 接管等待。

源码节选:makeMergeSinkSupplier:工厂、共享队列和两个闭包

OperatorSupplier makeMergeSinkSupplier(
    const core::PlanNodePtr& planNode,
    bool supportsDrain) {
  auto planNodeId = planNode->id();
  auto outputType = planNode->outputType();
  return [planNodeId = std::move(planNodeId),
          outputType = std::move(outputType),
          supportsDrain](int32_t operatorId, DriverCtx* ctx) {
    auto mergeSource = ctx->task->addLocalMergeSource(
        ctx->splitGroupId,
        planNodeId,
        outputType,
        ctx->queryConfig().localMergeSourceQueueSize());

    Consumer consumerCb;
    if (supportsDrain) {
      consumerCb =
          [mergeSource](
              RowVectorPtr input, bool drained, ContinueFuture* future) {
            return mergeSource->enqueue(std::move(input), future, drained);
          };
    } else {
      consumerCb =
          [mergeSource](
              RowVectorPtr input, bool drained, ContinueFuture* future) {
            VELOX_CHECK(!drained);
            return mergeSource->enqueue(std::move(input), future);
          };
    }

    auto startCb = [mergeSource](ContinueFuture* future) {
      return mergeSource->started(future);
    };

    return std::make_unique<CallbackSink>(
        operatorId,
        ctx,
        std::move(consumerCb),
        std::move(startCb),
        planNodeId,
        core::PlanNode::Boundary::kInput);
  };
}

这段工厂同时服务 LocalMerge 与 MixedUnion;当前 LocalMerge 选择不支持 drain 的分支,检查 drained 为 false 后调用普通 enqueue。构造函数把两个闭包分别 move 到 startedCb_ 和 consumeCb_,因此闭包捕获的 shared_ptr 也随 CallbackSink 保留。前面讲的首次 isBlocked 会消耗启动回调;消费回调则持续服务后续批次,直到关闭。

从 sink supplier 创建和注册队列,到 LocalMerge 找到同一组 source 并启动、消费它们。
图 4:从 sink supplier 创建和注册队列,到 LocalMerge 找到同一组 source 并启动、消费它们。 打开 SVG 原图

Task 的注册表按 splitGroupId + LocalMerge planNodeId 分组,每个 CallbackSink 往对应 vector 追加自己的 source。LocalMerge 首次通过 addMergeSources 取数时,用相同的两个标识调用 getLocalMergeSources,并把 shared_ptr 集合保存到 sources_。这就是上下游“接线”的位置:两边拿到的不是两个内容相似的队列,而是同一批共享对象。

源码节选:Task::addLocalMergeSource:创建并登记

std::shared_ptr<MergeSource> Task::addLocalMergeSource(
    uint32_t splitGroupId,
    const core::PlanNodeId& planNodeId,
    const RowTypePtr& rowType,
    int queueSize) {
  auto source = MergeSource::createLocalMergeSource(queueSize);
  splitGroupStates_[splitGroupId].localMergeSources[planNodeId].push_back(
      source);
  return source;
}

源码节选:LocalMerge::addMergeSources:取得消费端连接

BlockingReason LocalMerge::addMergeSources(ContinueFuture* /* future */) {
  if (sources_.empty()) {
    sources_ = operatorCtx_->task()->getLocalMergeSources(
        operatorCtx_->driverCtx()->splitGroupId, planNodeId());
  }
  return BlockingReason::kNotBlocked;
}

队列本身不拥有 TableScan 或 LocalMerge 的业务逻辑。生产端只需知道 started / enqueue,消费端只需知道 start / next;Task 负责使它们在同一个计划节点下相遇。这样文件扫描、归并轮次和等待通知可以独立实现,第二阶段的 spill reader 也能复用同一种本地队列接口。

源码:CallbackSink 的两个回调成员、Task::getLocalMergeSources。

3.2 为什么启动门放在 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 取数。时序竖线仅表示参与者,不代表专属线程。
图 5:下游 LocalMerge 打开启动门,上游 Driver 才继续向 TableScan 取数。时序竖线仅表示参与者,不代表专属线程。 打开 SVG 原图

3.3 下一组 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::runInternal,TableScan::initialize,TableScan::close,组启动控制。

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

4.1 本地 MergeSource 为什么这样设计

准确地说,MergeSource 是接口,当前 LocalMergeSource 才是本文使用的线程安全有界队列实现;远程 MergeExchange 有不同的实现。通常是一份 producer Driver 向一份本地 source 入队,一份 LocalMerge Driver 从它取数。因为两边可能由不同 worker 并发运行,队列中的批次、started / atEnd 状态以及等待 promises 必须一起受同步保护。

folly::Synchronized<LocalMergeSourceQueue> 把状态判断和登记等待者放在同一个队列锁下。例如 consumer 观察到队列空时,在释放锁前就登记 consumer promise;producer 后续入队时,在同一把锁下更新数据并取走应通知的 promises。这样不会在“检查为空”与“开始等待”之间漏掉一次入队通知。promise 移出临界区后才兑现,以免持锁触发 continuation、增加重入和锁顺序问题。

将启动、数据可用和背压收进同一个对象,也让三个边界一致:未 start 时 producer 等消费端许可;队列空时 consumer 等生产端进度;队列满时 producer 等消费端释放槽位。EOF 用独立 atEnd 标志表示,不能挤占本来就满的队列,也不能提前覆盖已有批次。默认容量是两个批次,它限制 producer 的领先量,却不是整个 reader 或归并工作集的字节上限。

isBlocked 之后怎样让出 CPU? CallbackSink 或 LocalMerge 把 reason 与 future 返回给 Driver;Driver 创建 BlockingState、退出 runInternal,并通过 Task::leave 变为 off-thread。等待期间,这份 Driver 不占 worker 去调用 future.wait()。source.start()、enqueue() 或 next() 改变状态后兑现相应 promise,future 的成功 continuation 才会按 Task 状态重新入队;executor 运行工作项并通过 Task::enter 后,Driver 才重新 on-thread。若算子在 getOutput 中才发现等待,框架会在空返回后再次检查 isBlocked 来接走这份 future。

这里的“让出”是协作调度,不是 isBlocked 函数内部让 OS 线程睡眠;future 满足也不意味着恢复原来的 C++ 调用栈。完整的退出、排队与执行资格交接见 Task & Driver:从 off-thread 到 on-thread。下面按本地队列的实际字段和每个操作逐项展开。

4.2 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 不消耗容量;消费者先排空已有批次,然后才观察结束。
图 6:queueSize = 2 的往返。EOF 不消耗容量;消费者先排空已有批次,然后才观察结束。 打开 SVG 原图

4.3 为什么 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. 归并内核:比较、游标、批次复制

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

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

每次 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,不能把“参考输入编码建输出”写成整个归并是零拷贝。

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

源码节选:每个 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 行上限。

源码:TreeOfLosers,SourceStream 比较,copyToOutput,createOutputVector,outputBatchRows。

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

6.1 字段组合表达执行进度

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

6.2 为什么最后一轮也落盘

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

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

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

6.3 阻塞 future 怎样回到 Driver

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

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

源码:getOutputFromSource,finishMergeSourceGroup,Merge::isBlocked,MergeSpiller。

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

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

文件流模块的类关系及数据方向。BatchStream 的同步接口与 SpillMerger 的异步调度分处两层。
图 9:文件流模块的类关系及数据方向。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、反序列化和输出向量成本。

7.1 一条 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> 混为一谈,这条路径没有通过那个类实现异步化。

源码:setupSpillMerger,ConcatFilesSpillBatchStream,SpillReadFile 构造,SerializedPageFileReader,FileInputStream。

SpillMergeStream::current() 当前返回 const RowVectorPtr&,调用方通过 .get() 借用 RowVector 指针,或者复制 shared_ptr 保活批次;它不再返回 const RowVector&。无论是哪种接口,gather 前仍要遵守 currentIndex 的批次结束边界:如果还借用了当前批次的数据,必须先复制或取得有效所有权,再让 pop 推进到可能替换 rowVector_ 的下一批。见 current / currentIndex、gatherMerge。

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

8.1 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。
图 10:异步调度解耦了多路读取和单线程归并;每路仍然串行调用自己的 BatchStream。 打开 SVG 原图

8.2 快路径直接读,满队列后挂 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 游标。两个层次各管一段状态,职责比较清楚。

8.3 有界队列约束的是什么

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

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

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

源码:初始调度与后续 continuation。

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

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

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;
}

9.2 同步读取失败与 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 的原始动态异常信息,但那属于后续代码改进。

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

错误先记录到共享槽,再通过每个 source 的 EOF 收敛;归并端到达结束时统一重抛。箭头中的 EOF 实际经各自队列传递。
图 11:错误先记录到共享槽,再通过每个 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 不提供回滚,宿主必须把整个查询的最终失败与结果提交语义衔接起来。

9.4 为什么在 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. 对象所有权、回调寿命与提前取消

10.1 谁负责让对象继续存在

关系持有方式它解决的事
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 保存,避免把一个裸池指针交给异步读回后完全失去对象寿命支撑。

10.2 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 取消和未执行回调,需要一起审视。这里指出的是当前实现的证明边界,不是宣称已经复现了一个取消崩溃。

10.3 析构和 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. 启用条件、配置与运行观测

11.1 配置 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 数验证行为,不能只看配置已经写入。

11.2 怎样读两个阶段的统计

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

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

源码:LocalMerge 配置,makeSpillConfig,LocalMerge 构造,阶段统计。

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 销毁、取消交错都已经得到形式化保证。读实现时需要把测试覆盖到的具体时序与未覆盖的边界分开。

测试源码:MergeTest,MergerTest。

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

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

把单个原始 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 能力测量。

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

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

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

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

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

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

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

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

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

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

把一条流的启动、满队列等待和 EOF 串起来看,lazy start 与 bounded queue 限制的是不同阶段的资源:前者限制重 reader 何时开始工作,后者限制生产者相对消费者可以领先多少批。两者叠加仍不等于精确内存上限,因为当前批次、解码缓冲、预取和输出都在队列之外。调参必须把这些对象的实际生命周期放进估算,而不能只乘 queueSize。

14. 参考文章与源码入口

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