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,归并出最终结果。某一组完全没有行就不会生成文件组。
这是多个局部归并轮次,加一次最终归并。当前实现不会反复递归,把第二阶段也一直限制成 F 路。N ≤ F 时走普通流式 LocalMerge,无需为了“多轮模式”强制落盘。N > F 时所有原始输入组的结果都要 spill,最后一组同样如此。
2. 算子、队列、游标与文件流各做什么
规划层负责建立跨 pipeline 的连接,Merge 负责阶段与轮次,SourceMerger 负责有序归并,MergeSource 负责同步。SpillMerger 把文件读回转成 MergeSource 所提供的相同接口,因此最终归并可以继续使用 SourceMerger。
| 模块 | 主要类或函数 | 职责 |
|---|---|---|
| 计划与执行连接 | LocalMergeNode / LocalPlanner / Task | 定义排序规则,建立 producer pipeline,登记并共享 MergeSource |
| 生产入口 | CallbackSink | 启动门检查、批次入队、生产结束通知 |
| 跨线程同步 | MergeSource / LocalMergeSource | started、队列、EOF、producer / consumer promises |
| 算子控制 | LocalMerge → Merge → SourceOperator | 选择 source 组、调度 spill、处理 Driver 阻塞和结束 |
| 归并内核 | SourceMerger / SourceStream / TreeOfLosers | 比较排序键、推进游标、组装输出向量 |
| 中间写出 | MergeSpiller | 接收已经有序的 RowVector,写出同一 run 的文件 |
| 最终读回 | SpillMerger / BatchStream / SpillReadFile | 并行读取不同文件组,经队列交给单线程归并 |
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 声明。
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 让它再次变成可运行状态。
下一组 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. 一份队列协议,连接启动、数据和结束
LocalMergeSource 保存哪些状态
LocalMergeSource 的内部队列由 folly::Synchronized 保护,批次存于 boost::circular_buffer<RowVectorPtr>。started_ 表示生产是否获准启动,atEnd_ 表示生产端声明不会再提供数据;producerPromises_ 和 consumerPromises_ 分别承接两个方向的等待。容量以批次数计量,默认是 2,而不是按字节计量。
| 操作与条件 | 状态变化 | 调用者怎样继续 |
|---|---|---|
| started:尚未启动 | 保存 producer promise | 返回 kWaitForConsumer,等待 start |
| start | started_ = 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。两条路径使用同一套背压规则。
为什么 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 自动通知全部等待者”的路径。
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 字段加入排序键。
不是每次选中一行就立即复制整行
每次 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 行上限。
源码:TreeOfLosers,SourceStream 比较,copyToOutput,createOutputVector,outputBatchRows。
6. 多轮控制:从当前 source 组切换到最终归并
字段组合表达执行进度
| 字段 | 含义 | 关键变化 |
|---|---|---|
| sources_ / numStartedSources_ | 全部原始输入 / 已启动前缀的长度 | sources_ 保留原始连接;计数按组推进 |
| sourceMerger_ | 当前原始输入组的归并状态 | 一组一个实例,组结束后清空 |
| mergeOutputSpiller_ | 当前组的输出 writer | 首次出现输出才创建,finishSpill 后清空 |
| spillFileGroups_ | 已经完成的非空有序 runs | 每结束一个非空组追加一项 |
| spillMerger_ | 第二阶段的整体读回管理器 | 所有原始组完成后创建一次 |
| sourceBlockingFutures_ / finished_ | 待处理的阻塞通知 / 算子完成标记 | 等待与结束分别表达 |
流程伪代码:展示阶段切换;省略统计、检查和移动语义
// 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 尚未把读回批次放入队列。
源码:getOutputFromSource,finishMergeSourceGroup,Merge::isBlocked,MergeSpiller。
7. 文件组读回:有序拼接与真正的资源边界
一轮的归并输出本身全局有序。写满一个 spill 文件再写下一个,只是把同一条有序 run 分段,因此文件组内部可以直接按原次序拼接;不同组的 key 范围可能重叠,才需要在组之间做有序归并。
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> 混为一谈,这条路径没有通过那个类实现异步化。
源码:setupSpillMerger,ConcatFilesSpillBatchStream,SpillReadFile 构造,SerializedPageFileReader,FileInputStream。
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 执行。
快路径直接读,满队列后挂 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 保证。
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 的原始动态异常信息,但那属于后续代码改进。
一条流出错以后,其余流如何停止
失败的回调先 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 → SpillMerger | shared_ptr | 让最终归并管理器可以在执行回调期间保活 |
| SpillMerger → MemoryPool / SpillStats | shared_ptr | 异步工作仍可依赖池对象和统计对象的生命周期 |
| SpillMerger → BatchStream | vector<unique_ptr<…>> | 集中管理每个文件组的读取状态 |
| SpillMerger → MergeSource | vector<shared_ptr<…>> | 持有各路队列,让 consumer 与 producer 共享连接 |
| SpillMerger → SourceMerger | unique_ptr | 唯一拥有最终归并状态;其 cursor 借用 sources |
| SpillMerger → executor | folly::Executor* | 非拥有引用,executor 必须由外部正确管理寿命 |
| 读回回调 → SpillMerger | weak 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_enabled | false | 全局允许 spill;还需要宿主提供 spill 目录 |
| local_merge_spill_enabled | false | 允许有排序键的 LocalMerge 创建 spill 配置 |
| local_merge_max_num_merge_sources | uint32_t 最大值 | 第一阶段每组的 F;要求大于 0,1 是合法值 |
| local_merge_source_queue_size | 2 | 每路可排队的 RowVector 批次数,两个阶段都会使用 |
| preferred_output_batch_rows | 1024 | 当前 Merge 构造传入的输出行数上限 R |
| preferred_output_batch_bytes | 10 MiB | 结合首批估计行宽缩小输出行数,不是硬内存限额 |
| max_output_batch_rows | 10000 | 其他 outputBatchRows 分支使用;当前 Merge 无参数调用不读取此上限 |
| max_spill_bytes | 100 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 配置,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. 设计取舍:用阶段边界管理工作集
内存模型需要同时计算两阶段
把单个原始 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. 参考文章与源码入口
- Multi-Round Lazy Start Merge,官方背景与最初设计说明。
- Merge.cpp / Merge.h:控制流、归并内核、SpillMerger。
- MergeSource.cpp / CallbackSink.cpp:启动、队列、EOF 与通知。
- Spill.cpp / SpillFile.cpp:文件组拼接与 reader。
- PR #13139:merge source 启动机制;PR #13634:异步 SpillMerger。当前行为以本文固定快照为准。
源码节选来自 Apache-2.0 许可的 Velox 项目,原始版权和许可见各链接文件头及 项目 LICENSE。本文的流程图、类图和时序图按当前实现重新绘制。