1. 1. 1. 分布式 Query 与 Velox 执行片段的边界
  2. 2. 2. 层级结构总览
  3. 3. 3. 完整端到端数据流
    1. 3.1. 3.1 背压、结束和失败是同一条生命周期链
  4. 4. 4. Worker 接收查询:HTTP → Velox Task
    1. 4.1. 4.1 Worker 入口:HTTP 更新如何变成 Velox 计划
    2. 4.2. 4.2 HTTP 端点(TaskResource)
    3. 4.3. 4.3 Plan Fragment 转换(PrestoToVeloxQueryPlan)
    4. 4.4. 4.4 Task 创建流程(TaskManager)
  5. 5. 5. Split 分发与消费机制
    1. 5.1. 5.1 Split 分发:批次、watermark 与结束信号
    2. 5.2. 5.2 全链路概览
    3. 5.3. 5.3 第 1 步:Worker 接收与转换
    4. 5.4. 5.4 第 2 步:addSplitWithSequence —— 为什么需要 sequence
    5. 5.5. 5.5 第 3 步:Velox 侧 split 队列结构
      1. 5.5.1. 5.5.1 为什么拆成 SplitsState / SplitsStore 两层
      2. 5.5.2. 5.5.2 非 grouped 执行(ungrouped)长什么样
    6. 5.6. 5.6 第 4 步:Source operator 领取 split
    7. 5.7. 5.7 第 5 步:noMoreSplits 的两个层级 + 延迟处理
    8. 5.8. 5.8 第 6 步:Split group(bucketed/grouped 执行)
    9. 5.9. 5.9 Split 分发关键代码位置
  6. 6. 6. Query 如何并行执行
    1. 6.1. 6.1 Query 的三个并行尺度
    2. 6.2. 6.2 Stage 级:query 切成多个 plan fragment
    3. 6.3. 6.3 Task 级:每个 stage 分发到多个 worker
    4. 6.4. 6.4 Driver 级:每个 task 内的并行靠 driver
    5. 6.5. 6.5 Splits 与 plan fragment 分开传输
    6. 6.6. 6.6 并行度小结
  7. 7. 7. Task 内部:Pipeline × Driver × Operator
    1. 7.1. 7.1 Task 内部:Pipeline 与调度
    2. 7.2. 7.2 Task 内结构
    3. 7.3. 7.3 Driver 执行循环(Driver::runInternal)
    4. 7.4. 7.4 阻塞与调度
    5. 7.5. 7.5 Pipeline 间同步:HashJoinBridge
  8. 8. 8. Task 内数据交换:LocalPartition / LocalExchange
    1. 8.1. 8.1 Task 内数据交换:LocalPartition / LocalExchange
    2. 8.2. 8.2 整体结构
    3. 8.3. 8.3 LocalExchangeMemoryManager
    4. 8.4. 8.4 LocalExchangeQueue
    5. 8.5. 8.5 LocalPartition 算子(Producer,sink)
    6. 8.6. 8.6 LocalExchange 算子(Consumer,source)
  9. 9. 9. Task 间数据交换:Remote Exchange
    1. 9.1. 9.1 跨 Task 交换:从输出页到下游 RowVector
    2. 9.2. 9.2 整体数据流
    3. 9.3. 9.3 上游:PartitionedOutput
    4. 9.4. 9.4 上游:OutputBuffer / OutputBufferManager
    5. 9.5. 9.5 下游:ExchangeClient
    6. 9.6. 9.6 下游:ExchangeSource 抽象接口
    7. 9.7. 9.7 下游:PrestoExchangeSource(HTTP 实现)
    8. 9.8. 9.8 下游:Exchange 算子
    9. 9.9. 9.9 序列号机制(可靠传输)
    10. 9.10. 9.10 两种 sequence,不能混在一起
    11. 9.11. 9.11 流控总结
  10. 10. 10. 内存管理
    1. 10.1. 10.1 QueryCtx 和 Task:资源共享不是全局 Query 对象
  11. 11. 11. 从控制信息、数据流与资源归属看 Worker 集成
  12. 12. 12. 关键代码位置
    1. 12.1. 12.1 版本变化与阅读入口
Macduan Notes

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、输入分发、执行与交换的顺序展开实现。并行度和内存限制放回各自作用的执行边界来解释。

图 1:Presto 管理分布式 Query/Stage 和协议,Velox 管理 fragment 内的执行与资源。
图 1:Presto 管理分布式 Query/Stage 和协议,Velox 管理 fragment 内的执行与资源。

本文于 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 保证。

图 2:计划与 split 持续更新;批次 watermark、队列等待和 noMore 信号承担不同职责。
图 2:计划与 split 持续更新;批次 watermark、队列等待和 noMore 信号承担不同职责。

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 全链路概览

宿主如何给 Task 添加 Split
图 3:宿主如何给 Task 添加 Split。已按当前实现修正标注,具体约束见相邻正文。
Fig. Worker split 处理管线:从 Coordinator 的 POST 到 TableScan 取 split

边界说明:哪个 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() 处理两阶段:

  1. 合并:同一请求里可能有多个相同 planNodeId 的 source,先按 planNodeId 合并 splits、noMoreSplitsForLifespan,noMoreSplits 做 OR
  2. 转换:toVeloxSplit(ScheduledSplit) 把 Presto 协议 split 转成 velox::exec::Split
    • RemoteSplit → 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
图 4:数据从上游输出缓冲流向下游;拉取、acknowledge 与背压沿反向协调。
图 4:数据从上游输出缓冲流向下游;拉取、acknowledge 与背压沿反向协调。

当前 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 整体数据流

远程 Exchange 的两端状态
图 5:远程 Exchange 的两端状态。已按当前实现修正标注,具体约束见相邻正文。
Fig. 上下游 Exchange:下游 ExchangeSource 通过 HTTP GET 拉上游 PartitionedOutput 的分区数据

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