Velox Query In the Presto Worker
在 Presto 原生 Worker 中,一次查询并不是先“进入 Velox Query 对象”再统一执行。Coordinator 持续发送 task 更新,Worker 转换计划、复用 query 资源、创建 Velox Task,并继续投递 split;数据则沿本地算子和远程 exchange 两条路径流动。
执行流程与相关实现核对于 2026-09-29,Velox 源码版本为 48883e8521b2。下文源码节选、历史资料和宿主集成引用各自注明版本;教学输入用于解释状态变化,未作为性能基准运行。
1. 分布式 Query 与 Velox 执行片段的边界
这篇文章讨论 Presto 原生 Worker 如何调用 Velox 执行查询片段。分布式 Query 和 Stage 的计划、调度及跨 worker 协议由宿主组织;Velox Task 执行一个 PlanFragment,并在内部管理 pipeline、Driver、Operator 和资源。理解这条边界,才能判断某个配置、等待或状态变化属于哪一侧。
| 概念 | 描述的对象 | 对应的工作 |
|---|---|---|
| Query / Stage | 分布式查询及其 fragment 依赖 | 宿主决定任务放置、输入分发和跨任务连接 |
| PlanFragment | 一个片段内怎样计算 | 转换为 Task 内的 pipeline 和 Operator 链 |
| Task | 一个片段的执行实例 | 接收更新和 split,管理执行与结束流程 |
| Split | 一份可被 source 消费的输入工作描述 | 扫描 split 标识待读取工作,远程 split 可描述上游数据来源 |
| Driver | 一份推进算子链的执行实例 | 消费可用输入,遇到 I/O、反压或 barrier 时等待 |
| Exchange | 数据跨执行边界的传输 | 本地交换批次;远程交换还涉及序列化、传输、缓冲与确认 |
计划描述“怎样处理”,split 描述“还要处理哪些输入”,数据页才携带实际传输的行。它们到达的时间和结束信号并不相同。没有更多 split、上游没有更多数据以及某个 Task 结束,也分别属于不同的协议状态。
下面先用一条 Join 后聚合的查询串起计划更新、split、pipeline 和远程输出,再按请求进入 Worker、输入分发、执行与交换的顺序展开实现。并行度和内存限制放回各自作用的执行边界来解释。
本文于 2026-09-19 重写。Velox 基线为 1d1b76567870;Presto 宿主代码单独核对 d621a736e77c。 两者是独立源码快照,不宣称它们已组成一次经过构建验证的发行版本。正文明确标注接口归属,重点解释常规 interactive、流式 HTTP exchange 路径。
代码片段分别标明源码节选或流程示意;流程示意省略统计、异常包装与无关分支,不是可独立编译的程序。历史资料与宿主集成保留各自版本,不能据此推断它们组成了经过构建验证的发行版本。
2. 层级结构总览
Query
└── Stage 0 (plan fragment 0: TableScan + partial agg)
│ ├── Worker A → Task-0-A
│ │ ├── Pipeline 0: Driver0, Driver1, ... DriverN
│ │ └── Pipeline 1: Driver0 (HashBuild side)
│ ├── Worker B → Task-0-B
│ └── Worker C → Task-0-C
│
└── Stage 1 (plan fragment 1: Exchange + HashJoin + PartitionedOutput)
├── Worker A → Task-1-A
├── Worker B → Task-1-B
└── Worker C → Task-1-C
关键约束:
- worker/stage 到 task 的关系由宿主决定;本文示例为每个参与 worker 放置一个 task,实际应结合 taskId、attempt 和调度模式识别。
- 在本文常规放置示例中,不同 stage 可在同一 worker 各放一个 Task;实际数量和重试 attempt 由宿主调度决定。
- Task 内的并行度靠 Driver 数量,不靠 Task 数量
- Coordinator 按目标 Task 与 source PlanNode 增量发送 splits;对应 Task 内的 Driver 消费该来源的 split store。一个 worker 上可有不同 stage / attempt 的多个 Task,不能把整个 worker 的 splits 混进同一个 Task。
3. 完整端到端数据流
考虑 SELECT f.k, sum(f.v) FROM fact f JOIN dim d ON f.k = d.k GROUP BY f.k。下面只展示一种分阶段方式中的一个分区:远程 build 输入来自 A,B 本地扫描 probe 输入并做部分聚合,C 合并聚合中间态。多个 B Task 时,两侧必须按同一 Join key 分区或采用有效的广播安排;这不是声称优化器总会生成这个计划。
| 阶段 | 具体输入与输出 | 控制信息怎样到达 |
|---|---|---|
| A 读 build 输入 | dim.k = [1, 2] → 序列化数据页 | Coordinator 把 scan splits 发给 A;A 的扫描 pipeline 以 PartitionedOutput 结束。 |
| B 建立 Join 状态 | Exchange 反序列化 A 的页 → HashBuild | 远程 split 告诉 B 从哪个 Task / destination 拉数据;HashBuild 汇合后经 HashJoinBridge 发布共享表。 |
| B 探测并聚合 | fact = [(1,4), (1,5), (3,7)] → 匹配前两行 → (1,9) | B 的另一条 pipeline 为 TableScan → HashProbe → partial aggregation → PartitionedOutput;probe 等 build 时让出 Driver。 |
| C 汇总并输出 | Exchange → final aggregation → 结果 (1,9) | 对聚合中间态进行 final 合并,再通过输出协议交给消费者。 |
| 结束与背压 | 没有更多输入、缓冲排空、引用释放是不同条件 | split 的 noMore 与数据页的 complete 分开;消费确认解除上游容量等待,Task 结束还要满足自身完成条件。 |
这个例子中的 HashBuild 是一条 pipeline 的 sink,不能把它与 PartitionedOutput 串在同一条数据链上。build 输出是通过 bridge 共享的表状态,真正的 Join 行由 HashProbe 产生。见 pipeline 切分、建表发布、探测输出。
3.1 背压、结束和失败是同一条生命周期链
慢 consumer 可以让下游交换队列积压、减少网络请求,上游未被确认的输出增加,最终使 PartitionedOutput 等待。Task 内 local exchange 还有自己的容量和 future。排查“CPU 不忙但查询卡住”时,应沿等待链寻找谁负责唤醒,而不是只增加 Driver 数。
| 等待位置 | 唤醒/解除条件 |
|---|---|
| split store | 新 split 或 noMore |
| connector | IO future 完成或异常 |
| HashProbe | bridge 发布 build result 或取消 |
| local producer | consumer 释放队列占用 |
| remote producer | 输出被确认、删除或 task 终止 |
取消还要传播到这些等待点,释放 source、buffer 和 Operator 状态;一个 Task 已不再执行,不意味着远端不再可能发来迟到请求。清理逻辑必须考虑异步回调仍持有引用。
Coordinator
POST /v1/task/A0 + dim file splits + noMoreSplits
POST /v1/task/B1 + remote split(A0) + fact file splits
POST /v1/task/C2 + remote split(B1)
Worker A, Stage 0:
TableScan(dim splits) → PartitionedOutput → OutputBuffer[destination]
↑ HTTP GET / acknowledge
Worker B, Stage 1: PrestoExchangeSource
build pipeline: Exchange(A0) → HashBuild ──publish──→ HashJoinBridge
│ table / future
probe pipeline: TableScan(fact splits) → HashProbe ←────────┘
→ PartialAgg → PartitionedOutput
→ OutputBuffer
↑ HTTP GET
Worker C, Stage 2: PrestoExchangeSource
Exchange(B1) → FinalAgg → PartitionedOutput → OutputBuffer
↑ GET results / acknowledge
Coordinator / consumer
NoMoreSplits ends input-source discovery; page completion ends the data stream.
HashJoinBridge is shared inside B's Task; it is not the HTTP exchange transport.
4. Worker 接收查询:HTTP → Velox Task
4.1 Worker 入口:HTTP 更新如何变成 Velox 计划
Presto 的 TaskResource::registerUris 注册创建/更新 task、状态查询、结果读取、acknowledge 和删除等端点。其中 POST /v1/task/<taskId> 进入 createOrUpdateTask;请求可携带计划、source splits 和输出缓冲更新。
createOrUpdateTask 的当前顺序是:解析请求;若携带 fragment,则获取 QueryCtx、构造 VeloxInteractiveQueryPlanConverter、转换并验证计划;然后交给 TaskManager。后续只投递 splits 的请求不需要重新转换一遍完整计划。
PrestoToVeloxQueryPlan.cpp 是文件名;当前类名不是旧文中笼统使用的 “PrestoToVeloxQueryPlan 对象”。它的 toVeloxQueryPlan 转换表达式、schema、join/aggregation 等节点以及 fragment 输出分区;RemoteSource 转成 Exchange 或 MergeExchange 等输入形式。
这个转换保留宿主语义和分布式边界。Velox LocalPlanner 在后面处理 Task 内 pipeline,并不重新执行 Coordinator 的成本优化。
4.2 HTTP 端点(TaskResource)
Coordinator 通过 REST API 驱动每个 Worker:
POST /v1/task/{taskId} 创建/更新 task(含 plan fragment + splits)
GET /v1/task/{taskId}/results/{bufferId}/{token} 下游拉取输出数据
GET /v1/task/{taskId}/status 轮询 task 状态(支持 long-polling)
DELETE /v1/task/{taskId} 终止 task
GET /v1/task/{taskId}/results/.../acknowledge ACK 已收数据
DELETE /v1/task/{taskId}/results/{bufferId} 释放 output buffer
4.3 Plan Fragment 转换(PrestoToVeloxQueryPlan)
TaskUpdateRequest 中的 fragment 字段是 Base64 编码的 JSON Plan:
Presto JSON PlanFragment
↓ VeloxQueryPlanConverter(历史示例旧名:VeloxInteractiveQueryPlanConverter)
Velox PlanFragment(PlanNode 树)
关键 PlanNode 映射:
| Presto PlanNode | Velox PlanNode | 说明 |
|---|---|---|
TableScanNode |
TableScanNode |
读存储 |
RemoteSourceNode |
ExchangeNode |
跨 worker 拉数据 |
PartitionedOutputNode |
PartitionedOutputNode |
推数据给下游 |
AggregationNode |
AggregationNode |
聚合 |
ExchangeNode |
LocalExchangeNode |
task 内重分区 |
4.4 Task 创建流程(TaskManager)
POST /v1/task/{taskId}
↓
findOrCreateTask() → 创建 PrestoTask(状态 kPlanned,暂无 Velox Task)
↓
createOrUpdateTaskImpl()
├── exec::Task::create(planFragment, queryCtx, ExecutionMode::kParallel)
├── execTask->addSplitWithSequence() ← 注册 splits
└── startTaskLocked()
└── execTask->start(maxDrivers, concurrentLifespans)
PrestoTask 是 Velox exec::Task 的包装,额外持有:
- 待处理的
ResultRequest(下游长轮询的 HTTP 请求) - 状态变更 Promise(用于 long-polling)
5. Split 分发与消费机制
5.1 Split 分发:批次、watermark 与结束信号
当前 Presto TaskSource 合并与投递 先按 planNodeId 合并同一请求的 sources,再逐个转换 protocol split:
对当前 source 的所有 split:
addSplitWithSequence(planNodeId, split, sequenceId)
记录本批最大 sequenceId
本批添加完成后:
setMaxSplitSequenceId(planNodeId, batchMax)
再处理 group 结束与全局 noMoreSplits
Velox addSplitWithSequence 只接受大于已经记录的 maxSequenceId 的 split;setMaxSplitSequenceId 单调推进 watermark。关键是批次结束后才更新 watermark,避免一个批次内乱序编号的较小 split 被提前排除。
这不是一个无限制的已见 ID 集合。仅调用 addSplitWithSequence、却不遵守宿主的投递顺序和 watermark 协议,不能获得任意乱序、任意重复请求下的 exactly-once 保证。
Velox source 通过 getSplitOrFuture 获取 split;暂时为空时等待 future,收到 noMore 后才知道不会再有工作。
noMoreSplitsForGroup 关闭某个 lifespan/group 的输入;noMoreSplits 关闭整个 source 的后续工作/后续 group。Presto 在 task 尚未启动时会暂存全局 noMore 信号,启动后补交。不能把 group 完成、source 完成、Driver 完成和结果 buffer 消费完成合并成一个事件。
上面讲了 split 和 plan fragment 分开传输,这里展开完整链路:从 coordinator 下发,到 worker 入队,到 source operator 领取消费。
5.2 全链路概览
边界说明:哪个 split 分给哪个 worker,是 coordinator(Java) 的调度决策,不在 C++ worker 代码里。Worker 侧只负责接收分配给自己的 split 并消费。
5.3 第 1 步:Worker 接收与转换
TaskResource::createOrUpdateTask() 从 TaskUpdateRequest 解析出 sources[],每个元素(protocol::TaskSource)含:
源码核对:宿主协议中的 fragment 可随首次创建或后续更新携带;是否省略应按该版本 coordinator/worker 协议判断。Velox Task 本身不规定 HTTP 请求的重传和字段省略策略。
流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。
struct TaskSource {
PlanNodeId planNodeId; // 关联到 fragment 中的 source 节点
vector<ScheduledSplit> splits; // 本批 split
bool noMoreSplits; // 该节点是否不再有 split
vector<Lifespan> noMoreSplitsForLifespan; // grouped 执行的 per-group 结束标记
};
TaskManager::createOrUpdateTaskImpl() 处理两阶段:
- 合并:同一请求里可能有多个相同
planNodeId的 source,先按 planNodeId 合并splits、noMoreSplitsForLifespan,noMoreSplits做 OR - 转换:
toVeloxSplit(ScheduledSplit)把 Presto 协议 split 转成velox::exec::SplitRemoteSplit→RemoteConnectorSplit(下游拉上游用)EmptySplit→ 空 split(占位,什么都不读)- 其他 → 委托给 connector 特定的
toVeloxSplit()(如 Hive split) lifespan.groupid→Split.groupId(-1 表示 ungrouped)
5.4 第 2 步:addSplitWithSequence —— 为什么需要 sequence
Coordinator 可能重发同一批 split(网络重试、task update 重叠),所以用 sequence id 去重:
当前源码摘录(1d1b76567870):velox/exec/Task.cpp:1708。省略外围声明;此片段未作为独立程序编译。
bool Task::addSplitWithSequence(
const core::PlanNodeId& planNodeId,
exec::Split&& split,
long sequenceId) {
RECORD_METRIC_VALUE(kMetricTaskSplitsCount, 1);
std::vector<ContinuePromise> promises;
bool added = false;
bool isTaskRunning;
bool shouldLogSplit = false;
{
std::lock_guard<std::timed_mutex> l(mutex_);
isTaskRunning = isRunningLocked();
if (isTaskRunning) {
// The same split can be added again in some systems. The systems that
// want 'one split processed once only' would use this method and
// duplicate splits would be ignored.
auto& splitsState = getPlanNodeSplitsStateLocked(planNodeId);
if (sequenceId > splitsState.maxSequenceId) {
shouldLogSplit = true;
addSplitLocked(splitsState, split, promises);
added = true;
}
}
}
for (auto& promise : promises) {
promise.setValue();
}
if (!isTaskRunning) {
// Safe because 'split' is moved away above only if 'isTaskRunning'.
// @lint-ignore CLANGTIDY bugprone-use-after-move
addRemoteSplit(planNodeId, split);
}
if (shouldLogSplit) {
onAddSplit(planNodeId, split);
}
return added;
}
maxSequenceId是水位线:seqId ≤ 水位的 split 视为重复直接丢弃- 注意:
addSplitWithSequence本身不更新水位,而是由setMaxSplitSequenceId(nodeId, maxSeqId)在一批 split 加完后统一设置水位 - 对比
addSplit():无 sequence 检查,无条件入队(用于不需要幂等投递的场景)
5.5 第 3 步:Velox 侧 split 队列结构
流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。
// 每个 plan node 一份状态
struct SplitsState {
bool sourceIsTableScan{false};
bool noMoreSplits{false};
long maxSequenceId{LONG_MIN}; // 去重水位
// splitGroupId → 该 group 的 split 队列
unordered_map<uint32_t, unique_ptr<SplitsStore>> groupSplitsStores;
};
// 单个 (planNodeId, splitGroupId) 的 split 队列
class SplitsStore { // 实现 QueueSplitsStore
deque<Split> splits_; // 待消费 split
bool noMoreSplits_{false};
vector<ContinuePromise> promises_; // 等待 split 的 driver
};
splitsStates_:map<PlanNodeId, SplitsState>,每个 source 节点一份- ungrouped 执行:每个节点只有一个 store(key =
kUngroupedGroupId) - grouped 执行:每个 split group 一个 store
5.5.1 为什么拆成 SplitsState / SplitsStore 两层
这个两层嵌套不是随意设计,而是两个正交维度被拆开:
- SplitsState = per plan node(每个 source 节点一份)
- SplitsStore = per **(plan node, split group)**(每个 split group 一个队列)
有四个独立的设计力推动这个拆分:
1. 去重水位是 node 级,消费队列是 group 级——作用域不同。
maxSequenceId 去重水位必须是 node 级:coordinator 给每个 plan node 的 split 分配连续 sequence id,与 group 无关。而"排队 + 阻塞 + 消费"是 group 级的。两者作用域天然不同。
2. Bucketed 执行要求"一个节点 N 个独立队列"(最硬的需求)。
同一个 TableScan 节点产出的 split 带不同 groupId,每个 group 的 driver 组只能拉自己 group 的 split,且要独立阻塞——group 5 队列空了阻塞,不能影响 group 1 继续跑。所以每个 group 必须有自己的队列和自己的阻塞 promises,强制了 map<groupId, SplitsStore>。
不可复制与其中独占状态及 promise 的所有权语义一致;不能仅从 deleted copy constructor 推断等待方保存了指向 promise 容器元素的地址。Future 的关联状态与容器对象地址不是一回事,应按实际传递的对象和引用检查生命周期。
4. SplitsStore 是抽象类,还抽象了"split 从哪来"。
store 带虚函数(nextSplit / requestBarrier),注释原文:either can accumulate splits through addSplit() or generating splits from its own source。即 split 来源可替换:QueueSplitsStore 被动接收 coordinator push 的 split,或自生成 source。用类继承而非结构体字段来解耦"消费接口"与"来源实现"。
| SplitsState | SplitsStore | |
|---|---|---|
| 作用域 | 每个 plan node | 每个 (node, split group) |
| 装什么 | 去重水位、node 级 noMoreSplits、是否 TableScan | split 队列、barrier split、阻塞 promises |
| 生命周期 | 整个 task | 随 split group 动态创建/清理 |
| 形态 | 不可拷贝 struct | 带虚函数的抽象类 |
一句话:SplitsState 管"节点整体元数据",SplitsStore 管"某 group 具体怎么排队和阻塞"。根本原因是 grouped 执行要求同一节点下多个 group 各自独立排队与阻塞。
5.5.2 非 grouped 执行(ungrouped)长什么样
绝大多数普通查询(非 bucketed)走的是 ungrouped 执行,此时这套结构退化成最简单的形态:
ungrouped 执行(所有 split.groupId == -1)
SplitsState (planNode "0")
├── maxSequenceId
├── noMoreSplits
└── groupSplitsStores
└── [kUngroupedGroupId] → SplitsStore ← 整个节点只有这一个 store
├── splits_ (一条队列)
└── promises_
所有 driver(Driver 0..N)都从这唯一一个 store 抢 split
| ungrouped 执行 | grouped 执行 | |
|---|---|---|
split.groupId |
-1 | bucket id(0..N-1) |
| 每节点 store 数 | 1(kUngroupedGroupId) |
每 bucket 一个 |
| driver 创建时机 | task start 时一次性建好 | 每个 group 到达时按需创建 |
| split 消费 | 所有 driver 抢同一队列,谁空闲谁拿 | 每 group 的 driver 只拿本 group split |
| 阻塞粒度 | 整个节点一个 promises 集合 | 每 group 独立阻塞 |
noMoreSplits |
通知唯一的 store | 通知该节点所有 group 的 store |
| 适用场景 | 普通表扫描、shuffle 后的 join | bucketed 表的 grouped/colocated join |
关键点:ungrouped 不是"另一套代码路径",而是 grouped 设计的特例——map<groupId, store> 里只有一个固定 key kUngroupedGroupId = UINT32_MAX 的 entry。这样一套数据结构同时覆盖两种执行模式,无需为常规查询单独写一条路径。多个 driver 并行消费同一队列,靠的就是 SplitsStore 内部对队列和 promises 的加锁访问(实际锁在 Task::mutex_)。
5.6 第 4 步:Source operator 领取 split
TableScan 在 getOutput() 中调用 Task::getSplitOrFuture() 领取 split:
流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。
BlockingReason Task::getSplitOrFuture(
driverId, splitGroupId, planNodeId, maxPreloadSplits, preload,
Split& split, ContinueFuture& future) {
// 取 (planNodeId, splitGroupId) 对应的 SplitsStore
// 调 store->nextSplit(...)
}
SplitsStore::nextSplit() 三种结果:
| 情况 | 返回 |
|---|---|
splits_ 非空 |
弹出 split,kNotBlocked |
队列空但已 noMoreSplits_ |
无 split,kNotBlocked(operator 据此 finish) |
| 队列空且未结束 | 创建 ContinuePromise 存入 promises_,返回 kWaitForSplit + future |
阻塞时 driver 退出线程池;当新 split 入队或 noMoreSplits 触发时,fulfill promises_ 唤醒 driver 重新调度。
5.7 第 5 步:noMoreSplits 的两个层级 + 延迟处理
两个层级:
流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。
// 单个 group 结束(grouped 执行)
void Task::noMoreSplitsForGroup(planNodeId, splitGroupId);
// 整个 plan node 结束
void Task::noMoreSplits(planNodeId) {
splitsState.noMoreSplits = true;
// ungrouped: 通知唯一的 store
// grouped: 通知该节点所有 group 的 store
// 若所有 split group 都处理完 → terminate(kFinished)
}
延迟处理(delayedNoMoreSplitsPlanNodes_):
noMoreSplits 信号可能在 task->start() 之前就到达。若此时直接调 Velox 的 noMoreSplits 会出问题(task 还没起来)。所以:
// TaskManager.cpp
if (prestoTask->taskStarted) {
execTask->noMoreSplits(planNodeId); // 已启动,直接调
} else {
prestoTask->delayedNoMoreSplitsPlanNodes_.emplace(planNodeId); // 暂存
}
// 之后 startTaskLocked() 中补调:
for (auto& nodeId : delayedNoMoreSplitsPlanNodes_) {
execTask->noMoreSplits(nodeId);
}
delayedNoMoreSplitsPlanNodes_.clear();
否则 driver 会永久阻塞在 kWaitForSplit,等待一个永远不会来的 split。
5.8 第 6 步:Split group(bucketed/grouped 执行)
针对 bucketed 表的分组执行,split 带 groupId:
流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。
// Task.cpp,split 带 groupId != -1 时
if (!seenSplitGroups_.contains(groupId)) {
seenSplitGroups_.insert(groupId);
queuedSplitGroups_.push(groupId);
ensureSplitGroupsAreBeingProcessedLocked(); // 为新 group 创建 driver
}
addSplitToStore(planNodeId, groupId, split);
- 每个 split group 独立创建一组 driver,
concurrentSplitGroups控制同时处理的 group 数 kUngroupedGroupId = UINT32_MAX表示非分组执行- 同一 split group 的 build/probe 与各 peer 通过该 group 的共享状态协调;不同 group 通常各有状态。Task 还存在用于执行协议的 barrier,不能将这些机制一概解释为所有 group 互相等待。
5.9 Split 分发关键代码位置
| 环节 | 文件 |
|---|---|
| 解析 sources[] | presto_cpp/main/TaskResource.cpp createOrUpdateTask() |
| 入队 + 去重 + noMoreSplits | presto_cpp/main/TaskManager.cpp createOrUpdateTaskImpl() |
| 协议 split 转换 | presto_cpp/main/PrestoToVeloxSplit.cpp toVeloxSplit() |
| 延迟标记 | presto_cpp/main/PrestoTask.h delayedNoMoreSplitsPlanNodes_ |
| addSplitWithSequence / 去重 | velox/exec/Task.cpp |
| SplitsState / SplitsStore | velox/exec/TaskStructs.h |
| getSplitOrFuture | velox/exec/Task.cpp |
| noMoreSplits / ForGroup | velox/exec/Task.cpp |
6. Query 如何并行执行
6.1 Query 的三个并行尺度
Coordinator 把查询切成多个 stage,每个 stage 执行一个 plan fragment;同一 stage 可以在不同 worker 创建多个 task。一个 Velox Task 内部再有多个 pipeline,每条 pipeline 创建若干 Driver。
| 尺度 | 谁决定 | 主要约束 |
|---|---|---|
| Stage 之间 | Presto 计划与调度 | fragment 依赖和 exchange |
| 同一 stage 的多个 task | Coordinator | worker、分区、资源策略 |
| Task 内 Driver | Velox 与宿主传入配置 | maxDrivers、plan node 约束、split group |
| 某个 Driver 的 Operator 链 | Velox Driver | 串行推进;遇到 future 暂停 |
split 数决定待处理工作量,Driver 数决定 Task 内可同时推进的执行实例数量,executor 决定可用线程资源。三者不能互换。
一个 query 的并行度由三个层级叠加而成,从粗到细:
Stage 并行(query 切成多个 plan fragment)
└── Task 并行(每个 stage 分发到多个 worker)
└── Driver 并行(每个 task 内多个 driver)
6.2 Stage 级:query 切成多个 plan fragment
Coordinator 把 query 的执行计划按 shuffle 边界(exchange)切成多个 plan fragment,每个对应一个 stage:
Query
├── Fragment 0 (Stage 0): TableScan → partial agg → PartitionedOutput
├── Fragment 1 (Stage 1): RemoteSource → HashJoin → PartitionedOutput
└── Fragment 2 (Stage 2): RemoteSource → final agg → output
- **Plan fragment 描述"怎么算"**:算子树结构,静态,task 创建时发一次
- 上下游 stage 通过
PartitionedOutput(上游 sink)↔RemoteSource/Exchange(下游 source)衔接
6.3 Task 级:每个 stage 分发到多个 worker
Coordinator 把每个 stage 的 plan fragment 分发到参与该 stage 的所有 worker,每个 worker 为该 fragment 创建一个 task:
Stage 0 (Fragment 0) Stage 1 (Fragment 1)
├── Worker A → Task-0-A ├── Worker A → Task-1-A
├── Worker B → Task-0-B ├── Worker B → Task-1-B
└── Worker C → Task-0-C └── Worker C → Task-1-C
这里用每个 stage 在一个 worker 上放一个 task 的常规示例说明并行层次;这不是 Velox 对 (worker, stage) 的全局唯一性约束。Task 由宿主以 taskId 管理,重试、attempt、执行策略及调度实现都可能改变映射。分析时应读取实际 taskId、fragment 与 attempt,而不能仅从 worker/stage 二元组推断。
| 说法 | 对错 |
|---|---|
| 一个 query 每个 worker 只收到一个 plan fragment | ❌ |
| 常规示例中一个 stage/worker 放置一个 task;是否唯一由宿主调度策略决定,并非 Velox 类型或接口强制。 | ✅ |
| 网络重试可以重复发送更新,task attempt 也可能重新调度;要结合 taskId 和协议检查是否复用已有 task,不能假设请求永不重复。 | ✅ |
常规放置示例:(worker, stage) → 一个 fragment 的一个 task
实际 taskId/attempt 由宿主管理;重试或其他策略可改变这一映射
(worker, query) → 可能多个 plan fragment → 多个 task(每参与一个 stage 一个)
同一 worker 可以同时运行一个 query 的多个 task;它们有各自的 Task 状态和子内存池,但通常共享该 worker 内的 QueryCtx、query root pool 与 executor。不同 task 不等于永久拥有不同的一组 OS 线程。
6.4 Driver 级:每个 task 内的并行靠 driver
Task 内的并行度不靠 task 数量,靠 driver 数量:
Worker A, Task-0-A
└── Pipeline 0
├── Driver 0 → 处理 split[0], split[3], split[6]...
├── Driver 1 → 处理 split[1], split[4], split[7]...
└── Driver 2 → 处理 split[2], split[5], split[8]...
(N = maxDrivers,通常 ≈ CPU 核数)
Coordinator 把属于 Worker A 的所有 splits 都增量塞进同一个 Task-0-A,task 内的多个 driver 从共享的 split queue 里抢着消费。
6.5 Splits 与 plan fragment 分开传输
Plan fragment 不包含 splits。两者通过同一个 POST /v1/task/{taskId} 的不同字段传输:
TaskUpdateRequest {
fragment: <Base64 PlanFragment> ← 算子树,只在首次创建 task 时发
sources: [ ← splits,运行时可多次增量追加
{ planNodeId: "0", ← 关联到 fragment 中的某个 source 节点
splits: [split1, split2, ...],
noMoreSplits: false } ← 全部下发完后置 true
]
}
- **Plan fragment 描述"怎么算"**:静态,发一次
- **Splits 描述"算什么数据"**:动态,分批追加,可能 task 启动后才陆续到达
- 在 Velox 侧:
Task::create(planFragment)一次性建树;addSplitWithSequence(planNodeId, split)增量喂数据;noMoreSplits(planNodeId)标记结束 - Source operator(
TableScan/Exchange)无 split 时返回kWaitForSplit阻塞,等待 coordinator 下发
6.6 并行度小结
| 维度 | 由谁决定 | 数量 |
|---|---|---|
| Stage 数 | query plan 的 shuffle 边界 | 由计划、exchange 边界与宿主优化决定,没有通用的固定范围。 |
| 每 stage 的 task 数 | coordinator 选中的 worker 数 | = 参与该 stage 的 worker 数 |
| 每 worker 的 task 数 | 该 worker 参与的 stage 数 | 每参与一个 stage 一个 |
| 每 task 的 pipeline 数 | plan 中 join 数 | 由 LocalPlanner 的 pipeline 切分决定;Join、LocalPartition/Exchange 等均可能形成边界,不能仅用 Join 数加一。 |
| 每 pipeline 的 driver 数 | maxDrivers 配置 |
由 maxDrivers、DriverFactory 与算子能力共同决定;可与 CPU 核数不同。 |
| 每 driver 的线程数 | 运行时至多一个 | 同一 Driver 不并发进入算子链;可多次退出并在不同 executor 线程续调,没有固定 OS 线程绑定。 |
7. Task 内部:Pipeline × Driver × Operator
7.1 Task 内部:Pipeline 与调度
以下是普通 pipeline 执行路径。最新 Task 还可把包含 FixedPointNode 的计划交给 FixedPointLoop,由它组织子 Task;不能把这个专门分支理解为 LocalPlanner 直接创建一个普通循环算子。普通查询的主循环仍由 Driver::runInternal 推进。空 getOutput 后再次检查 isBlocked;真正退出线程后才安装 future continuation,事件完成后重新 enqueue 并经过 Task::enter。
Velox LocalPlanner 把计划拆成 DriverFactory;factory 根据运行模式创建 Driver 和 Operator。HashJoin 的 build/probe、LocalPartition 的 producer/consumer 等边界需要不同的共享结构。
Driver 在 executor 上推进算子链。没有 split、connector IO 未完成、join build 未发布、输出缓冲已满,都可能产生 BlockingReason 和 future。future ready 后 Driver 重新入队,不代表原线程一直阻塞等候,也不保证下一次由同一线程执行。
Operator 可以暂时返回空输出,Driver 要结合 needsInput、isBlocked、noMoreInput、isFinished 决定动作。具体循环和取消逻辑见 Task & Driver;这里不把抽象推拉模型当成源码中的全部调度分支。
7.2 Task 内结构
Task
├── Pipeline 0 (probe side: Exchange → HashProbe → Agg → PartitionedOutput)
│ ├── Driver 0 ── 可被不同 executor 线程续调,处理若干 splits
│ ├── Driver 1
│ └── Driver N (N = maxDrivers,通常 ≈ CPU 核数)
│
└── Pipeline 1 (build side: Exchange → HashBuild)
└── Driver 0 (build side 通常只需 1 个 driver)
- Pipeline 数量由 plan 决定(每个 HashJoin/CrossJoin 会拆出 build/probe 两条 pipeline)
- Driver 数量(每条 pipeline 内的并行度)由
maxDrivers配置决定 - 每个 Driver 独立运行,处理不同的 split 或不同分区的数据
- 每个 Driver 内部单线程串行执行算子链
7.3 Driver 执行循环(Driver::runInternal)
Driver 通常从下游向上游检查可运行条件,拿到上游的 getOutput 结果后调用下游 addInput 传递。控制扫描方向和数据方向不同,也没有按传统 Volcano 形式由算子递归调用孩子 next。精确推进和重新扫描位置应结合 runInternal 的需求检查、停止原因和输出状态理解。
loop:
从 sink 向上游遍历每个 operator:
if 下游.needsInput():
output = 上游.getOutput()
if output: 下游.addInput(output)
if operator.isBlocked(&future):
// 注册 future,退出线程池;future 完成后重新入队
break
Operator 核心接口(velox/exec/Operator.h):
流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。
RowVectorPtr getOutput(); // 产生一批行(RowVector)
void addInput(RowVectorPtr input); // 消费上游输出
bool needsInput(); // 是否能接收更多输入
BlockingReason isBlocked(ContinueFuture* future); // 是否阻塞,设置 future
bool isFinished(); // 是否全部完成
数据流示例(Pipeline 0):
TableScan.getOutput() → RowVector(一批行)
↓ addInput
Filter.addInput/getOutput → 过滤后的行
↓ addInput
HashProbe.addInput → 关联 HashTable
HashProbe.getOutput() → join 结果行
↓ addInput
PartitionedOutput.addInput → 序列化,写入 OutputBuffer
7.4 阻塞与调度
Driver 运行于共享线程池(folly::CPUThreadPoolExecutor)。算子阻塞时(如等待 HashTable 建好、等待 Exchange 数据到来),Driver 持有 ContinueFuture 后退出线程;Future 完成后重新入队:
Driver → isBlocked() → 得到 future → 返回执行器,归还本次线程
↓ future 完成
Driver::enqueue() → 重新调度
阻塞原因(BlockingReason):
| 原因 | 触发场景 |
|---|---|
kWaitForSplit |
TableScan 等待 coordinator 下发 split |
kWaitForProducer |
Exchange/LocalExchange 等待上游数据 |
kWaitForConsumer |
PartitionedOutput OutputBuffer 已满 |
kWaitForJoinBuild |
HashProbe 等待 HashBuild 完成 |
| 内存仲裁的等待需按实际路径区分:普通算子 future blocked 与 ScopedMemoryArbitrationContext 下保留调用栈的 suspended 不是同一种协议。当前 BlockingReason 名称和 StopReason 见相邻源码链接。 | 内存仲裁,等待 spill 释放内存 |
7.5 Pipeline 间同步:HashJoinBridge
Pipeline 1 (build): TableScan → HashBuild ──写入──► HashJoinBridge
│ build done
Pipeline 0 (probe): Exchange → HashProbe ◄──读取──────────┘
HashBuild 完成后通过 HashJoinBridge 通知 HashProbe 的所有 Driver 解除阻塞。
8. Task 内数据交换:LocalPartition / LocalExchange
8.1 Task 内数据交换:LocalPartition / LocalExchange
LocalPartition 根据分区函数把输入行送入多个 LocalExchangeQueue。一般通过 dictionary indices 引用原始列;消费者 LocalExchange 从对应队列取出 RowVector,不走 HTTP 或页反序列化。
enqueue 的背压发生在数据入队与计入占用之后。超过 LocalExchangeMemoryManager 的门槛时 producer 得到等待 future;consumer 取走数据、降低占用后唤醒 producer。这个门槛不应被误读为物理内存绝不会瞬时越过的上限。
队列为空时 consumer 等待 producer;只有 producer 集合关闭、所有 producer 报告结束且队列排空,才是最终结束。dictionary 复用减少复制,但可能保留大块 base buffer,逻辑输出大小与 retained memory 不同。
适用场景:同一 task 内,不同 pipeline 之间需要重分区(如两阶段聚合、repartitioned join probe 侧重分区)。纯共享内存,无序列化,无网络。
"local" 的含义是 local to the task(同进程同 task),不是 local to the machine。
8.2 整体结构
Task 内部
Pipeline 0(上游,N 个 producer driver)
Driver 0 → ... → LocalPartition ┐
Driver 1 → ... → LocalPartition ┼─→ LocalExchangeQueue[0] ─→ LocalExchange → Driver 0
Driver 2 → ... → LocalPartition ┘─→ LocalExchangeQueue[1] ─→ LocalExchange → Driver 1
(按 key 哈希分桶) (按 partition 分桶) (source 算子)
Pipeline 1(下游,M 个 consumer driver,M = numPartitions)
- 每个
LocalExchangeQueue对应一个分区,有且只有一个 consumer driver 消费 - 同一组 LocalExchange 队列共享 LocalExchangeMemoryManager,按聚合的 bufferedBytes 触发 producer 背压;数据已入队才检查阈值,因此这不是不可瞬时越过的物理内存上限。
8.3 LocalExchangeMemoryManager
职责:跨所有 LocalExchangeQueue 追踪总内存用量,在超限时阻塞 producer。
流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。
// velox/exec/LocalPartition.h
int64_t maxBufferSize_; // 总缓冲上限
int64_t bufferedBytes_; // 当前已用字节(原子)
std::vector<ContinuePromise> promises_; // 被阻塞的 producer
- **
increaseMemoryUsage(added)**:producer 入队后调用。若bufferedBytes_ >= maxBufferSize_,创建ContinuePromise存入promises_,返回 true(producer 阻塞) - **
decreaseMemoryUsage(removed)**:consumer 出队后调用。若降到maxBufferSize_以下,取出所有promises_并 fulfill,唤醒所有阻塞的 producer
8.4 LocalExchangeQueue
职责:单个分区的 FIFO 队列,协调多 producer 写入 / 单 consumer 读取。
流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。
folly::Synchronized<Queue> queue_; // std::queue<pair<RowVectorPtr, int64_t>>
std::vector<ContinuePromise> consumerPromises_; // consumer 等待数据
int pendingProducers_{0}; // 尚未完成的 producer 数量
bool noMoreProducers_{false}; // 不再有新 producer
bool closed_{false};
入队(producer 调用):
BlockingReason enqueue(RowVectorPtr input, int64_t inputBytes, ContinueFuture* future) {
// 1. 追加到 queue_
// 2. 唤醒等待的 consumer(fulfil consumerPromises_)
// 3. 调用 memoryManager_->increaseMemoryUsage()
// → 若超限,返回 kWaitForConsumer(producer 阻塞)
// → 否则返回 kNotBlocked
}
出队(consumer 调用):
BlockingReason next(ContinueFuture* future, memory::MemoryPool*, RowVectorPtr* data, bool& drained) {
// 1. 队列为空:
// - 若所有 producer 已完成 → 返回 finished
// - 否则创建 consumerPromise → 返回 kWaitForProducer
// 2. 有数据:
// - 弹出队头
// - 调用 memoryManager_->decreaseMemoryUsage() → 可能唤醒 producer
// - 返回 kNotBlocked
}
Producer 生命周期管理:
| 方法 | 时机 | 作用 |
|---|---|---|
addProducer() |
producer 构造时 | pendingProducers_++ |
noMoreData() |
producer 处理完所有输入后 | pendingProducers_--;若归零且 noMoreProducers_ 则唤醒 consumer |
noMoreProducers() |
task 确认不再添加新 pipeline 后 | 设标志;若 pendingProducers_==0 则立即唤醒 consumer |
drain() |
barrier 处理(grouped execution) | drainedProducers_++;全部 drain 后唤醒 consumer |
8.5 LocalPartition 算子(Producer,sink)
构造:从 task 获取所有分区队列,并对每个队列调用 addProducer()。
addInput() 分区逻辑:
1. partitionFunction_->partition(*input, partitions_)
→ partitions_[i] = 第 i 行应去的分区号
2. 若所有行去同一个分区:直接调用 queue[p]->enqueue(input)
3. 若行分散到多个分区:
a. 统计每分区行数,构建索引数组 rawIndices_[partition]
b. 根据 singlePartitionBufferSize_ 判断用哪种策略:
- Buffer 模式(numPartitions 较多):
populatePartitionBuffer() 把行 copy 进 partitionBuffers_[p]
累积到 singlePartitionBufferSize_ * numPartitions 时统一 flush
- 即时模式(numPartitions 少或 eagerFlush):
wrapChildren() 用 dictionary 索引包装列,零拷贝创建 RowVector
立即 enqueue
wrapChildren() 通过 dictionary indices 引用输入列,省去这一分区步骤的 payload 复制;索引和包装对象仍有成本,还可能延长整个 base vector 的生命周期。与复制小分区到独立 buffer 的路径相比,选择应结合批次大小、分区数和 retained memory,不能只看逻辑输出字节数。
**copy() / populatePartitionBuffer()**:把行物理复制进 partitionBuffers_[p],适合分区数多(字典开销大)的场景。
阻塞:若任意队列 enqueue() 返回 kWaitForConsumer,收集对应 future,在 isBlocked() 时返回,driver 退出线程池。
**noMoreInput()**:flush 剩余 buffer,对所有队列调用 noMoreData(),通知 consumer 不再有数据。
8.6 LocalExchange 算子(Consumer,source)
LocalExchange 从 Task 为其分配的 partition 获取对应队列;partition 与 Driver 的映射由 DriverFactory、pipeline 和局部分区布局确定。不能把 driverId 取模当成所有路径通用的构造公式。
流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。
// getOutput()
BlockingReason reason = queue_->next(&future_, pool(), &data, drained);
if (reason != kNotBlocked) {
blockingReason_ = reason;
return nullptr; // driver 将在 isBlocked() 中返回 future_
}
return data;
isBlocked() 直接返回 blockingReason_ 和 future_,无额外逻辑。
9. Task 间数据交换:Remote Exchange
9.1 跨 Task 交换:从输出页到下游 RowVector
常规流式路径分成三段:
| 位置 | 关键对象 | 数据表示 |
|---|---|---|
| 上游 Task | PartitionedOutput、OutputBuffer | 按目标分区的序列化页 |
| 宿主传输 | PrestoExchangeSource、HTTP endpoints | 响应 body 与协议 headers |
| 下游 Task | InMemoryExchangeClient、ExchangeQueue、Exchange | 缓冲页,随后反序列化为 RowVector |
当前 InMemoryExchangeClient 管理 sources、请求和队列容量。类名中的 InMemory 不代表排除了远程网络;真正从哪里获取数据,由 ExchangeSource 的具体实现决定。
Presto 注册 HTTP 版本的 PrestoExchangeSource::request。客户端可先请求待取数据大小,再根据可用容量选 source;单 source 等条件下也有优化路径。不要把流控简化成“每个 source 永远固定并发拉取 N 个 page”。
OutputBufferManager 在当前 Velox 中是接口,默认实现位于 DefaultOutputBufferManager;它不是所有 exchange transport 的唯一实现。materialized exchange、file exchange 或其他 transport 需要另看对应路径。
本节展开的是 Presto 注册 PrestoExchangeSource 的常规 HTTP 流式交换路径,同机 task 也可以经由该路径传输。Velox 的 ExchangeSource 是可替换接口,还支持其他 transport/materialized/file exchange;跨 task 并不在接口层强制等于 HTTP。
9.2 整体数据流
9.3 上游:PartitionedOutput
addInput() 分区流程:
1. 估算每行序列化大小(用于 flush 时机判断)
2. partitionFunction_->partition(*input, partitions_)
3. 将每行的索引分配到 Destination[p].rows_
getOutput() 序列化 + flush 流程:
for each Destination:
advance():
收集 rows_ 中的行索引
用 VectorStreamGroup 序列化(Presto / CompactRow / UnsafeRow 格式)
累积到 maxPageSize 或行数阈值时 flush()
flush():
VectorStreamGroup → IOBufOutputStream → IOBuf chain(Presto Page 格式)
bufferManager_.enqueue(taskId, partition, page)
→ OutputBuffer::enqueue()
→ DestinationBuffer[p].data_.push_back(page)
→ 若 bufferedBytes_ 达到该输出缓冲实现的背压条件 → 返回 blocking future(producer 阻塞)
随机 flush 阈值(targetSizePct_ = 70~120%):各 Driver 的 flush 时机错开,避免同时写入造成突发压力。
9.4 上游:OutputBuffer / OutputBufferManager
流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。
// velox/exec/OutputBuffer.h
OutputBuffer {
std::vector<DestinationBuffer> buffers_; // 每个下游分区一个
int64_t maxSize_; // 总缓冲上限
int64_t continueSize_; // 恢复阈值(= maxSize_ * 0.9)
int64_t bufferedBytes_; // 当前已缓冲字节
std::vector<ContinuePromise> promises_; // 被阻塞的 producer
}
DestinationBuffer {
std::deque<SerializedPageBase> data_; // 待下游拉取的 pages
int64_t sequence_; // data_[0] 对应的全局序列号
std::function<...> notify_; // 有数据时唤醒等待的 HTTP 请求
}
背压:bufferedBytes_ 达到该输出缓冲实现的背压条件 时,enqueue() 返回 ContinueFuture,producer(PartitionedOutput)阻塞。当下游 ACK 消费数据后,bufferedBytes_ 降到 continueSize_ 以下时 fulfill promises,解除阻塞。
getData() 回调机制:下游 HTTP GET 到来时若无数据,注册 notify_ 回调;有数据入队时回调触发,立即响应 HTTP 请求。
9.5 下游:ExchangeClient
在这里的普通内存 Exchange 路径中,同一 pipeline 的消费者共享 Task 创建的 InMemoryExchangeClient,再从其队列消费。materialized / file exchange 有其他实现,不能把这项共享布局当成全部 exchange transport 的规定。
流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。
ExchangeClient {
ExchangeQueue queue_; // 多 source 入队,多 driver 消费
std::vector<ExchangeSource> sources_; // 每个上游 task 一个
int64_t maxQueuedBytes_; // 本地缓冲上限(默认 32MB)
int64_t minOutputBatchBytes_; // 最小返回批次(防止消费者空跑)
std::deque<ProducingSource> producingSources_; // 有剩余数据的 source
std::deque<ExchangeSource*> emptySources_; // 尚未探测数据量的 source
}
**addRemoteTaskId()**:为每个上游 task 创建一个 ExchangeSource,加入 emptySources_ 队列,启动首次数据量探测请求。
**next(consumerId, maxBytes)**(由 ExchangeNode.isBlocked() 调用):
1. queue_->dequeueLocked(maxBytes) → 若有足够数据立即返回
2. 若数据不足:pickSourcesToRequestLocked()
a. 对 emptySources_ 中的 source 发 requestDataSizes()(HEAD 请求,探测剩余量)
b. 按剩余容量(maxQueuedBytes_ - queue_.totalBytes() - totalPendingBytes_)
分配请求量给 producingSources_
c. 特殊处理:若单个 page 超过容量上限,强制发一次请求防止死锁
**ExchangeQueue**:
ExchangeQueue {
std::deque<SerializedPageBase> queue_;
int64_t totalBytes_;
int64_t minOutputBatchBytes_; // 消费者唤醒阈值
std::map<int, ContinuePromise> promises_; // consumerId → promise
}
**enqueueLocked()**(ExchangeSource 收到数据后调用):
nullptr入队表示某个 source 数据结束- 计算 unassigned bytes,若
>= minOutputBatchBytes_则 fulfill 等待中的消费者 promise minOutputBatchBytes_自适应:min(配置值, 已接收总量 / 100),防止小流量场景下消费者频繁空唤醒
**dequeueLocked(maxBytes)**:取出页面直到达到 maxBytes 或队列耗尽;数据不足时返回 ContinueFuture(消费者 driver 阻塞)。
9.6 下游:ExchangeSource 抽象接口
流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。
// velox/exec/ExchangeSource.h
ExchangeSource {
std::string remoteTaskId_;
int destination_; // 拉取上游哪个分区
ExchangeQueue* queue_;
int64_t sequence_; // 下次请求的序列号(可靠传输)
bool requestPending_; // 防止并发发起重复请求
bool atEnd_;
}
virtual Response request(uint32_t maxBytes, uint32_t maxWaitMs);
virtual void requestDataSizes(uint32_t maxWaitMs);
virtual void pause(); // 提示暂停(队列满时)
virtual void close();
Response 含:
bytes:本次响应数据量atEnd:source 是否已全部传完remainingBytes:上游 DestinationBuffer 中各 page 的剩余大小列表(用于 ExchangeClient 调度决策)
9.7 下游:PrestoExchangeSource(HTTP 实现)
shouldRequestLocked() → true(无 pending 请求且未 atEnd)
↓
request(maxBytes, maxWaitMs)
↓
doRequest(sequence_, maxBytes, maxWaitMs)
HTTP GET /v1/task/{remoteTaskId}/results/{destination}/{sequence_}
Headers: X-Presto-Max-Size: maxBytes
X-Presto-Max-Wait: maxWaitMs
↓
processDataResponse(response)
读 response headers:
X-Presto-Page-Next-Token → 更新 sequence_(下次请求用)
X-Presto-Buffer-Complete → 若 true,入队 nullptr(end marker)
X-Presto-Buffer-Remaining-Bytes → remainingBytes(返回给 ExchangeClient)
IOBuf chain → SerializedPage → queue_->enqueueLocked()
fulfill 请求 promise → ExchangeClient 调度下一批请求
重试:网络失败时指数退避(100ms → 10s,带 jitter),超过 exchangeMaxErrorDuration 后 queue_->setError(),task 报错。
9.8 下游:Exchange 算子
流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。
// isBlocked():
// 1. 若 currentPages_ 有数据 → kNotBlocked
// 2. getSplits() 获取上游 task ID → 调用 client_->addRemoteTaskId()
// 3. client_->next() → 若无数据返回 kWaitForProducer + future
// getOutput():
// 根据 serdeKind_ 分发:
// columnar(Presto): getOutputFromColumnarPages()
// row(CompactRow/UnsafeRow): getOutputFromRowPages()
**getOutputFromColumnarPages()**:用 VectorSerde::deserialize() 增量反序列化,每次读满 preferredOutputBatchBytes_ 后返回,跨 currentPages_ 中多个 page 累积。
行式页读取通过当前 Exchange 的相应 serde/row page 分支转换为 RowVector。原有流程片段用于说明“收集页面、创建行迭代器、批量解码”的阶段;具体函数名、输入页组织和批次限制应按当前 Exchange.cpp 核对。
9.9 序列号机制(可靠传输)
9.10 两种 sequence,不能混在一起
split sequence 是工作投递去重协议;exchange token 是某个输出 destination 的结果读取进度。二者编号空间、更新时机和处理对象不同。
上游 DestinationBuffer::acknowledge 可以释放 token 之前的页;带更大 sequence 的 getData 也隐式 acknowledge 更早的数据。下游读取响应的 next token,再更新自己的请求位置。空响应并不总携带 next token,当前 Presto 特别避免因缺失 header 把 sequence 重置。参见 next token 处理。
收到数据后先入 ExchangeQueue,complete 信号也会进入队列;锁外通知等待者。入队与通知 与 acknowledgeResults 形成缓冲回收协议。
这提供重取与确认的基础,但不能将其扩大成“整个分布式 query 自动获得持久化 exactly-once”。task 重试、物化结果和故障恢复仍由宿主策略决定。
ExchangeSource.sequence_ → HTTP GET .../results/{partition}/{sequence}
↓
上游 DestinationBuffer[partition].sequence_ → data_[0] 的全局编号
(验证请求序列号是否与当前可用数据对齐)
↓
响应 header X-Presto-Page-Next-Token → 下次请求的 sequence
(ExchangeSource 更新 sequence_ 用于下次请求)
ACK:GET .../results/{partition}/{sequence}/acknowledge → OutputBuffer::acknowledge() → 删除 sequence 之前的所有 pages,释放内存。
9.11 流控总结
| 层级 | 机制 | 上限 |
|---|---|---|
| OutputBuffer(上游) | bufferedBytes_ 达到该输出缓冲实现的背压条件 时阻塞 producer |
QueryConfig::maxOutputBufferSize() / max_output_buffer_size(默认 32MB) |
| ExchangeQueue(下游本地缓存) | totalBytes_ 控制发起的请求量 |
maxQueuedBytes_(默认 32MB) |
| ExchangeClient 请求调度 | maxQueuedBytes_ - totalBytes_ - pendingBytes = 可发请求量 |
动态计算 |
| 消费者唤醒阈值 | minOutputBatchBytes_ 防止小批次频繁唤醒 |
自适应(已收量 / 100) |
10. 内存管理
10.1 QueryCtx 和 Task:资源共享不是全局 Query 对象
Presto 的 QueryContextManager 按 queryId 缓存当前 Worker 内的 QueryCtx。该 query 在这个 Worker 上的多个 task 可以共享它;其他 Worker 有自己的进程、MemoryManager 和 QueryCtx,不存在一个跨机器共享指针。
findOrCreateQueryCtxLocked 在缓存未命中时调用:
MemoryManager::addRootPool(
queryId + 唯一 poolId,
queryConfig.queryMaxMemoryPerNode(),
...)
→ 将 pool 传给 QueryCtx::create
→ 缓存 QueryCtx
这也回答了“query pool 是不是在 Velox 层 add”:pool 的实现和 addRootPool API 在 Velox;在这条 Presto 路径中,由宿主 QueryContextManager 显式调用并传给 QueryCtx。 单独使用 Velox 而未传 pool 时,QueryCtx 的 Builder 也支持自动创建 root。两条路径都存在,不能只说其中一种。
TaskManager 在锁内建立 exec::Task::create,传入 fragment、QueryCtx、并行执行模式和 spill 选项,再启动执行。Task 的 aggregate pool 来自 query root;plan node 与 Operator 继续建立子层次。参见 Task 内存层次。
更多内存语义见 Memory Pool and Arbitrator。
内存池树把 leaf 的分配与上层聚合记账联系起来,回收行为由挂接在不同层级的 reclaimer 委托到算子。各层并不都是独立分配器,也不都拥有可自行释放的 payload;query root capacity 的再分配还由仲裁器管理。
QueryCtx::MemoryPool(query 级,所有 task 共享上限)
└── Task::MemoryPool(task 级)
├── Operator pools(每个 operator 实例独立池)
├── ExchangeClient pool(每条 pipeline 一个)
└── LocalExchange pool
Task::MemoryReclaimer在内存压力下触发 spill(HashAggregation、HashJoin 等支持 spill 的算子)- Operator 通过
addOperatorPool(planNodeId, splitGroupId, pipelineId, driverId)在构造时申请子池
11. 从控制信息、数据流与资源归属看 Worker 集成
Worker 集成需要同时推进三条不同的路径:控制更新告诉系统还应执行什么,数据流提供实际输入和输出,资源管理约束哪些执行可以继续。将它们分开能支持逐步投递和异步执行,但也意味着完成条件不能只靠一个消息或一个计数器判断。
| 设计选择 | 带来的能力 | 对宿主与 Velox 协作的要求 |
|---|---|---|
| fragment 执行与分布式调度分开 | 执行库可以服务不同宿主,宿主保留调度策略 | 计划转换、配置映射和生命周期边界必须明确 |
| split 增量投递 | 输入发现与执行可以重叠,不要求先收集全部工作 | 去重、无更多输入信号、重试和终止由各自协议处理 |
| 本地与远程 exchange 分别实现 | 本地批次交换可以避免承担完整网络协议成本 | 分区语义要一致,远程还需面对序列化、缓冲和上游结束 |
| 异步等待与缓冲反压 | 等待资源或数据时释放执行线程,并限制在途数据 | 唤醒、取消、缓冲释放和生产者结束都要有可达路径 |
| 查询资源上下文与 Task 执行实例分开 | 多个 Task 可以共享查询配置和资源归属 | 共享上下文的存活期不应由单个 Task 的状态简单推断 |
这里的设计品味体现在归属清晰:某个状态由谁写、哪个事件使谁继续、什么条件允许释放缓冲。抽象层次只有在这些问题有明确答案时才有价值。单凭类名相似,不能认为 Presto Task、Velox Task 和 QueryCtx 是同一个层次的执行对象。
排查时也应沿同样边界取证:先确认宿主发出了哪些任务和 split,再看 Velox 的 Driver 正在运行还是等待,最后追踪等待条件对应的数据与资源。增加 Driver 数只能作用于其中一部分;当输入供给、带宽或输出消费受限时,更多并发还可能增加在途内存。
这条端到端链有三种不同的进度标识:split sequence 防止输入工作重复投递,exchange token 标识已传输和确认的数据页位置,Operator 状态表示本地计算推进到哪里。把它们分开可以独立处理重试、背压和计算生命周期;代价是取消必须穿过每一个边界,不能只设置一个 Query 状态就认为所有远端请求与本地 waiter 已经收束。
12. 关键代码位置
12.1 版本变化与阅读入口
阅读时区分 Presto 的计划 converter 与 Velox 的 LocalPlanner、OutputBufferManager 接口与实现、split sequence 与 transport token。Worker 转换计划并推进 split watermark;Velox 在 Task 内组织 Driver 与交换状态;网络传输由宿主实现的 ExchangeSource 接入。
建议按以下源码顺序读:TaskResource → TaskManager → Velox Task → LocalPlanner → Local Exchange → OutputBuffer → PrestoExchangeSource。
本文完成了两套固定提交的源码核对,没有启动 Presto 集群或执行端到端查询;网络协议和版本适配的结论限于上述代码路径。
| 组件 | 文件 |
|---|---|
| HTTP 端点 | presto_cpp/main/TaskResource.cpp |
| Task 生命周期管理 | presto_cpp/main/TaskManager.cpp |
| Plan 转换 | presto_cpp/main/PrestoToVeloxQueryPlan.cpp |
| Velox Task | velox/exec/Task.h, velox/exec/Task.cpp |
| Driver 执行循环 | velox/exec/Driver.cpp → runInternal() |
| Operator 接口 | velox/exec/Operator.h |
| LocalPartition/LocalExchange | velox/exec/LocalPartition.h, velox/exec/LocalPartition.cpp |
| Remote Exchange 算子 | velox/exec/Exchange.h, velox/exec/Exchange.cpp |
| ExchangeClient + ExchangeQueue | velox/exec/ExchangeClient.h/cpp, velox/exec/ExchangeQueue.h/cpp |
| ExchangeSource 抽象接口 | velox/exec/ExchangeSource.h |
| HTTP 数据拉取实现 | presto_cpp/main/PrestoExchangeSource.h/cpp |
| PartitionedOutput | velox/exec/PartitionedOutput.h, velox/exec/PartitionedOutput.cpp |
| OutputBuffer / Manager | velox/exec/OutputBuffer.h/cpp, velox/exec/OutputBufferManager.h/cpp |
源码核对(2026-09-20):本轮按 Velox 1d1b76567870 核对关键接口、控制流、默认值与边界条件。当前源码摘录附固定版本链接;流程伪代码用于说明分支,不是可直接编译的程序。未对全文示例做独立编译或性能复测。涉及宿主集成与历史实验的数据,按各节标注的来源理解。