1. 1. 1. Task、Pipeline、Driver 与线程的职责边界
  2. 2. 2. runInternal 调度实例:三算子 pipeline 全过程
    1. 2.1. 2.1 算子接口语义
    2. 2.2. 2.2 阶段一:冷启动——从 sink 往 source "找活干"
    3. 2.3. 2.3 阶段二:数据推进——i += 2 的精妙之处
    4. 2.4. 2.4 阶段三:稳态循环
    5. 2.5. 2.5 阶段四:阻塞场景——blockDriver 让出线程
      1. 2.5.1. 2.5.1 场景 A:OutputBuffer 满(下游消费跟不上)
      2. 2.5.2. 2.5.2 场景 B:TableScan 等 split
    6. 2.6. 2.6 阶段五:收尾——noMoreInput 链式传播 → close
    7. 2.7. 2.7 时间线汇总
  3. 3. 3. Task 的设计与实现
    1. 3.1. 3.1 执行架构:谁管理谁,谁推进数据
    2. 3.2. 3.2 Task 是什么
    3. 3.3. 3.3 Task 状态机
    4. 3.4. 3.4 核心数据成员
    5. 3.5. 3.5 Task 的执行模式
    6. 3.6. 3.6 Task 启动流程
    7. 3.7. 3.7 从 PlanFragment 到可调度的 Driver
      1. 3.7.1. 3.7.1 Pipeline 在哪里切开
      2. 3.7.2. 3.7.2 并行模式的启动顺序
      3. 3.7.3. 3.7.3 Serial 模式仍可能有多条 pipeline
    8. 3.8. 3.8 Split 分发
    9. 3.9. 3.9 allPeersFinished —— Hash Join Build 同步屏障
    10. 3.10. 3.10 多个 Driver 如何协调
      1. 3.10.1. 3.10.1 Split 分发与唤醒粒度
      2. 3.10.2. 3.10.2 Grouped execution 限制同时活跃的 split group
      3. 3.10.3. 3.10.3 Join Build 的 peer barrier 与 JoinBridge
    11. 3.11. 3.11 Task::terminate —— 统一终止路径
    12. 3.12. 3.12 普通 pipeline 与 FixedPointLoop 的执行边界
  4. 4. 4. runInternal 的执行模型:不是 Volcano,也不是教科书式 push
    1. 4.1. 4.1 数据如何在一条 pipeline 中前进
      1. 4.1.1. 4.1.1 算子接口表达不同的事情
      2. 4.1.2. 4.1.2 从下游检查需求,拿到一批后向下游推进
      3. 4.1.3. 4.1.3 Sink 与结束传播
      4. 4.1.4. 4.1.4 第二次 isBlocked 把“没有输出”变成可调度的等待
    2. 4.2. 4.2 控制流:Driver 是唯一的调度器
    3. 4.3. 4.3 调度扫描方向:从 sink 往 source 找"谁能干活"
    4. 4.4. 4.4 demand-gated(需求门控)
    5. 4.5. 4.5 为什么"扁平循环 + Driver 调度"是协作式调度的前提
    6. 4.6. 4.6 小结
  5. 5. 5. Driver 的设计与实现
    1. 5.1. 5.1 Driver 是什么
    2. 5.2. 5.2 Driver 的核心组成
    3. 5.3. 5.3 ThreadState —— Driver 的状态
    4. 5.4. 5.4 Driver 的调度状态与 Task 的暂停协议
      1. 5.4.1. 5.4.1 ThreadState 是一组字段,不是单一状态枚举
      2. 5.4.2. 5.4.2 enter / leave 保护的是执行资格与线程计数
      3. 5.4.3. 5.4.3 Yield、Pause、Suspend 分别解决什么问题
      4. 5.4.4. 5.4.4 Suspended 保留线程,但暂时允许 Task 停稳
    5. 5.5. 5.5 Driver 状态流转
    6. 5.6. 5.6 Driver::runInternal —— 核心执行循环
    7. 5.7. 5.7 blockDriver —— 进入 Blocked 状态
    8. 5.8. 5.8 CancelGuard —— 线程安全的清理保证
    9. 5.9. 5.9 Driver::close —— 正常关闭路径
  6. 6. 6. isBlocked 与三种让出机制:结合典型算子解读
    1. 6.1. 6.1 阻塞与唤醒:future 怎样接回执行循环
      1. 6.1.1. 6.1.1 先退出 Driver,再安装恢复回调
      2. 6.1.2. 6.1.2 等待什么,由算子决定
      3. 6.1.3. 6.1.3 从 off-thread 到 on-thread:完整的调度交接
      4. 6.1.4. 6.1.4 future ready、queued 和 on-thread 是三个时刻
    2. 6.2. 6.2 isBlocked 的角色:声明阻塞意愿 + 交出唤醒句柄
    3. 6.3. 6.3 三种让出线程的方式
    4. 6.4. 6.4 典型算子的 isBlocked 解读
      1. 6.4.1. 6.4.1 (1) CallbackSink —— sink 的下游反压(延迟报告模式)
      2. 6.4.2. 6.4.2 (2) LocalExchange —— 本地交换 source(getOutput 探测 + isBlocked 报告)
      3. 6.4.3. 6.4.3 (3) Exchange —— 远程交换 source(isBlocked 里干实活 + collectAny)
      4. 6.4.4. 6.4.4 (4) Merge / LocalMerge —— 多源归并(kWaitForProducer)
      5. 6.4.5. 6.4.5 (5) LocalPartition —— 多下游 sink(collectAll)
      6. 6.4.6. 6.4.6 (6) HashProbe —— join probe(状态机驱动 + kWaitForJoinBuild)
    5. 6.5. 6.5 两种 future 组合语义:collectAny vs collectAll
    6. 6.6. 6.6 设计模式总结
  7. 7. 7. Task 与 Driver 的交互与生命周期
    1. 7.1. 7.1 生命周期:引用、算子资源与内存池
      1. 7.1.1. 7.1.1 Task 与 Driver 为什么需要显式拆开引用
      2. 7.1.2. 7.1.2 Future 的完成责任属于产生它的组件
      3. 7.1.3. 7.1.3 MemoryPool 的生命期可以长于 Operator
    2. 7.2. 7.2 引用关系图
    3. 7.3. 7.3 完整生命周期
    4. 7.4. 7.4 Task::enter / leave —— Driver 线程注册协议
    5. 7.5. 7.5 Suspended 状态的用途
    6. 7.6. 7.6 Split Group 并发控制
  8. 8. 8. Task 的资源管理:terminate 流程与资源回收
    1. 8.1. 8.1 Task 何时结束,终态之后还要做什么
      1. 8.1.1. 8.1.1 Finished 的条件包含输出消费
      2. 8.1.2. 8.1.2 terminate 先决定终态,再分路径清理
      3. 8.1.3. 8.1.3 三个完成边界
    2. 8.2. 8.2 Memory Pool 树:层级、所有权与"刻意保活"
    3. 8.3. 8.3 Happy Path:自然完成
    4. 8.4. 8.4 异常 Path:错误 / 取消 / 中止
    5. 8.5. 8.5 terminate:统一清理漏斗(happy/error 共用)
    6. 8.6. 8.6 最终回收:~Task()
    7. 8.7. 8.7 资源回收时机汇总
    8. 8.8. 8.8 设计要点
  9. 9. 9. Future/Promise 管理:避免 Driver 泄漏的设计与实现
    1. 9.1. 9.1 带着状态与等待关系排查问题
    2. 9.2. 9.2 附录:源码阅读路线
    3. 9.3. 9.3 为什么 future/promise 在这里特别危险
    4. 9.4. 9.4 全景:Task/Driver 中的 promise/future
    5. 9.5. 9.5 陷阱一:持锁兑现 promise → 死锁
    6. 9.6. 9.6 陷阱二:future 永不兑现 → Driver 泄漏 → terminate 作"总清算"
    7. 9.7. 9.7 陷阱三:double-resume / 两个线程进同一个 Driver
    8. 9.8. 9.8 可观测性:让泄漏"看得见"
    9. 9.9. 9.9 设计哲学总结
  10. 10. 10. 为什么产出统一从 getOutput 发起
    1. 10.1. 10.1 核心原因:解耦"消费输入速率"与"产生输出速率"
    2. 10.2. 10.2 getOutput 空返回与阻塞、完成分别判断
    3. 10.3. 10.3 统一性:source / 中间 / 攒批算子,Driver 一视同仁
    4. 10.4. 10.4 与 demand-gated push 的关系:pull 触发,push 传递
    5. 10.5. 10.5 小结
  11. 11. 11. 架构设计品味、代码品味与多线程编程范式
    1. 11.1. 11.1 架构设计品味
      1. 11.1.1. 11.1.1 "协作式调度"而非"抢占式线程"
      2. 11.1.2. 11.1.2 单一终止漏斗(terminate as a funnel)
      3. 11.1.3. 11.1.3 所有权方向与生命周期的"单调性"
      4. 11.1.4. 11.1.4 循环引用的"受控泄漏"
      5. 11.1.5. 11.1.5 窄接口统一异构等待——isBlocked 作为唯一的"挂起-唤醒"抽象
    2. 11.2. 11.2 代码品味
      1. 11.2.1. 11.2.1 快路径无锁,慢路径加锁
      2. 11.2.2. 11.2.2 "Locked 后缀"约定——把锁契约编码进函数名
      3. 11.2.3. 11.2.3 RAII 兜底一切
      4. 11.2.4. 11.2.4 统计的双保险
      5. 11.2.5. 11.2.5 主循环写得极致紧凑
      6. 11.2.6. 11.2.6 延迟报告模式——阻塞在数据操作里发现,在下一轮 isBlocked 报告
      7. 11.2.7. 11.2.7 future 组合语义必须匹配等待语义——collectAny vs collectAll
    3. 11.3. 11.3 涉及的多线程编程范式
      1. 11.3.1. 11.3.1 一条贯穿全局的主线
  12. 12. 12. 从执行进度、等待关系与所有权看调度设计
  13. 13. 13. 附录:TaskDriverOperatorLifecycle.md 中文翻译
  • Task、Driver、Operator 生命周期
    1. 0.1. 13.1 OutputBufferManager
    2. 0.2. 13.2 Operator
    3. 0.3. 13.3 Driver 与 Task
  • Macduan Notes

    Velox Task & Driver

    Velox 的 Task 管理一个计划片段的执行,Driver 则推进其中一条 pipeline 的算子链。LocalPlanner 将计划划成 pipeline,每条链可以有多份 Driver 实例;它们共享 Executor 的工作线程,以 RowVectorPtr 传递批次,并在等待输入或下游容量时保存进度、退出线程。

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

    1. Task、Pipeline、Driver 与线程的职责边界

    Task 是一个计划片段的执行实例;Driver 是推进一条 pipeline 中算子链的执行实例。Pipeline 描述可以连续传递批次的一段处理流程,Operator 完成具体操作,Executor 提供运行 Driver 的线程资源。一个 Driver 包含自己的算子实例,但并不固定占有一个线程。

    对象 它表示什么 生命周期或执行边界
    QueryCtx 查询配置和共享资源上下文 可被多个 Task 共同引用;不是一条算子执行链
    Task 一个 PlanFragment 的执行与协调 管理 split、Driver、跨 pipeline 共享状态和结束流程
    Pipeline / DriverFactory 一段算子链及其创建方式 同一 pipeline 可以创建多个 Driver
    Driver 一份可暂停、可继续的执行进度 当前算子、批次和等待状态跨线程调度保留
    Operator 扫描、过滤、聚合、Join、输出等处理单元 接收输入、提供输出、声明阻塞与完成状态
    Executor 线程 实际执行 C++ 指令的资源 Driver 阻塞或让出后,可以运行其他可执行工作

    数据沿算子链以 RowVectorPtr 传递;调度需要同时处理“有数据可取”“下游可接收”“等待外部事件”“输入已结束”等状态。返回空批次不自动等于算子结束,阻塞也不等于线程应当原地等待。后面的 needsInput、getOutput、isBlocked 和 noMoreInput 分析,都在解释这些接口怎样共同表达执行进度。

    Task 的终态、执行线程停止和对象最终析构也要分开。取消请求需要先阻止继续执行、唤醒或结束等待、关闭算子并释放引用,资源才能最终回收。先建立这些概念,才容易理解 runInternal 的循环、future 回调和 terminate 中看似分散的保护逻辑。

    Task 管理计划片段、Driver 和共享状态;每个 Driver 拥有自己的算子实例,Executor 在线程上运行 Driver。
    图 1:Task 管理计划片段、Driver 和共享状态;每个 Driver 拥有自己的算子实例,Executor 在线程上运行 Driver。 打开原图

    本文对照 Velox 1d1b765678702e6b5d1b3c582373c618813c3c85(2026-09-17),于 2026-09-19 复核。源码链接固定到该提交;图示侧重普通 CPU 执行路径,特殊的 trace、barrier 和自定义算子分支在相关位置补充说明。

    代码片段分别标明源码节选或流程示意;流程示意省略统计、异常包装与无关分支,不是可独立编译的程序。历史资料与宿主集成保留各自版本,不能据此推断它们组成了经过构建验证的发行版本。

    2. runInternal 调度实例:三算子 pipeline 全过程

    先跟随一次普通并行扫描走完启动、取数、等待、恢复和结束,再展开各对象的字段与实现。假设查询从一个 split 读取列 x,过滤 x > 10,输出 x + 1;某一批输入为 [4, 12, 9, 20],过滤与投影后得到 [13, 21]。这里的四行是教学输入,不依赖具体文件格式。

    阶段谁推进数据或状态怎样交接
    建立执行实例宿主 → Task::start → LocalPlanner → DriverFactory创建 Scan → FilterProject → PartitionedOutput 三个算子,Driver::enqueue 提交执行项。
    获得执行资格Executor worker → Driver::run → Task::enterTask 允许进入后才能执行算子;被排队不等于已经 on-thread。
    需求驱动取数Driver 从 sink 检查到 sourceScan 交出四行,FilterProject 交出 [13, 21],PartitionedOutput 将结果序列化到输出缓冲。
    没有 split 或下游满算子 → isBlocked → blockDriver交出有效 future;runInternal 返回后 CancelGuard 析构调用 Task::leave 撤销线程登记,再安装恢复回调。
    事件完成并恢复promise → continuation → enqueue → worker → Task::enterfuture 成功只使 Driver 可以重新竞争执行;pause 时交由 Task::resume,异常则进入 Task error。
    完成及释放noMoreInput → sink isFinished → Driver::close关闭算子并移除 Driver;带 PartitionedOutput 的 Task 还要满足输出消费条件,最终拆开引用并释放资源。

    后面的逐步推演展开同一条链。完整分支见 Driver::runInternal;执行实例的建立见 Task::start。

    用一条最经典的流水线 TableScan → FilterProject → PartitionedOutput 来逐步 trace runInternal(Driver.cpp:519)的调度过程。

    operators_[0] = TableScan         (source,从 split 读数据)
    operators_[1] = FilterProject     (中间算子,过滤+投影,逐批无状态)
    operators_[2] = PartitionedOutput (sink,写 OutputBuffer)
    N = 3,startingOperator = N-1 = 2
    

    2.1 算子接口语义

    方法 TableScan FilterProject PartitionedOutput
    isBlocked() 无 split 可读且 noMoreSplits 未到 → kWaitForSplit 不阻塞 OutputBuffer 满 → kWaitForConsumer
    needsInput() —(source 无上游) 缓存空时 true buffer 未满时 true
    getOutput() 读出一批;读完返回 null 处理完输入产出一批;无输入返回 null 并行模式写 buffer 后返回 null
    addInput() —— 接收一批待处理 接收一批待写出
    isFinished() noMoreSplits 且无数据 → true 收到 noMoreInput 且缓存空 → true noMoreInput 且全部写出 → true

    2.2 阶段一:冷启动——从 sink 往 source "找活干"

    Driver 刚上线,所有算子都还没有数据。内层循环 for (i = 2; i >= 0; --i):

    i = 2(sink, PartitionedOutput) → i == N-1,走 else 分支

    getOutput(sink) → null   // 还没人喂它
    isFinished()?   → false
    continue → --i → i=1
    

    i = 1(FilterProject) → i < N-1,走 if 分支

    nextOp = operators_[2] (sink)
    nextOp->isBlocked()?  → 否
    nextOp->needsInput()? → true   // sink 要数据
    getOutput(FilterProject) → null   // FilterProject 还没收到输入
      → else 分支:isFinished(FilterProject)? → false
      → 既无数据也未结束,不 break,继续 --i → i=0
    

    i = 0(TableScan) → if 分支

    nextOp = operators_[1] (FilterProject)
    nextOp->needsInput()? → true
    getOutput(TableScan) → batch1   // 终于从 split 读出真实数据!
      intermediateResult 非空:
      addInput(FilterProject, batch1)   // ← 数据被 push 给 FilterProject
      i += 2 → i=2; continue → --i → i=1
    

    关键点:冷启动时扫描从 sink 一路走到 source(i: 2→1→0),目的是**自下而上确认"谁此刻能产出数据"**。 最终落到 TableScan 这个真正的数据来源,拿到第一批数据。

    2.3 阶段二:数据推进——i += 2 的精妙之处

    刚才在 i=0 喂完 FilterProject,i += 2 然后 for 的 --i,净效果 i = 1,回头去看刚收到数据的 FilterProject:

    i = 1(FilterProject)再次

    nextOp = sink, needsInput? → true
    getOutput(FilterProject) → batch1'   // 它有输入了,过滤投影后产出
    addInput(sink, batch1')              // ← 数据 push 给 sink
    i += 2 → i=3; continue → --i → i=2
    

    i = 2(sink)

    getOutput(sink) → null   // PartitionedOutput 把 batch1' 写进 OutputBuffer,并行模式不返回数据
    isFinished()? → false
    continue → --i → i=1
    

    i += 2 的设计意图:一批数据刚喂给 operators_[i+1],下一步立刻回到 operators_[i+1] 看它能否把这批数据继续往下游推。 这样一批数据会被"连续地"沿 pipeline 推到底,而不是每次都从头扫描。数据流向始终是 0 → 1 → 2(push),而扫描指针在局部上下跳动。

    2.4 阶段三:稳态循环

    batch1 推到底后,i 回到 1,FilterProject 已无缓存数据:

    i=1: getOutput(FilterProject) → null → isFinished? false → --i → i=0
    i=0: getOutput(TableScan) → batch2 → addInput(FilterProject) → 推进...
    

    于是稳定在 i=1 ↔ i=0 ↔ i=2 之间往复,每轮从 TableScan 拉一批、过滤、写出。**注意每个算子调用前循环顶部都会执行 task()->shouldStop()(检查终止/暂停/yield)和 op->isBlocked()**——这是协作式调度的轮询点。

    2.5 阶段四:阻塞场景——blockDriver 让出线程

    2.5.1 场景 A:OutputBuffer 满(下游消费跟不上)

    某轮 i=1,准备喂 sink 时:

    流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

    i=1: nextOp = sink
         nextOp->isBlocked(&future) → kWaitForConsumer   // buffer 满了
         blockingReason_ != kNotBlocked
         → return blockDriver(self, /*blockedOperatorId=*/i+1=2, future, ...)
    

    blockDriver 构造 BlockingState 并设置 hasBlockingFuture,随后返回 kBlock。先退出 runInternal,由 CancelGuard 调用 Task::leave,使 Driver off-thread;外层 Driver::run 才安装 BlockingState::setResume。工作线程随后可以执行其他工作。

    当下游消费了 buffer → future 兑现 → setResume 回调(Driver.cpp:227)→ Driver::enqueue 重新入队 → 再次进入 runInternal,从 startingOperator=2 重新扫描。

    2.5.2 场景 B:TableScan 等 split

    某轮 i=0,循环顶部:

    流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

    i=0: op = TableScan
         op->isBlocked(&future) → kWaitForSplit   // 无 split 且 noMoreSplits 未到
         → return blockDriver(self, /*blockedOperatorId=*/0, future, ...)
    

    同样下线,等 Task::addSplit() 兑现 promise 后被唤醒。

    这两个就是第 4 章讲的两类反压:B 是上游缺数据,A 是下游消费慢,都通过同一套 kBlock + ContinueFuture 机制把整个 Driver 挂起。

    2.6 阶段五:收尾——noMoreInput 链式传播 → close

    假设 noMoreSplits 已到,TableScan 读完最后一批。某轮 i=0:

    流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

    i=0: nextOp=FilterProject, needsInput? true
         getOutput(TableScan) → null              // 没数据了
         → else 分支:
           isBlocked(TableScan)? → 否              // 不是阻塞,是真的结束
           isFinished(TableScan)? → true
           nextOp->noMoreInput()  → FilterProject.noMoreInput()  // 通知下游:上游收工
           break                  → 跳出内层 for,回到 for(;;),重新 i=2
    

    重新 i=2(sink)

    getOutput(sink) → null; isFinished? false → --i → i=1
    

    i=1(FilterProject)

    nextOp=sink, needsInput? true
    getOutput(FilterProject) → null              // 收到 noMoreInput,已无缓存可 flush
    → else:isFinished(FilterProject)? → true     // 它也结束了
      nextOp->noMoreInput() → PartitionedOutput.noMoreInput()  // 继续向下游传播
      break → 回到 for(;;),i=2
    

    i=2(sink)

    getOutput(sink) → null     // PartitionedOutput flush 剩余、标记 buffer 结束
    isFinished(PartitionedOutput)? → true
    close();                   // 关闭所有算子、上报统计、Task::removeDriver()
    return StopReason::kAtEnd;  // driver 正常终结
    

    kAtEnd 让 Driver::run 直接 return(Driver.cpp:877),driver 彻底退出。close() 里的 Task::removeDriver 会触发 checkIfFinishedLocked,若是最后一个 driver 且 output 已被消费,则 Task::terminate(kFinished)。

    结束信号的传播链很清晰:noMoreInput 沿 pipeline 自上游向下游逐级传递 (TableScan→FilterProject→PartitionedOutput),每个算子先把自己的剩余缓存 flush 完(isFinished 返回 true 前的 getOutput),再把"收工"信号交给下一级。每传一级就 break 重启外层循环,确保从 sink 重新评估全局状态。

    2.7 时间线汇总

    轮次 i 序列 发生的事
    冷启动 2→1→0 自下而上找到 TableScan,读出 batch1,喂 FilterProject
    推进 1→2 FilterProject 产出 batch1',喂 sink,sink 写 buffer
    稳态 1→0→1→2… 反复拉取-过滤-写出
    阻塞 A 1 sink buffer 满 → kBlock 下线,等消费
    阻塞 B 0 TableScan 无 split → kBlock 下线,等 addSplit
    收尾 0(break)→2→1(break)→2 noMoreInput 逐级传播 → isFinished → close → kAtEnd

    一句话概括这套调度:指针从 sink 往 source 扫描以"定位能产出数据的最上游算子",一旦取到数据就靠 i+=2 把它连续推向 sink;每个算子边界都轮询 stop/block,需要等待时整条 Driver 作为无栈协程挂起,被唤醒后从 sink 重新扫描。


    3. Task 的设计与实现

    3.1 执行架构:谁管理谁,谁推进数据

    对象 职责 与其他对象的关系
    QueryCtx 查询配置、查询内存池、executor / spill executor 等上下文 Task 持有 shared_ptr<QueryCtx>;可为多个 Task 提供共同上下文
    Task 执行一个 PlanFragment,管理 Driver、split、共享状态和结束流程 拥有 DriverFactory 和 Driver,协调跨 pipeline 的连接
    DriverFactory 描述一条 pipeline 的计划节点、算子创建方式和并行度 创建多份 Driver;pipeline 是逻辑执行链,不是一个独立线程对象
    Driver 调用算子接口、传递批次、响应阻塞和控制请求 独占 DriverCtx 与 vector<unique_ptr<Operator>>
    Operator 完成扫描、过滤、聚合、Join、输出等工作 通过上下文访问 Task 与内存池,向 Driver 暴露输入、输出和等待状态
    Executor 排队并执行 Driver::run 一个 Driver 可多次入队,每次执行不必落在同一个工作线程上

    同一 pipeline 的多个 Driver 拥有各自的 Operator 实例;需要共享的数据通过 Task 管理的 JoinBridge、LocalExchange 等结构连接。Driver 的执行受线程注册协议约束,同一时刻只有获准进入它的线程能推进其算子状态。

    这里还要区分三个数量:计划创建的 Driver 数、尚未关闭的 Driver 数、当前活跃 Driver 线程数。一个等待 split 的 Driver 仍然存在,但已经不占用 executor 的执行线程。

    源码:Task 接口、Driver 与上下文、DriverFactory、ThreadState。

    3.2 Task 是什么

    Task 代表一个查询片段(PlanFragment)的完整执行单元,是 Velox 执行层最顶层的对象,负责:

    • 将 PlanFragment 翻译成若干 Pipeline(DriverFactory);
    • 创建并管理所有 Driver;
    • 管理 Split(数据分片)的分发;
    • 维护跨 Driver 共享的桥接结构(JoinBridge、LocalExchange 等);
    • 协调终止、暂停、取消等控制流。

    定义在 velox/exec/Task.h:43。

    3.3 Task 状态机

    TaskState(velox/exec/TaskStructs.h:44)
    
    kRunning(0) ──── 正常结束 ──────────────────► kFinished(1)
        │          ──── requestCancel() ──────────► kCanceled(2)
        │          ──── requestAbort()  ──────────► kAborted(3)
        │          ──── setError()      ──────────► kFailed(4)
        └─── 所有路径均通过 Task::terminate(terminalState) 实现
    

    状态转换核心代码(Task.cpp:2516):

    流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

    state_ = terminalState;
    terminateRequested_ = true;
    

    3.4 核心数据成员

    成员 类型 作用
    drivers_ vector<shared_ptr<Driver>> 拥有所有 Driver(Task.h:1250)
    driverFactories_ vector<unique_ptr<DriverFactory>> 每条 Pipeline 一个,描述如何创建 Driver
    splitGroupStates_ unordered_map<uint32_t, SplitGroupState> 每个 split group 的跨算子共享状态(JoinBridge、LocalExchange、Barrier)
    splitsStates_ unordered_map<PlanNodeId, SplitsState> 每个 source node 的 Split 队列
    numThreads_ int32_t 当前在线程上运行的 Driver 数量
    terminateRequested_ atomic_bool 是否请求终止,Driver 会频繁轮询
    pauseRequested_ atomic_bool 是否请求暂停
    mutex_ timed_mutex 保护相应 Task 复合状态;ThreadState、统计和独立子组件仍有各自的 atomic、锁与访问协议。

    3.5 Task 的执行模式

    Task.h:46 定义了两种模式:

    • kParallel:通过 Task::start() 启动,把 Driver 提交到 Executor 并行执行(生产用途)。
    • kSerial:通过 Task::next() 在调用线程上单步执行,用于测试或嵌入式查询。

    3.6 Task 启动流程

    3.7 从 PlanFragment 到可调度的 Driver

    3.7.1 Pipeline 在哪里切开

    LocalPlanner::detail::plan 递归遍历计划,把同一执行链上的 PlanNode 放入同一个 DriverFactory。当前 mustStartNewPipeline 的主要规则是:

    • LocalMerge、MixedUnion、LocalPartition 的各个 source 单独形成 pipeline。
    • 一般多输入节点的非第一个 source 单独形成 pipeline。例如 Hash Join 的 build 侧形成以 HashBuild 结束的 pipeline,probe 侧包含 HashProbe。
    • IndexLookupJoin 有专门处理,其索引 source 由算子内部管理,不按普通二输入 Join 直接展开。

    PlanNode 和 Operator 也不是固定的一一对应:相邻 Filter 与 Project 可以合成一个 FilterProject。因此读 Driver 中的 operator ID 时,应对照实际创建出的算子序列。

    每条 pipeline 的实际 Driver 数为 min(factory->maxDrivers, start 参数 maxDrivers)。factory->maxDrivers 还受单线程算子、LocalExchange 分区数、writer 配置及扩展算子约束。maxDrivers 因而是每条 pipeline 的请求上限,不能直接当成整个 Task 的线程数。

    源码:切分规则、并行度计算、创建算子及融合。

    3.7.2 并行模式的启动顺序

    下面是流程示意,省略锁、异常和统计:

    Task::create → initTaskPool
    Task::start(maxDrivers, concurrentSplitGroups)
      createDriverFactoriesLocked → LocalPlanner::plan
      initializePartitionOutput
      createAndStartDrivers
        createSplitGroupStateLocked
        createDriversLocked → DriverFactory::createDriver
        Driver::enqueue → queryCtx()->executor()->add(...)

    Driver::enqueue 先记录排队状态和起始时间,再向 executor 提交捕获 shared_ptr<Driver> 的 lambda。工作线程随后调用 Driver::run,由 runInternal 向 Task 注册进入,并首次调用各算子的 initialize()。算子对象创建与执行期初始化分处两个阶段;当前 Task 代码还检查创建 Driver 时 task pool 没有产生意外的内存预留。

    源码:Task::start、创建与入队、Driver::enqueue、initializeOperators。

    3.7.3 Serial 模式仍可能有多条 pipeline

    kSerial 使用调用线程上的 Task::next()。Task 初始化时以 maxDrivers=1 规划,并为每条 pipeline 创建一个 Driver;next() 依次推进它们,一个 Driver 阻塞时可以继续运行其他 Driver。它并不等价于整张计划只有一个 Driver。

    当前实现要求 ungrouped execution,并检查各 pipeline 支持串行执行。普通 split 模式下,调用 next() 前要发出 noMoreSplits;请求了 barrier 时有另一套分批处理约束。Serial 不使用 consumer callback 交付结果,而是让末端算子的一批输出沿 Driver::next → Task::next 返回给调用方。

    当所有剩余 Driver 都被外部事件阻塞时,Task::next(&future) 可返回空并通过 collectAny 给出等待句柄;等待结束后由调用方再次调用 next()。因此“返回空”需要结合 future 和 barrier 状态解释。Driver 内部此时还有一个特殊约定:返回一批 serial 结果也使用 StopReason::kBlock,但没有对应的 BlockingState。

    源码:Serial 初始化、Task::next、Driver::next。

    调用链:Task::start() → createDriverFactoriesLocked() → initializePartitionOutput() → createAndStartDrivers()

    流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

    // Task.cpp:982 — Task::start()
    void Task::start(uint32_t maxDrivers, uint32_t concurrentSplitGroups) {
      // 1. 通过 LocalPlanner::plan() 将 PlanFragment 切分成 DriverFactory 列表
      createDriverFactoriesLocked(maxDrivers);
      // 2. 初始化 PartitionedOutput buffer(如果有 PartitionedOutputNode)
      initializePartitionOutput();
      // 3. 为 ungrouped execution 创建 Driver,为 grouped execution 按需创建
      createAndStartDrivers(concurrentSplitGroups);
    }
    

    createAndStartDrivers()(Task.cpp:1070):

    // 为 ungrouped execution 创建并入队所有 Driver
    for (auto it = drivers_.end() - numDriversUngrouped_; it != drivers_.end(); ++it) {
      ++numRunningDrivers_;
      Driver::enqueue(*it);   // 提交到线程池
    }
    // 如果有 grouped execution,调用 ensureSplitGroupsAreBeingProcessedLocked()
    // 按 concurrentSplitGroups_ 上限,为已有 split 的 group 创建 Driver
    

    drivers_ 的内存布局(Grouped Execution):

    drivers_[0 .. numDriversPerSplitGroup_*concurrentSplitGroups_-1]   ← grouped execution 槽位
    drivers_[numDriversPerSplitGroup_*concurrentSplitGroups_ ..]       ← ungrouped execution Driver
    

    Grouped execution 的槽位在 split group 完成后置 nullptr,新 split group 的 Driver 复用这些槽位。

    3.8 Split 分发

    // 普通 queue-backed SplitsStore 路径。
    Task::addSplit -> addSplitLocked -> SplitsStore::addSplit
      单个普通 split:放入队列,取出一个等待 promise,在锁外兑现
      barrier split:有面向同 source 多个 Driver 的广播处理
      noMoreSplits:通知仍等待的消费者重新检查结束条件
    

    Driver 通过 Task::getSplitOrFuture()(Task.h:442)从 SplitsStore 取 split。如果没有 split,返回 kWaitForSplit,Driver 进入 Blocked 状态。

    3.9 allPeersFinished —— Hash Join Build 同步屏障

    3.10 多个 Driver 如何协调

    算子侧的 peer barrier 入口为 OperatorCtx::allPeersFinished。它调用 Task::allPeersFinished 汇集 Driver,再用 shared_ptr 的 aliasing constructor 返回相同 operatorId 的 peer Operator:指针指向 Operator,控制块仍持有 Driver。这既保留 Task 层的 Driver 屏障,也让 HashBuild / HashProbe 不必自行遍历 Driver 查找算子。

    3.10.1 Split 分发与唤醒粒度

    Task 按 source PlanNode 和 split group 管理 SplitsStore。getSplitOrFuture 尝试取 split;尚无 split 且来源未结束时,返回 kWaitForSplit 和等待 future。普通 addSplit 将数据放入队列,取出一个等待 promise,出锁后兑现;普通的单个 split 不会广播唤醒所有等待 Driver。Barrier split 面向同一 source 的每个 Driver,才有相应的广播处理;noMoreSplits 也需要通知仍在等待的消费者。

    当前 SplitsStore 是抽象接口,有队列与其他 split 来源实现。队列实现还可配合 preload 选择已经准备好的 split,所以不能无条件把所有 scan split 调度理解为严格 FIFO。

    源码:Task::addSplit、SplitsStore 的唤醒策略、getSplitOrFuture。

    3.10.2 Grouped execution 限制同时活跃的 split group

    Grouped pipeline 为每个 split group 创建独立的一组 Driver。Task 最多同时推进 concurrentSplitGroups_ 个 group;Driver 槽位数量与生命周期总 Driver 数不同。前部槽位为同时运行的 grouped Driver 预留,后部可以放 ungrouped Driver。

    某 group 的所有 Driver 关闭后,Task 清理对应共享状态、记录完成,并为 queued group 创建新 Driver、复用空槽位。混合 grouped / ungrouped Join 还会延长连接两侧的 bridge 生命周期,避免一侧先结束就清掉另一侧仍需使用的状态。

    源码:启动和槽位安排、split group 调度、SplitGroupState。

    3.10.3 Join Build 的 peer barrier 与 JoinBridge

    allPeersFinished 按 splitGroupId + planNodeId 汇集同一 pipeline 的 peer Driver。常规 HashBuild 路径中,非最后的 Driver 得到等待 future;最后到达者拿走 peers 和 promises,并负责合并 build 结果、发布 JoinBridge,随后通知其他参与者。

    因此要分开理解两件事:peer barrier 协调 build Driver 之间的汇合;JoinBridge 向 probe 侧提供 build 结果。最后一个 peer 到达只说明参与者已汇齐,调用方还需要完成实际的建表与发布。

    allPeersFinished 也允许 future == nullptr 的非等待调用方式,并非只供 Hash Join 使用。Task 的 requestBarrier / Driver drain 则是另一套执行阶段屏障,不能与这个 peer barrier 混用。

    源码:allPeersFinished、HashBuild 建表与发布、执行 barrier 接口。

    Task.h:526:所有并发的 build driver 都执行完后,最后一个 driver 收到 true,其他 driver 通过 ContinueFuture 阻塞等待。这是 HashJoin 多线程 build 的协调机制。

    SplitGroupState::barriers(TaskStructs.h:213)存储这些屏障状态,持有 Driver 的 shared_ptr,形成一个引用环—— 该引用环在 allPeersFinished 正常路径或 Task::terminate 异常路径中被清理。

    3.11 Task::terminate —— 统一终止路径

    Task.cpp:2502 是所有终止的汇聚点:

    terminate(terminalState) {
      1. 加锁,将 state_ 设置为终态
      2. terminateRequested_ = true(Driver 循环会看到)
      3. numRunningDrivers_ = 0
      4. 对所有 drivers_:
           enterForTerminateLocked() → 不在线程上的 Driver 直接标 isTerminated,
           收集到 offThreadDrivers(需在锁外 closeByTask())
      5. 解锁后:
           offThreadDrivers.closeByTask()        // 关闭不在线程的 Driver
           maybeRemoveFromOutputBufferManager()  // 清理输出 buffer
           关闭 ExchangeClient
           清理 splitGroupStates(bridges、barriers、localExchanges)
           兑现所有 splitPromises(让等待 split 的 Driver 感知终止)
           cancel 所有 JoinBridge
      6. 返回 future,调用方 await 等待所有线程退出
    }
    

    3.12 普通 pipeline 与 FixedPointLoop 的执行边界

    包含 FixedPointNode 并由 FixedPointLoop 接管的计划是一个明确例外。Task::create 建立 loop,Task::init 跳过普通 Driver 编译;start / next 委托给 loop,由它调度子 Task,addSplit / noMoreSplits 也交给 loop 路由。因而“每个 Task 都直接把整张计划编译为一组 Driver”只适用于普通 pipeline 路径。子 Task 仍可以使用本文的 Driver 机制;loop 完成或失败再汇入拥有它的 Task 终态。

    源码:Task::create 的路径选择、串行入口、并行入口。

    4. runInternal 的执行模型:不是 Volcano,也不是教科书式 push

    4.1 数据如何在一条 pipeline 中前进

    4.1.1 算子接口表达不同的事情

    接口 含义
    needsInput() 当前是否能接收下一批输入
    addInput(batch) 接收输入,并按算子实现进行计算、缓存或更新状态
    getOutput() 返回一批非空向量,或表示暂时没有可交付输出的 nullptr
    isBlocked(&future) 当前能否继续推进;等待时给出原因和有效 future
    noMoreInput() 上游不会再调用 addInput,算子可以完成剩余输出
    isFinished() 已完全结束,不再产生输出;也可能提前结束

    计算可以发生在 addInput 中。例如 HashAggregation::addInput 调用 GroupingSet::addInput,直接更新聚合状态。getOutput 的统一职责是向 Driver 交付输出批次,而不是包办所有计算。一次输入可以对应零批、一批或多批输出;部分聚合也可能在输入结束前 flush。

    源码:Operator 契约、聚合输入处理。

    4.1.2 从下游检查需求,拿到一批后向下游推进

    operators_[0] 是 source,最后一个是 sink。普通 runInternal 从 sink 开始,向 source 扫描。对于相邻的 op 与 nextOp,它先检查阻塞,再检查下游需求,最后决定是否从上游取数据。

    以 Scan → FilterProject → Sink 为例展示一批数据的传递。箭头表示方法调用或返回,调用者始终是 Driver。
    图 2:以 Scan → FilterProject → Sink 为例展示一批数据的传递。箭头表示方法调用或返回,调用者始终是 Driver。 打开原图

    下面是相邻算子的核心分支示意,省略计时、trace 和异常包装:

    // 流程示意:处于 for (...; i >= 0; --i) 中
    if (op->isBlocked(&future) != kNotBlocked) {
      return blockDriver(...);
    }
    if (nextOp->isBlocked(&future) != kNotBlocked) {
      return blockDriver(...);
    }
    if (!nextOp->needsInput()) {
      break;  // 停止本轮向上游扫描,外层循环重新从下游开始
    }
    if (auto batch = op->getOutput()) {
      nextOp->addInput(batch);
      i += 2; // 随后的 --i 使净效果为 i + 1,优先推进下游
      continue;
    }
    if (op->isBlocked(&future) != kNotBlocked) {
      return blockDriver(...);
    }
    if (op->isFinished()) {
      nextOp->noMoreInput();
      break;
    }
    // 当前算子需要更多输入,继续向上游扫描
    

    getOutput() 之后再次检查 isBlocked() 很重要:扫描器可能在取数据时才发现没有 split 或 connector 尚未就绪,并在本次调用中生成 future。返回 nullptr 后,Driver 据此选择阻塞、传播结束,或继续寻找上游输入。单独的 nullptr 不会自动让出线程。

    当前下游 needsInput()==false 分支会直接 break,让外层循环重新从下游开始。这个可传递的背压限制了更上游的提前读取;相关修复也防止尚未加载的 LazyVector 留在中间算子中,而 source reader 已经移到下一批数据。

    普通恢复仍从 sink 开始检查;curOperatorId_ 等字段主要记录当前位置和统计信息。Trace 输入写入受阻的特殊路径会保存中间批次 traceInput_,并从对应算子恢复,不应把它推广为所有 block 的恢复策略。

    源码:runInternal、背压分支、恢复起点。

    4.1.3 Sink 与结束传播

    在并行执行中,sink 的 getOutput() 用来推进末端算子,但结果通过 callback 或输出 transport 交付;Driver::run 明确检查它没有直接返回结果批次。Serial 模式才把末端批次交还调用方。

    上游 isFinished() 后,Driver 调用下游 noMoreInput(),再让它排出剩余结果。最终 sink 完成时,Driver 调用 close() 并返回 kAtEnd。Limit 等算子可提前完成,后续由 Task 的完成判定决定是否终止其他 pipeline。

    这种执行方式可以概括为:Driver 以批次为单位,按下游需求调用上游 getOutput,再调用下游 addInput。算子之间的常规传递由这一个调度循环完成。

    4.1.4 第二次 isBlocked 把“没有输出”变成可调度的等待

    getOutput() 之前的 isBlocked() 只能说明此刻没有需要交给 Driver 的已知等待,不能保证接下来的 getOutput 一定能拿到数据。取队列、请求 split、尝试 connector 读取时,算子才可能发现需要等待,并在本次调用中保存一个新 future。getOutput 返回 nullptr 后,Driver 必须再问一次 isBlocked,才能把刚发现的等待转交给框架。

    当前 LocalExchange 是一个直接的例子:isBlocked 先交出先前保存的 future;如果尚未记录阻塞,它返回 kNotBlocked。getOutput 才调用 queue->next,更新 blockingReason_ 与 future_。因此可能出现“第一次检查不阻塞 → getOutput 发现队列空 → 第二次检查报告 kWaitForProducer”的完整过程。不能将第一次检查当成对整次算子调用的可用性保证。

    源码节选:LocalExchange::isBlocked

    BlockingReason LocalExchange::isBlocked(ContinueFuture* future) {
      if (blockingReason_ != BlockingReason::kNotBlocked) {
        *future = std::move(future_);
        auto reason = blockingReason_;
        blockingReason_ = BlockingReason::kNotBlocked;
        return reason;
      }
    
      return BlockingReason::kNotBlocked;
    }

    源码节选:LocalExchange::getOutput

    RowVectorPtr LocalExchange::getOutput() {
      if (hasDrained()) {
        return nullptr;
      }
    
      RowVectorPtr data;
      bool drained{false};
      blockingReason_ = queue_->next(&future_, pool(), &data, drained);
      if (blockingReason_ != BlockingReason::kNotBlocked) {
        VELOX_CHECK(future_.valid());
        VELOX_CHECK(!drained);
        return nullptr;
      }
    
      if (data != nullptr) {
        VELOX_CHECK(!drained);
        auto lockedStats = stats_.wlock();
        lockedStats->addInputVector(data->estimateFlatSize(), data->size());
        return data;
      }
    
      if (drained) {
        VELOX_CHECK(!isDraining());
        operatorCtx_->driver()->drainOutput();
      } else {
        VELOX_CHECK(queue_->isFinished());
      }
      return nullptr;
    }
    getOutput 返回后的状态Driver 的动作线程去向
    isBlocked 返回等待原因,并交出有效 futureblockDriver 保存等待,返回 kBlock;退出 runInternal这份 Driver off-thread,worker 可执行其他工作
    未阻塞,isFinished 为 true调用 nextOp->noMoreInput,让下游排出剩余结果继续按框架结束传播,不因 nullptr 自动挂起
    未阻塞,也未结束当前算子可能还要上游输入,继续向 source 查找继续在当前调用中推进
    Task 已要求 pause / terminate 等先服从 shouldStop 的返回结果沿对应退出路径处理,而不是忽略 Task 状态

    严格说,isBlocked(&future) 本身没有把操作系统线程阻塞住,也没有把 Driver 从 executor 上摘走;它返回 BlockingReason,并在需要时交出有效的 ContinueFuture。Driver 检查的是 reason,而非通过轮询 future.isReady() 决定要不要停止。真正的 off-thread 发生在 blockDriver 返回 kBlock、runInternal 的栈退出、CancelGuard 调用 Task::leave 之后。

    这段“空输出后的二次检查”是相邻上游算子的取数分支。并行 pipeline 最末端 sink 的 getOutput 是另一条分支,用于推进 sink,不能直接返回批次给 Driver::run;sink 在数据操作中产生的背压,会在框架后续的 isBlocked 检查中被接走。Serial 模式还有把结果批次返回调用方的 kBlock,不能与这里持有 BlockingState 的异步等待混为一谈。

    完整源码:runInternal 的数据分支。后面的 7.1.3 将这一次返回接到 future 唤醒和重新入队。

    准确说法:这是一个由 Driver 集中调度的、单栈帧的、demand-gated(需求门控的)push 执行模型。 第 2.5 / 5.2.5 节用了"主循环"的笼统措辞,本章给出精确语义。

    4.2 控制流:Driver 是唯一的调度器

    operators_ 是一个扁平数组,索引含义固定:

    operators_[0]    = source(TableScan / Exchange)   ← 最上游
    ...
    operators_[N-1]  = sink(PartitionedOutput / CallbackSink) ← 最下游
    数据流向: 0 ──► 1 ──► 2 ──► ... ──► N-1
    

    在这条 Driver 数据推进路径中,runInternal 对 operators_ 调用 isBlocked、needsInput、getOutput、addInput、isFinished 和 noMoreInput,避免以一层 next 递归嵌套一层的方式传递批次。算子内部仍可以调用其他对象并执行计算,reclaim/abort 还有独立调用路径;“一个栈帧”只概括中央推进循环,不是程序真实调用栈始终只有一层。

    与 Volcano 的本质区别:

    维度 Volcano(经典 pull) Velox Driver
    谁驱动执行 根算子 next() 递归调用孩子的 next() 单个 Driver 循环调用所有算子
    调用栈 深递归:sink→…→source,每批一层套一层 扁平:一个栈帧 + 对算子数组的循环
    算子间耦合 算子 A 直接调用算子 B 算子互不调用,全由 Driver 中转
    数据传递 孩子把行返回给父亲(返回值上溯) 生产者 getOutput() → Driver → 消费者 addInput()

    4.3 调度扫描方向:从 sink 往 source 找"谁能干活"

    普通情况下 getStartingOperator 从末端算子开始扫描;它还要考虑特定执行状态,不能把单个示例中的起点作为所有重入场景的固定值。下面三算子 trace 展示没有这些附加状态的常规路径,完整分支以当前 getStartingOperator/runInternal 为准。

    源码核对:当前 getStartingOperator() 只有 traceInput_ 非空时返回 blockedOperatorId_;普通执行返回 operators_.size() - 1,从 consumer 一端推进。不要把输入追踪的特殊入口解释成所有 future 恢复都从 blockedOperatorId_ 开始。

    流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

    int32_t startingOperator = getStartingOperator();          // 普通路径为 N-1;traceInput_ 非空时返回 blockedOperatorId_
    for (;;) {
      for (int32_t i = startingOperator; i >= 0; --i) {        // sink → source 递减
        auto* op = operators_[i].get();
        if (i < operators_.size() - 1) {
          Operator* nextOp = operators_[i + 1].get();          // i+1 是 i 的下游
          if (nextOp->needsInput()) {                          // ← 下游发出"需求信号"
            getOutput(op, intermediateResult);                 // 从上游 op 取一批
            if (intermediateResult) {
              addInput(nextOp, intermediateResult);            // 推给下游
              i += 2; continue;                                // 回头看下游能否再往前推
            }
          }
        } else {
          // i == sink:串行模式产出结果 / 判断 isFinished
        }
      }
    }
    

    要分清两个方向:

    • 调度扫描方向:下游 → 上游(找活干);
    • 数据流动方向:上游 → 下游(push)。

    i += 2 配合 for 的 --i,净效果是 i+1——回头去看刚收到数据的那个下游算子,能否把它的产出再往前推一步。 于是数据一步步被 push 向 sink。数据是被 push 的,不是被 pull 上来的。

    4.4 demand-gated(需求门控)

    下游 needsInput 为 false 时,Driver 不应继续推进对应上游,否则可能过早读取 source 或 LazyVector 状态。是否需要输入由具体算子当前缓冲与阶段决定;HashBuild 在正常收集输入时通常仍需要数据,“正在攒数据”本身不是拒绝输入的通用理由。

    跨 Pipeline 的反压则通过 BlockingReason::kWaitForConsumer + ContinueFuture,让整个 Driver 下线等待。

    4.5 为什么"扁平循环 + Driver 调度"是协作式调度的前提

    扁平 Driver 循环使正常阻塞点能够把继续执行所需状态存回对象,然后返回并重新进入,便于采用无栈续调。递归执行器也可以设计显式状态、future 或有栈协程,因此“Volcano 无法挂起”并非普遍结论;这里应关注 Velox 当前具体保存了哪些状态、在哪些边界返回。

    扁平循环 + Driver 中央调度(无算子间递归)
            │
            ▼
    Driver 现场可压缩成几个标量(curOperatorId_ / blockedOperatorId_ /
                                blockingReason_ / traceInput_),可在算子边界随时 return
            │
            ▼
    StopReason 协议(kBlock/kYield/kPause/kTerminate)得以成立
            │
            ▼
    协作式调度:M 个 Driver 复用 N 个线程,阻塞即让出
    

    如果递归调用中需要在深层暂停,实现者必须决定如何保存上层现场,或让阻塞信息逐层返回。Velox 的中央循环把多数续调状态集中在 Driver/Operator 成员中,减少需要保存的隐式调用栈;这是一种实现取舍,不构成对所有递归 pull 引擎能力的限制。

    4.6 小结

    • Driver 是调度器,模型是 push,不是 Volcano 的顶层 next() 递归;
    • push 的同时有下游 needsInput() 的需求门控(单 Pipeline 内反压),是 demand-gated push 而非无脑灌;
    • 中央循环与对象化状态便于在明确边界退出和续调,同时统一检查暂停、取消和下游需求;这是当前结构的重要收益,但不能从源码直接断言它是唯一设计目的。

    5. Driver 的设计与实现

    5.1 Driver 是什么

    Driver 是单条 Pipeline 的单线程执行引擎。每个 Driver 持有一组串联的 Operator,从 source operator (TableScan、Exchange 等)拉数据,逐级传递给 sink operator(PartitionedOutput、CallbackSink 等)。

    定义在 velox/exec/Driver.h:353。

    5.2 Driver 的核心组成

    Driver
      ├── ctx_: unique_ptr<DriverCtx>               // 持有 Task shared_ptr(引用环的一侧)
      ├── operators_: vector<unique_ptr<Operator>>  // 独占所有算子
      ├── state_: ThreadState                       // 线程状态,由 Task::mutex_ 保护
      ├── barrier_: BarrierState                    // Barrier 处理状态
      └── blockingReason_: BlockingReason           // 当前阻塞原因
    

    DriverCtx(Driver.h:231):Driver 的执行上下文,是 Driver 与 Task 之间的接口层。它持有 shared_ptr<Task>,这是 Driver → Task 引用的唯一路径。

    5.3 ThreadState —— Driver 的状态

    5.4 Driver 的调度状态与 Task 的暂停协议

    5.4.1 ThreadState 是一组字段,不是单一状态枚举

    Driver 的常见调度状态示意。Suspended 保留线程;Closed 表示算子已关闭;isTerminated 是可能先于关闭设置的终止标志。
    图 3:Driver 的常见调度状态示意。Suspended 保留线程;Closed 表示算子已关闭;isTerminated 是可能先于关闭设置的终止标志。 打开原图
    状态或阶段 关键字段 谁让它继续前进
    Created / off thread 没有执行线程 启动或恢复路径
    Enqueued isEnqueued=true Executor 取出任务,随后 Task::enter
    执行算子 thread 已设置,numSuspensions==0 当前工作线程
    Blocked off thread,hasBlockingFuture=true Future continuation
    因 pause 停下 off thread,pause 请求仍有效 Task::resume;若仍有 blocking future,则继续等该事件
    Suspended thread 仍设置,numSuspensions>0 原线程退出 suspended section
    Closed closed_=true 不再执行算子;等待引用释放

    这些字段表达多个维度。比如排队中的 Driver 可能同时遇到 Task pause;终止标志也可能已经设置,而算子清理尚未完成。正常 Driver::close 与 ThreadState::isTerminated 的赋值位置不同,不能把“正常结束”机械翻译成某个统一的 Driver enum 迁移。

    5.4.2 enter / leave 保护的是执行资格与线程计数

    Task::enter 在锁内消费 isEnqueued,检查重复进入、终止、暂停和 yield 请求。返回 kNone 时才增加 numThreads_,设置 thread,并允许 Driver 执行算子。

    Task::leave 撤销这次进入。若发现终止请求,会先在锁外调用关闭回调,再减少线程计数;计数归零时取走 threadFinishPromises_,在锁外兑现。CancelGuard 即使调用过 notThrown(),析构时仍然会 leave;notThrown() 只取消异常路径的主动关闭,不会跳过线程登记的撤销。

    这一协议只在必要的共享状态修改处使用 Task 锁;正常的整段 Operator 执行并不一直持有该锁。

    源码:ThreadState、Task::enter / leave、CancelGuard 析构。

    5.4.3 Yield、Pause、Suspend 分别解决什么问题

    机制 是否退出本轮 runInternal 是否保留当前算子调用栈 恢复方式
    StopReason::kYield 是 否 立即重新 enqueue,重新竞争 executor 线程
    StopReason::kPause 是 否 等 Task 恢复后,根据当前字段决定是否 enqueue
    Block 是 否 Future 就绪后调度
    Suspended section 否 是 原线程执行 leaveSuspended 后继续

    Driver 时间片配置 driver_cpu_time_slice_limit_ms 默认为 0;非零时在循环检查点比较 execTimeMs() 与配置值。当前 execTimeMs() 由进入线程以来的时间差计算,并非精确的 CPU 消耗计数。Serial 模式关闭这一 Driver 时间片机制。Task 还可通过 requestYield / yieldIfDue 请求活跃 Driver 主动让出。

    这些检查发生在协作点。单个长时间运行的算子调用仍需返回,或在内部主动检查控制请求,调度器才能及时响应;时间片并不构成强制抢占。

    requestPause() 设置 pauseRequested_,返回等待 numThreads_==0 的 future。之后 reclaimer 才能在该协议下检查和回收算子内存。Task::MemoryReclaimer::reclaimTask 实际使用了 pause / resume,并通过 guard 保证恢复路径被执行。

    5.4.4 Suspended 保留线程,但暂时允许 Task 停稳

    enterSuspended 首次进入时减少 numThreads_,保留 thread 和 C++ 栈;递归进入只增加 numSuspensions。最后一层 leaveSuspended 恢复线程计数。如果 Task 还处于 pause,原线程在锁外短暂休眠并重查,直到允许继续。

    因此 suspended section 中的线程不能继续访问可能被回收的 Driver 内存。它仍占有原来的线程与栈,区别于通过 future 把线程交回 executor 的 Block。

    暂停与终止同时发生时,还要避免关闭算子与 reclaimer 并发访问。当前 shouldStopLocked 先检查 pause;enterForTerminateLocked 对已暂停、off-thread 的 Driver 延后关闭,交给 Task::resume 的终止分支处理。这个优先次序服务于资源访问协议。

    源码:shouldYield、配置默认值、暂停与 suspended、Task::resume、内存回收调用方。

    Driver.h:97 定义了 ThreadState:

    流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

    struct ThreadState {
      atomic<thread::id> thread;            // 正在运行的线程 id(非空 = 在线程上)
      atomic<bool> isEnqueued;              // 已提交到 Executor 但还未开始运行
      atomic<bool> isTerminated;            // 已终止,终态
      tsan_atomic<bool> hasBlockingFuture;  // 有待兑现的 future(Blocked 状态)
      atomic<uint32_t> numSuspensions;      // Suspended 层数(> 0 = suspended)
    };
    

    5.5 Driver 状态流转

    Driver 调度与保栈 Suspended
    图 4:Driver 调度与保栈 Suspended。已按当前实现修正标注,具体约束见相邻正文。
    Fig. Driver 状态流转:Created → Enqueued → On Thread,再分流到 Blocked / Off Thread / Terminated

    状态说明(对应代码注释 Driver.h:72):

    状态 判断条件 进入方式 退出方式
    Created 所有 flag 为 false Task::createDriversLocked enqueue
    Enqueued isEnqueued=true Driver::enqueue() Executor 调度到线程
    On Thread thread 字段非空 Task::enter() 返回 kNone 返回 kBlock/kYield/kPause/kTerminate
    Blocked hasBlockingFuture=true BlockingState 构造时设置 future 兑现 → BlockingState::setResume() → enqueue
    Suspended numSuspensions > 0 Task::enterSuspended() Task::leaveSuspended()
    Terminated isTerminated=true Task::enter() 返回 kTerminate,或 Task::leave() 无(终态)

    5.6 Driver::runInternal —— 核心执行循环

    Driver.cpp:519 是 Driver 的心脏(详细的执行模型分析见第 4 章):

    // runInternal 的主干流程伪代码;保留关键分支,不作完整函数使用。
    Task::enter(state) 获取本次执行资格;根据 StopReason 处理暂停 / 终止等
    安装 CancelGuard:析构总会 Task::leave;异常时还触发关闭
    首次 initializeOperators
    循环从 getStartingOperator() 选定的算子开始:
      检查 Task stop / yield、QueryCtx 仲裁与算子 isBlocked
      若下游 needsInput,调用上游 getOutput
        有输出:调用下游 addInput(这里也可计算)并向下游推进
        无输出:再次检查上游 isBlocked,必要时保存 future 并退出
        未阻塞且 isFinished:通知下游 noMoreInput
        未完成:继续依据 needsInput 等状态寻找可推进的上游
      sink 完成:close,返回 kAtEnd
      serial 产出一批结果:返回 kBlock,但不对应普通 BlockingState
    catch:向 Task 记录异常;按终止协议退出
    // 普通 async kBlock:离开 Driver 后才安装 BlockingState 恢复回调。
    

    5.7 blockDriver —— 进入 Blocked 状态

    流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

    // Driver.cpp:669(声明 Driver.h:669)
    StopReason blockDriver(...) {
      blockingState = make_shared<BlockingState>(self, move(future), op, reason);
      guard.notThrown();          // 避免 CancelGuard 触发 close
      return StopReason::kBlock;
    }
    // Driver::run() 中:
    case StopReason::kBlock:
      BlockingState::setResume(blockingState);  // 注册 future 回调
      return;
    

    BlockingState::setResume()(Driver.cpp:227):future 兑现时(通过 QueuedImmediateExecutor)记录算子 blocking 时间、hasBlockingFuture = false,若 task 未 pause 则 Driver::enqueue(driver) 重新入队。

    5.8 CancelGuard —— 线程安全的清理保证

    Driver.h:563,RAII 对象。runInternal 正常结束时调用 notThrown();若发生异常,guard 析构会调用 Task::leave() 和 Driver::close(),确保 Driver 干净退出。

    5.9 Driver::close —— 正常关闭路径

    当前源码摘录(1d1b76567870):velox/exec/Driver.cpp:1094。省略外围声明;此片段未作为独立程序编译。

    void Driver::close() {
      if (closed_) {
        // Already closed.
        return;
      }
      if (!isOnThread() && !isTerminated()) {
        LOG(FATAL) << "Driver::close is only allowed from the Driver's thread";
      }
      closeOperators();
      updateStats();
      closed_ = true;
      Task::removeDriver(ctx_->task, this);
    }
    

    closeOperators()(Driver.cpp:898):关闭所有算子、将 driver 生命周期统计(queued/on-thread/blocked 时间)与算子统计上报到 Task。


    6. isBlocked 与三种让出机制:结合典型算子解读

    6.1 阻塞与唤醒:future 怎样接回执行循环

    6.1.1 先退出 Driver,再安装恢复回调

    并行 Driver 等待 split 的一次往返。即使 future 提前就绪,也要先离开 Driver,再安装恢复回调;唤醒后重新入队。
    图 5:并行 Driver 等待 split 的一次往返。即使 future 提前就绪,也要先离开 Driver,再安装恢复回调;唤醒后重新入队。 打开原图

    并行路径按以下顺序执行:

    1. 算子报告 BlockingReason 和 future。blockDriver 检查 future 有效,创建持有 Driver 的 BlockingState。
    2. 构造 BlockingState 时,将 hasBlockingFuture=true,然后返回 kBlock。
    3. 离开 runInternal 时,CancelGuard 调用 Task::leave,减少活跃线程计数并清除 Driver 的 thread 字段。
    4. 回到 Driver::run,再调用 BlockingState::setResume。它检查 Driver 已经 off thread,并通过 QueuedImmediateExecutor 安装 continuation。
    5. Future 成功就绪后,回调在 Task 锁内清除 hasBlockingFuture。未暂停时调用 Driver::enqueue;暂停时暂不入队,交给后续 Task::resume 处理。

    回调可能在 future 已就绪时很快执行,因此第 3、4 步的顺序保证了重新进入 Driver 前,原来的算子执行区间已经结束。QueuedImmediateExecutor 用于 continuation 调度;真正执行 Driver 的 executor 来自 QueryCtx。

    Driver::run 还处理 block 与 Task 终止并发的情况。Future 以异常结束时,并行恢复路径把它转成 Task error。Serial 的 DriverBlockingState 则记录异常、通知等待者,并在后续检查中传播,不能把两条路径完全等同。

    源码:BlockingState 与回调、blockDriver、Driver::run、Serial 等待管理。

    6.1.2 等待什么,由算子决定

    例子 常见原因 恢复条件
    TableScan kWaitForSplit、kWaitForConnector、kWaitForScanScaleUp 收到 split、connector 完成异步工作、扫描并发获准增加
    LocalExchange kWaitForProducer 本地输入队列有数据或收到完成信号
    CallbackSink / PartitionedOutput kWaitForConsumer 下游释放缓冲容量或完成回调中的等待
    HashProbe kWaitForJoinBuild JoinBridge 上的 build 结果可用
    HashBuild / HashProbe 的 spill 协作 kWaitForJoinProbe 等 当前探测轮结束,可以进入下一轮
    Driver 自身 kWaitForArbitration 同一查询正在进行的内存仲裁完成

    有些阻塞状态在 addInput 或 getOutput 中产生,随后通过 isBlocked 交给 Driver;另一些 isBlocked 本身会推进工作。例如 Exchange::isBlocked 会获取 split、尝试取 page,并在同时等待 split 与数据时使用 collectAny:任一事件就绪即可重新检查。

    BlockingReason 是等待原因,StopReason 是本次 Driver 执行返回的控制结果,二者不应混成一张 enum 表。TableScan 就可能以 BlockingReason::kYield 和一个已就绪 future 报告主动让出;它仍通过 blockDriver 返回 StopReason::kBlock,随后很快重新排队。

    源码:BlockingReason、TableScan 主动让出、CallbackSink、Exchange 的 collectAny。

    6.1.3 从 off-thread 到 on-thread:完整的调度交接

    先分清两份信息:BlockingReason 说明现在为什么不能继续,ContinueFuture 提供将来重新尝试的通知。future 通常由队列、split store、JoinBridge 或 I/O 组件背后的 promise 完成;Driver framework 不需要知道事件的内部实现。它负责保存执行对象、让出 worker,并在通知到来后再次申请执行机会。

    getOutput 内发现等待,经第二次 isBlocked 交出 future,再通过退出、通知、排队和 Task::enter 完成一次执行往返。
    图 6:getOutput 内发现等待,经第二次 isBlocked 交出 future,再通过退出、通知、排队和 Task::enter 完成一次执行往返。 打开 SVG 原图

    第一步,保存等待,再退出算子执行区间。 blockDriver 检查 future 有效,记录 blockedOperatorId_,创建持有 Driver、Operator 和 future 的 BlockingState。BlockingState 构造把 hasBlockingFuture 设为 true,此时 Driver 仍可能在线程上,标志是在退出前先写好的。blockDriver 返回 kBlock;runInternal 的 CancelGuard 析构调用 Task::leave,正常路径减少 numThreads_ 并清除 thread 字段。

    源码节选:Driver::blockDriver

    StopReason Driver::blockDriver(
        const std::shared_ptr<Driver>& self,
        size_t blockedOperatorId,
        ContinueFuture&& future,
        std::shared_ptr<BlockingState>& blockingState,
        CancelGuard& guard) {
      auto* op = operators_[blockedOperatorId].get();
      VELOX_CHECK(
          future.valid(),
          "The operator {} is blocked but blocking future is not valid",
          op->operatorType());
      VELOX_CHECK_NE(blockingReason_, BlockingReason::kNotBlocked);
      if (blockingReason_ == BlockingReason::kYield) {
        recordYieldCount();
      }
      blockedOperatorId_ = blockedOperatorId;
      blockingState = std::make_shared<BlockingState>(
          self, std::move(future), op, blockingReason_);
      guard.notThrown();
      return StopReason::kBlock;
    }

    第二步,已经 off-thread 后才安装恢复回调。 外层 Driver::run 收到 kBlock,先处理与 Task 终止并发的情况,再调用 BlockingState::setResume。setResume 的第一行检查 Driver 已不在线程上。这个先后顺序很重要:future 可能在交出前已经就绪,也可能在 leave 与安装 continuation 之间完成;future 会保留完成状态,不会因为回调安装稍晚就丢失唤醒。反过来,若尚未退出旧执行区间就允许重新调度,可能让同一份算子状态被两个执行区间同时访问。

    第三步,promise 完成使 continuation 可运行。 setResume 使用 QueuedImmediateExecutor 衔接 continuation,它不是专门运行查询的线程池。成功回调在 Task mutex 内更新统计并清除 hasBlockingFuture;若 Task 正在 pause,暂时不入队,交给 Task::resume。否则调用 Driver::enqueue。若 future 以异常完成,走 thenError 并记录 Task error,而非当作正常数据就绪继续跑。

    源码节选:BlockingState::setResume

    void BlockingState::setResume(std::shared_ptr<BlockingState> state) {
      VELOX_CHECK(!state->driver_->isOnThread());
      auto& exec = folly::QueuedImmediateExecutor::instance();
      std::move(state->future_)
          .via(&exec)
          .thenValue([state](auto&& /* unused */) {
            auto& driver = state->driver_;
            auto& task = driver->task();
    
            std::lock_guard<std::timed_mutex> l(task->mutex());
            if (!driver->state().isTerminated) {
              state->operator_->recordBlockingTime(state->sinceUs_, state->reason_);
              // Accumulate driver-level blocked time using high_resolution_clock,
              // matching sinceUs_ and all other driver lifecycle timing.
              driver->addDriverBlockedTime(
                  (currentTimeMicrosHires() - state->sinceUs_) * 1'000);
            }
            VELOX_CHECK(!driver->state().suspended());
            VELOX_CHECK(driver->state().hasBlockingFuture);
            driver->state().hasBlockingFuture = false;
            if (task->pauseRequested()) {
              // The thread will be enqueued at resume.
              return;
            }
            Driver::enqueue(state->driver_);
          })
          .thenError(
              folly::tag_t<std::exception>{}, [state](std::exception const& e) {
                try {
                  VELOX_FAIL(
                      "A ContinueFuture for task {} was realized with error: {}",
                      state->driver_->task()->taskId(),
                      e.what());
                } catch (const VeloxException&) {
                  state->driver_->task()->setError(std::current_exception());
                }
              });
    }

    第四步,入队只是等待 worker。 Driver::enqueue 经 enqueueInternal 把 isEnqueued 设为 true,并把捕获 Driver 强引用的工作项提交到 QueryCtx::executor()。这是“等待外部事件”到“等待 CPU 调度”的转换;队列中还没有获得 worker 的 Driver,仍不是 on-thread。回调也不直接跳回之前 getOutput 的 C++ 栈。

    源码节选:Driver::enqueue

    void Driver::enqueue(std::shared_ptr<Driver> driver) {
      process::ScopedThreadDebugInfo scopedInfo(
          driver->driverCtx()->threadDebugInfo);
      // This is expected to be called inside the Driver's Tasks's mutex.
      driver->enqueueInternal();
      if (driver->closed_) {
        return;
      }
      driver->task()->queryCtx()->executor()->add(
          [driver]() { Driver::run(driver); });
    }

    第五步,worker 运行任务,并取得 Task 的执行资格。 executor 取出工作项后调用 Driver::run,再进入 runInternal。Task::enter 在锁内检查入队标志、终止、已在执行、pause / yield 等条件;只有返回 kNone 的路径,才增加 numThreads_、调用 state.setThread 并清除等待标志。这一步才完成真正的 on-thread。拿到 executor worker 并不意味着可以越过 Task 的状态检查运行算子。

    源码节选:Task::enter

    StopReason Task::enter(ThreadState& state, uint64_t nowMicros) {
      TestValue::adjust("facebook::velox::exec::Task::enter", this);
      std::lock_guard<std::timed_mutex> l(mutex_);
      VELOX_CHECK(state.isEnqueued);
      state.isEnqueued = false;
      if (state.isTerminated) {
        return StopReason::kAlreadyTerminated;
      }
      if (state.isOnThread()) {
        return StopReason::kAlreadyOnThread;
      }
      const auto reason = shouldStopLocked();
      if (reason == StopReason::kTerminate) {
        state.isTerminated = true;
      }
      if (reason == StopReason::kNone) {
        ++numThreads_;
        if (numThreads_ == 1) {
          onThreadSince_ = nowMicros;
        }
        state.setThread();
        state.hasBlockingFuture = false;
      }
      return reason;
    }

    6.1.4 future ready、queued 和 on-thread 是三个时刻

    典型阶段isOnThreadhasBlockingFutureisEnqueued
    正常执行算子truefalsefalse
    BlockingState 已创建,正在退出短暂仍为 truetruefalse
    Task::leave 后等待事件falsetruefalse
    成功 continuation 已处理,但 Task 正在 pausefalsefalsefalse
    已提交 QueryCtx executor,等待 workerfalsefalsetrue
    worker 进入且 Task::enter 返回 kNonetruefalsefalse

    这张表描述正常路径上的代表性瞬间,不把 ThreadState 的多个字段误当作一个互斥 enum。future 刚刚 ready 但 continuation 尚未执行时,hasBlockingFuture 仍可能为 true;入队后如果发生 pause 或 terminate,Task::enter 也可能不允许运行。Task::resume 只会为尚未排队、没有未完成阻塞等待等符合条件的 Driver 再次入队,避免双重调度。

    默认恢复后,getStartingOperator() 仍返回最后一个算子,从 sink 向 source 重新检查需求与阻塞。blockedOperatorId_ 用来记录是哪一个算子阻塞,并不表示所有恢复都从这个算子直接续跑。Trace 输入写入中断时保存 traceInput_ 的特殊分支才使用被阻塞的算子作为恢复起点。保留下来的是 Driver / Operator 对象中的执行进度,旧的 runInternal 栈已经退出;这也是 Driver 可以换一个 worker 继续执行的基础。

    如果算子给出一个已就绪 future,例如主动让出的路径,框架仍可先 off-thread、再很快排队,以便其他工作获得调度机会。“有 future”不意味着必须等待很久,也不意味着等待期间占着一个线程执行 future.wait()。

    源码:CancelGuard、Driver::run、Task::resume、getStartingOperator。Multi-Round LocalMerge 的 CallbackSink 与 MergeSource就是这套调度协议的具体使用者。

    6.2 isBlocked 的角色:声明阻塞意愿 + 交出唤醒句柄

    isBlocked 是 Velox 把"异步等待"嵌进协作式调度的核心接口。看签名(Operator.h:284):

    流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

    /// 返回 kNotBlocked 表示可以继续;否则返回原因,并把 'future' 设为一个在
    /// 阻塞原因消失时兑现的 future。调用者必须等 future 完成才能再次调用。
    virtual BlockingReason isBlocked(ContinueFuture* future) = 0;
    

    要点:**isBlocked 不是 Driver 探测算子状态的工具,而是算子主动声明"我要等外部事件、请让出线程"并交出唤醒句柄的接口。** 它干两件事:

    1. **返回 BlockingReason**(BlockingReason.h):语义标签,如 kWaitForSplit、kWaitForConsumer、 kWaitForProducer、kWaitForJoinBuild、kWaitForArbitration、kWaitForRPC 等;
    2. **通过 future 出参交出"唤醒句柄"**:这个 ContinueFuture 是算子与外部世界之间的"门铃",外部事件就绪时谁兑现它谁就负责唤醒 Driver。

    在普通异步 kBlock 路径上,future 把本次等待源与恢复回调接起来;生产者负责成功或异常完成它。Pause、yield、终止以及串行产出有另外的协议,不能把这个 future 当成所有 Driver 恢复动作的唯一入口。

    让出 CPU 的动作链(isBlocked 是起点和唤醒源,但"让出"由 Driver 的 return kBlock 完成):

    op->isBlocked(&future) → reason != kNotBlocked
            │
            ▼
    blockDriver(self, opId, future, ...)              // Driver.cpp:669
       构造 BlockingState(driver, future, op, reason)  // 持有 future + driver 的 shared_ptr
       return kBlock
            │
            ▼
    Driver::run 收到 kBlock → BlockingState::setResume  // Driver.cpp:227,861
       future.via(executor).thenValue([state]{
          driver->state().hasBlockingFuture = false;
          Driver::enqueue(driver);                      // future 兑现 → 重新入队唤醒
       })
       ← runInternal 已 return,线程归还线程池去跑别的 Driver
    

    6.3 三种让出线程的方式

    让出方式 触发 是否需要唤醒句柄 恢复方式 语义
    Block(isBlocked) 算子等外部事件 需要 future future 兑现 → enqueue 异步等待的核心
    Yield(shouldYield) CPU 时间片耗尽 不需要 立即重新 enqueue 到队尾 公平性,不是等待
    Suspend(enterSuspended) 等 IO / 内存仲裁,保留调用栈 不需要(同步等) leaveSuspended 原地恢复 暂时退出 Task 活跃计数,但保留线程与调用栈

    只有 Block 是真正的"挂起-异步唤醒"——它把 Driver 退出栈、压缩成几个标量、靠 future 回调复活,这才是异步编程的关键。 Yield 不是等待(马上回队列重排,防霸占线程);Suspend 保留 C++ 调用栈(不退出 runInternal),用于"必须原地等完"的同步阻塞。

    6.4 典型算子的 isBlocked 解读

    不同算子等待的"外部事件"不同,isBlocked 的写法因此分化出几种鲜明的模式。

    6.4.1 (1) CallbackSink —— sink 的下游反压(延迟报告模式)

    CallbackSink(CallbackSink.cpp:48)是 sink 算子,把数据交给 consumer 回调(下游 pipeline 的 LocalExchangeQueue 或 OutputBuffer):

    流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

    void CallbackSink::addInput(RowVectorPtr input) {
      // 推给 consumer;若 consumer 满了,consumeCb_ 返回 future 并记录到成员
      blockingReason_ = consumeCb_(std::move(input), false, &future_);
    }
    
    BlockingReason CallbackSink::isBlocked(ContinueFuture* future) {
      ...
      if (blockingReason_ != BlockingReason::kNotBlocked) {
        *future = std::move(future_);                 // 把上次 addInput 记录的 future 报出来
        blockingReason_ = BlockingReason::kNotBlocked;
        return BlockingReason::kWaitForConsumer;
      }
      return BlockingReason::kNotBlocked;
    }
    

    模式:延迟报告。 真正发现"下游满了"是在 addInput 里(拿到 future_ 和 blockingReason_ 存进成员),但要等 下一轮 Driver 循环调用 isBlocked 时才报告出来。原因见第 2 章:Driver 循环里 isBlocked 在 getOutput/addInput 之前调用,所以这一轮产生的阻塞只能下一轮才暴露。kWaitForConsumer 对应第 2 章「场景 A:OutputBuffer 满」。

    6.4.2 (2) LocalExchange —— 本地交换 source(getOutput 探测 + isBlocked 报告)

    LocalExchange(LocalPartition.cpp:275)从本地 LocalExchangeQueue 拉数据:

    流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

    RowVectorPtr LocalExchange::getOutput() {
      ...
      blockingReason_ = queue_->next(&future_, pool(), &data, drained);  // queue 空则记录 future
      if (blockingReason_ != BlockingReason::kNotBlocked) {
        return nullptr;                                                   // 没数据,先返回 null
      }
      ...
    }
    
    BlockingReason LocalExchange::isBlocked(ContinueFuture* future) {
      if (blockingReason_ != BlockingReason::kNotBlocked) {
        *future = std::move(future_);                                     // 报告上次记录的 future
        auto reason = blockingReason_;
        blockingReason_ = BlockingReason::kNotBlocked;
        return reason;                                                    // 通常是 kWaitForProducer
      }
      return BlockingReason::kNotBlocked;
    }
    

    模式:getOutput 探测、isBlocked 报告——同样是延迟报告。 与 CallbackSink 对称:sink 在 addInput 探测下游, source 在 getOutput 探测上游 queue。等待的事件是"本地上游 producer 还没产出数据"(kWaitForProducer)。

    6.4.3 (3) Exchange —— 远程交换 source(isBlocked 里干实活 + collectAny)

    Exchange(Exchange.cpp:142)从远程 worker 拉数据,是少数在 isBlocked 里直接干活的算子:

    流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

    BlockingReason Exchange::isBlocked(ContinueFuture* future) {
      if (!currentPages_.empty() || atEnd_) return BlockingReason::kNotBlocked;
    
      if (!splitFuture_.valid()) getSplits(&splitFuture_);               // 主动拉取 split
    
      ContinueFuture dataFuture;
      currentPages_ = exchangeClient_->next(driverId_, ..., &dataFuture);  // 主动拉一批 page
      if (!currentPages_.empty() || atEnd_) return BlockingReason::kNotBlocked;
    
      if (splitFuture_.valid()) {
        // 同时在等数据 OR 等更多 split —— 任一就绪即唤醒
        std::vector<ContinueFuture> futures;
        futures.push_back(std::move(splitFuture_));
        futures.push_back(std::move(dataFuture));
        *future = folly::collectAny(futures).unit();                     // ← collectAny
        return BlockingReason::kWaitForSplit;
      }
      *future = std::move(dataFuture);
      return BlockingReason::kWaitForProducer;
    }
    

    模式:isBlocked 即工作点 + 多 future 用 collectAny 合并。 它不只是检查,还触发异步预取。这里同时等两件事 (数据到达、新 split 到达),用 folly::collectAny——任一就绪即唤醒,因为有任何一边推进都值得 Driver 重新尝试。

    6.4.4 (4) Merge / LocalMerge —— 多源归并(kWaitForProducer)

    Merge(Merge.cpp:75)要从多个有序 source 做归并排序,必须每个 source 都有数据才能比较:

    流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

    BlockingReason Merge::isBlocked(ContinueFuture* future) {
      const auto reason = addMergeSources(future);          // 先确保所有 source 就位
      if (reason != BlockingReason::kNotBlocked) return reason;
      ...
      if (sourceMerger_ != nullptr) {
        sourceMerger_->isBlocked(sourceBlockingFutures_);    // 收集所有 source 的阻塞 future
      }
      if (sourceBlockingFutures_.empty()) return BlockingReason::kNotBlocked;
    
      *future = std::move(sourceBlockingFutures_.back());    // 任一 source 没数据就整体阻塞
      sourceBlockingFutures_.pop_back();
      return BlockingReason::kWaitForProducer;
    }
    

    模式:归并算子受限于最慢的 source。 任何一个参与归并的 source 没数据,整个归并都无法推进(否则破坏全局有序性), 所以报 kWaitForProducer 等待该 source。

    6.4.5 (5) LocalPartition —— 多下游 sink(collectAll)

    LocalPartition(LocalPartition.cpp:630)把数据按分区扇出到 N 个下游 queue,与 Exchange 形成鲜明对比:

    流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

    BlockingReason LocalPartition::isBlocked(ContinueFuture* future) {
      if (!futures_.empty()) {
        auto blockingReason = blockingReasons_.front();
        *future = folly::collectAll(futures_.begin(), futures_.end()).unit();  // ← collectAll
        futures_.clear();
        blockingReasons_.clear();
        return blockingReason;
      }
      return BlockingReason::kNotBlocked;
    }
    

    LocalPartition 用 collectAll 收集本次产生的等待 future。 当前本地 exchange 队列共享内存管理器,producer 在数据入队并计数后获得背压通知。collectAll 表示这批通知都已就绪,不保证恢复瞬间每个队列仍有独占空间;其他 producer 可能已入队,后续操作仍须重新走流控。

    6.4.6 (6) HashProbe —— join probe(状态机驱动 + kWaitForJoinBuild)

    HashProbe(HashProbe.cpp:651)的 isBlocked 是一个内部状态机的推进器:

    流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

    BlockingReason HashProbe::isBlocked(ContinueFuture* future) {
      switch (state_) {
        case ProbeOperatorState::kWaitForBuild:           // 等 build 侧建好哈希表
          if (!future_.valid()) { setRunning(); asyncWaitForHashTable(); }
          break;
        case ProbeOperatorState::kWaitForPeers:           // 等其他 probe 完成(spill 场景)
          if (!future_.valid()) setRunning();
          break;
        ...
      }
      if (future_.valid()) { *future = std::move(future_); return ...; }  // kWaitForJoinBuild 等
    }
    

    模式:isBlocked 驱动算子状态机。 probe 在 build 完成前处于 kWaitForBuild,它通过 asyncWaitForHashTable() 拿到一个挂在 HashJoinBridge 上的 future(与第 3 章 allPeersFinished 的屏障机制呼应), 建好后才转入 kRunning。

    6.5 两种 future 组合语义:collectAny vs collectAll

    上面 Exchange 与 LocalPartition 的对比揭示了一个值得记住的设计准则:

    组合 唤醒条件 适用场景 例子
    folly::collectAny 任一 future 就绪 等多个独立来源中任何一个推进就值得重试 Exchange:等数据 OR 等 split
    folly::collectAll 全部 future 就绪 本批收集的 future 全部就绪后重试;不保证共享条件永久成立 LocalPartition:收齐本次下游等待通知后续调

    选错会有真实后果:扇出 sink 若用 collectAny,会在某个仍满的下游上反复"唤醒→立刻又阻塞"空转;source 若用 collectAll,会在数据已到但 split 未到时白白多等。

    6.6 设计模式总结

    算子 等待事件 BlockingReason future 来源 模式
    CallbackSink 下游 consumer 满 kWaitForConsumer consumeCb_ 回调 延迟报告(addInput 探测)
    LocalExchange 本地上游无数据 kWaitForProducer queue_->next 延迟报告(getOutput 探测)
    Exchange 远程数据/split kWaitForSplit/Producer exchangeClient_ isBlocked 干活 + collectAny
    Merge/LocalMerge 任一归并源无数据 kWaitForProducer 各 MergeSource 受限于最慢 source
    LocalPartition 所有满的下游 kWaitForConsumer 类 各下游 queue 扇出 + collectAll
    HashProbe build 完成/peer 完成 kWaitForJoinBuild 等 HashJoinBridge 状态机驱动

    贯穿这些算子的共同点:

    • Driver 调用 isBlocked,算子根据自身状态与外部资源决定结果;有些实现仅报告此前保存的 future,有些会在调用中推进 I/O 或内部状态机。它既是调度器查询接口,也是算子表达等待的统一协议。
    • 普通 kBlock 路径由相应 future 完成触发续调;pause、yield、barrier 与终止有各自恢复或清理机制。因此 future 是异步 blocked 路径的连接点,不是整个 Driver 生命周期唯一的恢复来源。
    • 延迟报告是常态——多数算子在 addInput/getOutput 里就拿到了 future,存入成员,下一轮 isBlocked 才报告,这与第 2 章的 Driver 循环顺序(isBlocked 先于数据操作)严丝合缝;
    • future 组合语义要匹配语义——collectAny(任一推进即重试)vs collectAll(全部解除才继续)。

    这就是 isBlocked 作为"协作式异步关键接口"的全貌:它让每个算子用统一的 (BlockingReason, ContinueFuture) 二元组,把 五花八门的外部等待(split、远程数据、本地 queue、join build、内存仲裁、RPC)翻译成 Driver 能统一处理的"挂起-唤醒"信号。


    7. Task 与 Driver 的交互与生命周期

    7.1 生命周期:引用、算子资源与内存池

    实线容器表示独占包含,箭头注明强引用关系。Task 主动解除对 Driver 的引用后,还需等待 executor、future 等持有者释放;内存池另由 Task 保活。
    图 7:实线容器表示独占包含,箭头注明强引用关系。Task 主动解除对 Driver 的引用后,还需等待 executor、future 等持有者释放;内存池另由 Task 保活。 打开原图

    7.1.1 Task 与 Driver 为什么需要显式拆开引用

    Task 的 drivers_ 保存 shared_ptr<Driver>;Driver 独占 DriverCtx,后者持有 shared_ptr<Task>,形成强引用环。Task::removeDriver 在正常关闭时清掉槽位;terminate / resume 的关闭路径也会移走相应引用。Peer barrier 可能额外持有 Driver,需要随汇合或清理一并释放。

    Driver 自己还可能被 executor lambda、BlockingState continuation 或 serial 等待回调持有。移出 drivers_ 后,Driver 不一定立即析构;它析构时才释放 DriverCtx 中的 Task 引用。

    Task 中的 driversClosedByTask_ 是用于诊断的 weak_ptr 集合,本身不额外保活 Driver。closed_ 也只是关闭标志:当前 close() 并没有用 exchange(true) 抢占关闭权,实际并发安全来自 Task 的线程和终止协议。

    源码:Driver 成员、Task 成员与弱引用、关闭实现。

    7.1.2 Future 的完成责任属于产生它的组件

    Future continuation 捕获 Driver,为等待期间的生命期提供保证。产生 future 的组件因此必须处理成功、失败、取消和关闭路径,释放 continuation 持有的引用。Task 会清理它认识的 split、bridge、exchange、barrier 等等待源;自定义算子、connector 或外部服务返回的任意 future,仍需要其实现提供结束路径。

    Promise 兑现可能立即触发 continuation,尤其使用 inline / queued-immediate executor 时。因此常见写法是“锁内移出 promise,锁外 setValue”。这不是所有 setValue 都必须锁外的绝对规则:makeFinishFutureLocked 会立即兑现一个尚未交给调用方的新 future,此时没有用户安装的回调可以重入。

    源码:恢复回调的 Driver 引用、EventCompletionNotifier、finish future。

    7.1.3 MemoryPool 的生命期可以长于 Operator

    默认内存池路径是 QueryCtx pool → task aggregate → plan-node aggregate → operator leaf。Hash Join 的 node pool 还可能按 split group 区分;connector 有 aggregate child,exchange 有自己的 leaf pool。当前版本也支持从 QueryCtx 的 custom roots 镜像创建 task / node / operator 池。

    Task 用 childPools_ 和 customChildPools_ 持有这些池,裸指针 map 用于查找。算子 close 会释放它持有的执行资源,但 pool 对象被保留到 Task 的相应清理阶段,以支持向量和缓冲区跨 Driver 共享。保活 pool 对象与保留所有已分配数据是两件事:数据仍由实际 buffer / vector 的所有权控制。

    析构时 Task 按顺序清理 Driver、共享结构、node map、child pools、task pool 和 QueryCtx 等成员,最后兑现 deletion promises。正常所有 Driver 完成后也会尝试清理 spill 目录,析构路径另有兜底。

    源码:默认与自定义 task pool、Operator pool、池所有权成员、Task 析构。

    7.2 引用关系图

    Task ──── shared_ptr ────► Driver (drivers_[i])
      ▲                           │
      │                      unique_ptr
      │                           │
      │                        DriverCtx
      │                           │
      └─── shared_ptr ────────────┘ (ctx_->task)
    
    形成循环引用:Task → Driver → DriverCtx → Task
    
    解除循环引用的时机:
      drivers_[i] = nullptr  (Task::removeDriver 或 Task::terminate 中)
      → 此后 Driver 只剩 DriverCtx 持有的 Task 引用
      → Driver 析构后,Task 引用计数归零(若外部也无引用)
    

    额外的两个 Driver 引用:

    1. BlockingState 中:std::shared_ptr<Driver> driver_,在 future 未兑现期间延长 Driver 生命。
    2. Executor 队列 lambda 中:[driver]() { Driver::run(driver); }。

    7.3 完整生命周期

    Task::create()
        │
    Task::start() / Task::next()
        │  createDriverFactoriesLocked()   — LocalPlanner 把 PlanFragment 切成 DriverFactory
        │  initializePartitionOutput()     — 初始化输出 buffer
        │  createAndStartDrivers()         — 创建 Driver,调用 Driver::enqueue()
        ▼
    [所有 Driver 在 Executor 上并发运行]
        │
        ├── 正常路径:runInternal() → kAtEnd → Driver::close() → Task::removeDriver()
        │       → checkIfFinishedLocked() → 所有 driver 完成 + output consumed
        │       → Task::terminate(kFinished)
        │
        ├── 错误路径:operator 抛异常 → task()->setError() → Task::terminate(kFailed)
        │       → on-thread Driver 在下次 shouldStop() 看到 kTerminate,退出循环,
        │         CancelGuard 触发 close();off-thread Driver 直接 closeByTask()
        │
        └── 外部取消:requestCancel()/requestAbort() → terminate(kCanceled/kAborted)(同错误路径)
    
    Task 析构 → pool_ 等资源释放 → taskDeletionPromises_ 被 fulfill
    

    7.4 Task::enter / leave —— Driver 线程注册协议

    Task::enter(state)(Task.cpp:3341):

    加锁
    isEnqueued = false
    若 isTerminated → kAlreadyTerminated
    shouldStopLocked() 检查 terminateRequested_/pauseRequested_/toYield_
    若 kNone: ++numThreads_; state.setThread()
    返回 StopReason
    

    Task::leave(state, driverCb)(Task.cpp:3382):

    加锁,检查 shouldStop(此时可能收到 terminate 请求)
    若 kTerminate 且有 driverCb:
      解锁,调用 driverCb(kTerminate)   // 在 driver 线程上关闭 driver
      再加锁
    --numThreads_
    若 numThreads_ == 0 → fulfill threadFinishPromises_(解除 pause/terminate 的等待)
    state.clearThread()
    

    关闭必须与正在执行的 Driver 协调,保证不会同时修改其 Operator。在线程上的 Driver 通常在自己的 leave/终止路径关闭;Task 对已经不在线程上的 Driver 可以取得终止资格后通过 closeByTask 清理。因此安全条件是执行资格与状态互斥,不是永久绑定某个 OS 线程。

    7.5 Suspended 状态的用途

    一般异步 I/O 使用 isBlocked/future,让 Driver 退出当前调用并归还 executor 线程。某些必须保留调用栈的同步路径(特别是发起内存仲裁)使用 Suspended:线程仍在该调用栈中,只是暂时从 Task 活跃线程计数里扣除,以便 pause/reclaim 能停稳相关 Task。

    流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

    // Task::enterSuspended():
    //   --numThreads_(让 Task 认为没有线程在跑,满足 pause 条件)
    //   numSuspensions++(Driver 仍持有线程栈)
    // 等待外部事件...
    // Task::leaveSuspended():
    //   若 pause 中,sleep 等待 resume;否则 ++numThreads_,--numSuspensions
    

    7.6 Split Group 并发控制

    Grouped Execution 模式下,Task 按 concurrentSplitGroups_ 控制并发处理的 split group 数量。每当一个 split group 的所有 Driver 完成,Task::removeDriver() → ensureSplitGroupsAreBeingProcessedLocked() (Task.cpp:1503)为下一个 queued split group 创建新 Driver。


    8. Task 的资源管理:terminate 流程与资源回收

    8.1 Task 何时结束,终态之后还要做什么

    结束统计还区分 executionEndTimeMs 与 endTimeMs:前者在执行完成条件满足时记录,后者等输出缓冲被消费后再记录。这段差值可以包含下游消费尾延迟,不能全部算成 Driver 仍在计算。见 checkIfFinishedLocked。

    TaskState 只有 Running 和四种终态。暂停、等待 split、Driver yield 都不新增 TaskState;终态选择之后还要完成清理和引用释放。
    图 8:TaskState 只有 Running 和四种终态。暂停、等待 split、Driver yield 都不新增 TaskState;终态选择之后还要完成清理和引用释放。 打开原图

    8.1.1 Finished 的条件包含输出消费

    TaskState 从 kRunning 进入 kFinished、kCanceled、kAborted 或 kFailed,主要终结路径汇入 Task::terminate。Running 是逻辑任务状态,Task 尚未 start 或所有 Driver 都在等数据时也可能处于它。

    正常完成由 checkIfFinishedLocked 判断:

    1. 所有 Driver 已完成;或者在 ungrouped execution 中,输出 pipeline 的 Driver 已全部完成,可以让上游提前停止。
    2. 若存在 PartitionedOutput,还要求 partitionedOutputConsumed_ 为真。

    setAllOutputConsumed() 和 removeDriver() 都会重新做这个判断。因此最后一个 Driver 已结束,Task 仍可能等待下游确认消费。反过来,输出 pipeline 提前结束也可能触发 kFinished,然后终止仍在工作的上游 Driver。

    源码:TaskState、完成条件、removeDriver。

    8.1.2 terminate 先决定终态,再分路径清理

    terminate 首先在 Task 锁内检查是否仍 Running,写入终态和统计,设置 terminateRequested_,处理非正常结束的 cancellation,并收集后续清理工作。已经进入终态的 Task 不会被这次调用重新改成另一种终态。

    Driver 按当前访问状态分别处理:

    • 正在执行的 Driver:由控制检查、CancelGuard 和 Task::leave 在安全的位置关闭。
    • 未暂停的 off-thread Driver:Task 取得关闭资格、移出 drivers_,在锁外调用 closeByTask。这里由清理方线程执行,不要求曾经运行该 Driver 的同一个 OS 线程。
    • pause 中的 off-thread Driver:延后到 resume 的终止分支关闭,避免与正在 reclaim 的线程竞争。

    锁外工作还包括移除输出缓冲、关闭 InMemoryExchangeClient、取消 JoinBridge、清理 split group、通知已知的 split / barrier 等待者,以及关闭预加载的数据源。共享容器的搬移与 promise 的兑现分开,减少回调重入 Task 锁的风险。

    numRunningDrivers_ 在 terminate 中会被设为 0,这是逻辑计数更新,并不证明所有算子调用已经返回。

    源码:terminate、关闭资格、Driver::close / closeByTask。

    8.1.3 三个完成边界

    边界 能说明什么 对应入口
    Task 进入终态 执行结果已确定 TaskState、终结时的 completion 通知
    活跃线程计数归零 已登记的活跃 Driver 执行区间结束;suspended 已从此计数扣除 requestPause / terminate 返回的 finish future
    Task 对象析构 最后的持有者释放引用,进入最终成员清理 taskDeletionFuture

    terminate() 返回的 future 基于 numThreads_==0;它不是“所有 OS 线程已销毁”或“Task 对象已析构”的信号。taskCompletionFuture() 若在 Running 时注册,会被 terminate 的 completion notifier 通知;若调用时已处于终态,当前实现转而使用 makeFinishFutureLocked。读等待逻辑时应结合注册时机。

    stateChangeFuture() 的通知范围更广:除终态和 split group 完成外,当前 scan split 队列从非空变空、且之后还可能收到 split 时,也会通知状态等待者,帮助上层补充输入。它不只代表 TaskState enum 的变化。

    源码:finish future、completion / deletion future、split 队列进度通知。

    Task 持有或协调多类执行资源,包括 Driver、task / child pool、exchange client、bridge、spill 目录和等待状态。但 query root 属于 QueryCtx 及其引用链,输出缓冲也有 manager / consumer 的生命周期;这些资源不都由 Task 独占。下面分别看正常完成和终止如何解除关联。

    8.2 Memory Pool 树:层级、所有权与"刻意保活"

    Task 构造时建立内存池树(Task.cpp:710 起):

    queryCtx_->pool()                                    (query 级,aggregate)
       └── pool_  = "task.<taskId>"                       (task 级,aggregate)   initTaskPool()
             └── childPools_[k] = "node.<planNodeId>"      (node 级,aggregate)   getOrAddNodePool()
                   └── childPools_[m] = "op.<id>.<pipe>.<driver>.<type>"  (operator 级,leaf)  addOperatorPool()
    

    所有权要点:

    • pool_(Task.h:1205)是 task root,aggregate 类型(只聚合不直接分配);
    • childPools_ 持有默认执行树的子池;自定义路径还有 customChildPools_ 等状态。裸指针查找表不保活对象,向量 buffer 的引用也不等于独立持有 pool 所有权。
    • nodePools_(Task.h:1215,map<id, MemoryPool*>)只存裸指针,用于按 plan node 复用 node 池;
    • Operator / connector 池都是 leaf 或 aggregate child,由 addOperatorPool(Task.cpp:781)等创建后 push_back 进 childPools_。

    关键设计——刻意保活(Task.h:1207 注释):

    "Keep plan node and operator memory pools alive for the duration of the task to allow for sharing vectors across drivers without copy."

    Driver close 调用 Operator::close,释放该算子不再持有的状态与 buffer 引用;仍被其他向量或 Driver 共享的数据会继续存活。Task 保活的是 memory pool 对象及释放入口,不意味着所有在这个 pool 分配过的物理 buffer 都必须保留到 Task 析构。

    8.3 Happy Path:自然完成

    某 Driver 的 runInternal() → sink isFinished() → kAtEnd
            │
            ▼
    Driver::close()                                   // Driver.cpp:1063(仅允许在 driver 线程)
       closeOperators()        // 逐个 Operator::close(),释放算子逻辑状态、上报 stats
       updateStats()
       closed_ = true
       Task::removeDriver(task, this)                 // Task.cpp:1440
            │
            ▼
    Task::removeDriver
       drivers_[i] = nullptr                          // ← 切断 Task→Driver 的 shared_ptr(解循环引用)
       driverClosedLocked()  → ++numFinishedDrivers_
       splitGroupState.numRunningDrivers-- / numFinishedOutputDrivers++
       allFinished = checkIfFinishedLocked()          // 所有 driver 完成 + output 已消费?
            │
            ▼ (若 allFinished)
    Task::terminate(TaskState::kFinished)             // 见 10.4 的统一清理
    

    checkIfFinishedLocked(Task.cpp:2237)的判定:numFinishedDrivers_ == numTotalDrivers_(或 ungrouped 下输出 pipeline 的 driver 全部完成),且 !hasPartitionedOutput() || partitionedOutputConsumed_(结果已被下游取走)。 注意输出是否被消费由 OutputBuffer 经 setAllOutputConsumed() 异步通知——所以最后一个 driver 完成时若 output 还没被 取走,会先返回 false,待 setAllOutputConsumed 再触发 finish。

    8.4 异常 Path:错误 / 取消 / 中止

    三种非正常终结都走同一入口,最终汇聚到 terminate:

    算子抛异常 → runInternal catch → task()->setError(eptr)   // Driver.cpp:810
    用户取消    → Task::requestCancel() → terminate(kCanceled)
    外部错误    → Task::requestAbort()  → terminate(kAborted)
    
    setError(Task.cpp:3304):
       加锁;若已非 running 或已有 exception_ 则直接返回(保证只第一个错误生效)
       exception_ = exception
       解锁
       terminate(TaskState::kFailed)
       onError_(exception_)    // 回调通知应用层
    

    与 happy path 的关键差异:Driver 此刻可能正在线程上跑。terminate 因此要分两类处理 Driver(见下)。

    8.5 terminate:统一清理漏斗(happy/error 共用)

    terminate(terminalState)(Task.cpp:2502)是所有终结的汇聚点,分"锁内"与"锁外"两段:

    锁内(持 mutex_):

    1. 若已非 running,直接返回 makeFinishFutureLocked(幂等);
    2. state_ = terminalState;记录 terminationTimeMs;取消/中止时构造 exception_;
    3. terminateRequested_ = true(on-thread 的 Driver 下次 shouldStop() 会看到);numRunningDrivers_ = 0;
    4. 遍历 drivers_,对每个非空 driver 调 enterForTerminateLocked(Task.cpp:3367):
      • 不在线程上 → 返回 kTerminate,driver 被移入 offThreadDrivers,driverClosedLocked();
      • 在线程上 / 已 pause → 返回 kAlreadyOnThread / kPause,不在这里关,留给 driver 自己的 leave 路径关;
    5. 把 exchangeClients_、barrierFinishPromises_ swap 到局部变量;barrierRequested_ = false。

    锁外(避免持锁触发回调死锁,见 11.2.3):

    1. taskCompletionNotifier.notify() / stateChangeNotifier.notify()——兑现完成/状态变更 promise;
    2. 对 offThreadDrivers 逐个 driver->closeByTask()(Driver.cpp:1077,在"已 terminate + 仿在线程"语境下关闭算子);
    3. maybeRemoveFromOutputBufferManager()——删除 OutputBuffer(释放它对 Task 的反向引用);
    4. 关闭并清空 exchangeClients——停止远程取数、避免重发请求;
    5. 收集并清空所有 SplitsStore 的待发 split promise、处理剩余 remote split;
    6. splitGroupState.clear()——清 JoinBridge、LocalExchange、barriers;
    7. 兑现 splitPromises(唤醒等 split 的 driver,让它们感知终止)、bridge->cancel()、关闭 preloadingSplits_、兑现 barrierPromises;
    8. 返回基于 numThreads_ == 0 的 finish future,表示 Task 所计入的活跃 Driver 已离开;不表示 executor 的 OS 线程被销毁,也不等于 Task 析构。

    on-thread Driver 会在自己的退出路径完成相应清理;off-thread Driver 则由 Task 在取得终止资格后通过 closeByTask 关闭。两种分工都必须保证与 Operator 正在执行的工作互斥,不能概括为 close 永远在最初运行该 Driver 的线程上执行。

    8.6 最终回收:~Task()

    当最后一个引用 Task 的 shared_ptr 释放(drivers_ 已清空解了循环引用、外部持有者也释放后),~Task() (Task.cpp:452)执行:

    流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

    removeSpillDirectoryIfExists();        // 删除 spill 目录(若曾创建)  Task.cpp:691
    // 按固定顺序逐项 clear,便于定位析构期崩溃:
    threadFinishPromises_.clear(); splitGroupStates_.clear(); taskStats_ = {};
    stateChangePromises_.clear(); taskCompletionPromises_.clear(); splitsStates_.clear();
    drivers_.clear();              // 此时通常已空
    driverFactories_.clear(); exchangeClientByPlanNode_.clear(); exchangeClients_.clear();
    nodePools_.clear();            // 先清裸指针 map
    childPools_.clear();           // ← 释放所有 node/operator 内存池(10.1 的刻意保活在此终结)
    pool_.reset();                 // ← 最后释放 task root 池
    queryCtx_.reset();
    // 最后兑现 taskDeletionPromises_(taskDeletionFuture 的等待方在此被唤醒)
    

    内存池的销毁顺序很关键:必须先 childPools_.clear()(leaf/node 池)再 pool_.reset()(root),因为 aggregate 父池要求子池先于自己销毁。removeFromTaskList()(从全局运行任务表摘除)通过 SCOPE_EXIT 保证最先注册、最后执行。

    8.7 资源回收时机汇总

    资源 释放时机 代码位置
    算子逻辑状态(哈希表、缓冲行等) Driver close 时 Operator::close() Driver::closeOperators Driver.cpp:898
    Task→Driver 的 shared_ptr(解循环引用) removeDriver(happy)/ terminate(error) Task.cpp:1464,2559
    OutputBuffer terminate → maybeRemoveFromOutputBufferManager Task.cpp:2590
    ExchangeClient terminate 关闭并清空 Task.cpp:2592
    JoinBridge / LocalExchange / barriers terminate → SplitGroupState::clear() TaskStructs.h:248
    待发 split promise / preloading split terminate 兑现 / 关闭 Task.cpp:2675,2683
    算子 / node 内存池(buffer 所在) **~Task() 的 childPools_.clear()**(刻意保活至此) Task.cpp:490
    task root 内存池 pool_ **~Task() 的 pool_.reset()**(最后) Task.cpp:491
    spill 目录 ~Task() → removeSpillDirectoryIfExists Task.cpp:470,691

    8.8 设计要点

    • terminate 是 happy / error 共用的单一漏斗(呼应 5.1.2)——清理逻辑只写一遍,状态码不同而已;
    • 内存池与 Driver 生命周期解耦——Driver close 只释放算子逻辑状态,内存池对象保活到 ~Task(),支撑跨 Driver 零拷贝向量共享;
    • on-thread Driver 自己在自己的线程上 close——terminate 只关 off-thread 的,避免 close/abort 竞态;
    • 回收顺序受所有权约束——子池先于父池、childPools_ 先于 pool_、drivers_ 清空解循环引用后 Task 才可能析构;
    • 析构期带调试探针——~Task() 用 clearStage 字符串逐步标记,便于定位 jemalloc 析构崩溃(源码 TODO 注释)。

    9. Future/Promise 管理:避免 Driver 泄漏的设计与实现

    9.1 带着状态与等待关系排查问题

    现象 优先检查
    Driver 排队时间长 Executor 负载、Driver 并行度、时间片和长算子调用
    大量 kWaitForSplit Split 来源及 noMoreSplits、队列补充;Exchange 的组合等待还可能同时等 page
    大量 kWaitForConsumer 下游输出缓冲、消费速度和容量回收
    Pause 等待不结束 仍在执行的 Driver / Operator 调用、是否正确进入 suspended section
    Task 已 Finished 但对象仍存在 外部 Task 引用、输出持有者、executor 队列与 future continuation
    Task 的 Running Driver 数为 0 结合 TaskState、活跃线程计数及 terminate 阶段解释,不能单看一个计数

    Driver 记录 queued、on-thread、blocked 时间,Task 按 pipeline / partition 汇总;grouped execution 会合并相同 partition 索引的多轮 Driver。On-thread 是时间区间统计,不等于线程 CPU 消耗,也不适合与其他重叠统计简单相加。OpCallStatus 可以标出当前正在执行的算子方法及持续时间,帮助定位长调用。

    本次同时阅读了 DriverTest.pause、yield、driverCpuTimeSlicingCheck、recursiveSuspensionCheck,以及 TaskTest.stateChangeFutureOnSplitQueueDrain、driverEnqueAfterFailedAndPausedTask 等测试源码。它们是继续追踪暂停、重入、时间片与进度通知的入口;本文没有重新运行 C++ 测试。测试名也要结合测试体理解,例如当前 serialExecutionExternalBlockable 在外部阻塞分支前有提前 return,不能仅凭名称认定该分支已被执行验证。

    源码:Driver 调度测试、时间片测试、递归 suspended 测试、split 通知测试、pause 与终止竞态测试、Serial 测试边界。

    9.2 附录:源码阅读路线

    1. Task.h、Driver.h、Operator.h:先确定执行与接口边界。
    2. LocalPlanner:跟着一个含 Join 或 LocalExchange 的计划看 pipeline 切分与并行度。
    3. Driver::runInternal:用批次时序图对照下游需求、输出、阻塞和结束分支。
    4. BlockingState、Task::enter / leave、resume:核对“离开后才能恢复”的调度协议。
    5. checkIfFinishedLocked、terminate、析构:分别核对逻辑终态、算子清理和最终释放。

    仓库中的 TaskDriverOperatorLifecycle.md 可辅助理解强引用关系;其中部分路径说明仍保留早期版本的描述。本文的 pause / resume、正常结束和关闭流程以当前 C++ 实现及调用点为准。

    Task/Driver 里散布着十来组 promise/future。它们是协作式调度的"唤醒总线",但也是最危险的地方——一个永不兑现的 future 就能让一个 Driver(及它通过 DriverCtx 拽住的整个 Task)永久泄漏。本章解读这套机制的设计与背后的防泄漏思考。

    9.3 为什么 future/promise 在这里特别危险

    回顾第 7 章的引用关系:Task ↔ Driver 是循环引用,靠 drivers_[i] = nullptr 手动切断。但 Driver 的 shared_ptr **不只存在于 drivers_**,还藏在 future 的回调闭包里:

    • BlockingState(Driver.cpp:184,并行模式):成员 std::shared_ptr<Driver> driver_,并把整个 state 捕获进 setResume 的 thenValue 闭包(Driver.cpp:232);
    • DriverBlockingState(Task.cpp:3772,串行模式):闭包捕获 driverHolder = driver_->shared_from_this();
    • Executor 队列:[driver]() { Driver::run(driver); }(Driver.cpp:302)。

    也就是说,一个阻塞中的 Driver,其生命被它正在等待的那个 future 兜住。推出核心不变量:

    每个挂起 Driver 的 future 都必须有一条"终将兑现或出错"的路径。只要有一个 future 可能永不了结,对应的 Driver 和 Task 就会泄漏。 文档 TaskDriverOperatorLifecycle.md 末尾那句警告正是此意(见第 13 章)。

    这套设计的全部精力,就是保证这个不变量在所有路径(正常、阻塞、暂停、取消、错误、提前终止)下都成立。

    9.4 全景:Task/Driver 中的 promise/future

    promise 组 等待者 兑现者 兑现时机
    BlockingState / DriverBlockingState 阻塞的 Driver 外部事件(split/数据/buildbridge…) future 兑现 → 重新 enqueue
    SplitsStore::promises_ 等 split 的 Driver addSplit / noMoreSplits / terminate 来 split / 收工 / 终止
    BarrierState::allPeersFinishedPromises join build 非末位 Driver 末位 Driver(或 terminate) 屏障完成
    threadFinishPromises_ requestPause / terminate 调用方 allThreadsFinishedLocked numThreads_ 归零
    resumePromises_ pauseRequested 调用方 Task::resume 任务恢复
    taskCompletionPromises_ taskCompletionFuture 等待方 terminate 任务终结
    stateChangePromises_ stateChangeFuture 等待方 removeDriver / terminate 状态变更
    taskDeletionPromises_ taskDeletionFuture 等待方 ~Task() 任务析构
    barrierFinishPromises_ requestBarrier 调用方 末位 barrier driver / terminate barrier 完成

    表中各等待源必须有与其所有者相匹配的完成与取消路径。Task 自有的 promise 可由 terminate/析构处理,外部 future 则需要组件响应取消;不能仅因列出了一个 future,就认为 Task 一定能强行兑现它。

    9.5 陷阱一:持锁兑现 promise → 死锁

    promise.setValue() 可能立即执行相应 continuation,也可能通过绑定的 executor 排队;QueuedImmediateExecutor 等上下文尤其要求考虑重入。为避免把回调带进临界区,先在锁内提取状态,再在锁外通知,而不是依赖“这次回调大概异步”来保证安全。

    对策是贯穿全代码的**"锁内收集、锁外兑现"**模式。两种写法:

    (a) swap 到局部变量,出锁再兑现(allThreadsFinishedLocked,Task.cpp:3532):

    流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

    std::vector<ContinuePromise> Task::allThreadsFinishedLocked() {
      std::vector<ContinuePromise> threadFinishPromises;
      threadFinishPromises.swap(threadFinishPromises_);   // 锁内只搬走
      return threadFinishPromises;                          // 调用方出锁后才 setValue
    }
    

    (b) EventCompletionNotifier——两段式 + 析构兜底(Task.cpp:51):

    class EventCompletionNotifier {
      ~EventCompletionNotifier() { notify(); }              // 析构兜底,幂等
      void activate(std::vector<ContinuePromise> promises,  // 锁内:收集 promise + callback
                    std::function<void()> callback = nullptr);
      void notify() {                                       // 锁外:兑现 + 回调,只生效一次
        if (active_) { for (auto& p : promises_) p.setValue(); ... active_ = false; }
      }
    };
    

    用法(terminate,Task.cpp:2545):锁内 taskCompletionNotifier.activate(std::move(taskCompletionPromises_), ...), 出锁后 taskCompletionNotifier.notify()。即便中途异常提前 return,析构函数也会 notify(),且 active_ 保证不重复兑现。

    9.6 陷阱二:future 永不兑现 → Driver 泄漏 → terminate 作"总清算"

    这是最致命的陷阱。一个 Driver 阻塞在某个 future 上(等 split、等 consumer、等 join build),如果任务因错误/取消而中止, 那些事件可能永远不会自然发生——split 不会再来、consumer 已死、build 永远完不成。若放任不管,这些 Driver 的 BlockingState 闭包会永久持有 Driver→Task,泄漏。

    Task::terminate 清理其拥有的 promise 集合,并通知 bridge、队列、exchange 等组件取消等待。外部 I/O 或自定义算子产生的 future 仍需要生产方实现成功、异常和取消完成路径;Task 无法兑现任意不归它所有的 future。这里的“总清算”应限定到明确接入 Task 生命周期的资源。

    流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

    // terminate 锁外清理片段:完成 Task 拥有或接入清理协议的等待源。
    // 不覆盖任意外部 future,外部组件仍须响应 cancel / close。
    for (auto& promise : splitPromises)   promise.setValue();   // 唤醒等 split 的 driver
    for (auto& bridge  : oldBridges)      bridge->cancel();     // 取消 join,唤醒等 build 的 driver
    for (auto& barrierPromise : barrierPromises) barrierPromise.setValue();
    // SplitsStore::noMoreSplits() 把每个 store 的 promises_ 全部收集后兑现
    

    被唤醒的 Driver 重新上线后,循环顶部 shouldStop() 看到 terminateRequested_ → 返回 kTerminate → 干净退出 → close() → drivers_[i]=nullptr,闭包随 future 兑现而析构,Driver 引用计数归零。**"唤醒"不是为了让它继续干活,而是为了 让它有机会发现该退出、从而释放自己。** 这就是为什么 11.2 表格里每组 promise 都有 terminate 兜底。

    对应地,makeFinishFutureLocked(Task.cpp:2710)在 numThreads_ == 0 时立即兑现而非挂起——没有线程在跑就 没什么可等的,避免凭空制造一个永不兑现的 future。

    9.7 陷阱三:double-resume / 两个线程进同一个 Driver

    future 兑现是异步的,可能与 Task 的 pause/terminate 竞争。若处理不当,同一个 Driver 可能被两个线程同时 run。 BlockingState::setResume(Driver.cpp:227)有两道防线:

    流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

    void BlockingState::setResume(std::shared_ptr<BlockingState> state) {
      std::move(state->future_).via(&exec).thenValue([state](auto&&) {
        ...
        std::lock_guard<std::timed_mutex> l(task->mutex());     // ① 在锁内决策
        driver->state().hasBlockingFuture = false;
        if (task->pauseRequested()) {
          return;                                                // ② pause 中则不 enqueue,留给 resume
        }
        Driver::enqueue(state->driver_);
      })
      .thenError(folly::tag_t<std::exception>{}, [state](std::exception const& e) {
        ... task->setError(...);                                 // ③ future 出错也要兜底转成 task error
      });
    }
    

    三个要点:

    1. 在 Task 锁外注册 resume(Driver::run 里调 setResume,此时已不持锁)——注释(Driver.cpp:863)说明这样 "若 future 已经兑现,也不会有第二个线程进入同一个 Driver";
    2. pause 检查:若任务正在 pause,不重新 enqueue,把恢复权交给 Task::resume(Task.cpp:1204)——避免 pause 期间 Driver 偷偷回到线程;
    3. thenError 兜底:folly future 若以异常兑现(包括 promise 被析构导致的 BrokenPromise),thenError 把它转成 task->setError,绝不让一个 future 静默地烂尾。DriverBlockingState::setDriverFuture(Task.cpp:3808)有完全 对称的 thenError 分支。

    Driver::run 里还有一处微妙竞态处理(Driver.cpp:856):拿到 kBlock 后若发现 task->shouldStop()==kTerminate, 直接 return 而不进 resume 模式——否则会在已终止的任务上挂一个永不兑现的 future。

    9.8 可观测性:让泄漏"看得见"

    设计者清楚这类泄漏极难调试,于是埋了两个探针:

    • **numBlockedDrivers_**(Driver.cpp:205,BlockingState 构造 ++、析构 --):进程级"当前有多少 Driver 处于 block 态"的计数,通过 BlockingState::numBlockedDrivers() 暴露。泄漏的 Driver 会让这个数永不归零;
    • **driversClosedByTask_**(Task.h:1307,vector<weak_ptr<Driver>>):记录被 Task 强行关闭的 Driver。注释 (Task.h:1303)直言其目的——"当 race/bug 导致这些 Driver 被永久持有、进而把 Task 变成僵尸时,用这个 vector 辅助 调试僵尸 Task"。用 weak_ptr 是为了观测而不延长生命。

    9.9 设计哲学总结

    原则 实现
    每个 future 必有兑现路径 正常兑现 + terminate 总清算 + thenError 兜底,三重保证
    promise 兑现必在锁外 swap-to-local / EventCompletionNotifier 两段式
    兑现幂等、提前退出也兜底 EventCompletionNotifier 析构调 notify() + active_ 标志
    唤醒是为了让 Driver 发现该退出 terminate 唤醒所有等待者,它们上线即见 kTerminate 自行 close
    无线程可等就别造 future makeFinishFutureLocked 在 numThreads_==0 时立即兑现
    防双线程进同一 Driver setResume 锁外注册 + pause 检查 + run 里的 terminate 竞态短路
    泄漏可观测 numBlockedDrivers_ 计数 + driversClosedByTask_ 僵尸探针

    一句话:Velox 不试图"小心翼翼地不泄漏",而是建立一个强不变量——"任何挂起的 Driver 都被某个 future 兜住,而每个 future 都有正常兑现、terminate 兜底、thenError 转错三条出路之一"——再让 terminate 这个单一漏斗在所有异常路径上强制兑现 一切,从结构上消灭"future 永不兑现"的可能。 这与第 11 章的"把并发正确性从靠程序员小心转化为靠结构保证"一脉相承。

    10. 为什么产出统一从 getOutput 发起

    常规批次交付通过 getOutput 返回 RowVector,再由 Driver 调用下游 addInput。算子可以在 addInput 中完成大量计算或向外部 consumer 写出;HashBuild 等 sink 还通过 bridge 交付状态,所以这里不能把“对 Driver 返回输出批次”扩大为全部计算或副作用。

    流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

    // 中间算子(Driver.cpp:668 附近)
    getOutput(op, intermediateResult);    // 上游产出
    addInput(nextOp, intermediateResult); // 推给下游;addInput 也可能执行计算
    // sink(Driver.cpp:772 附近)
    getOutput(op, result);
    

    addInput 的返回类型不直接向 Driver 返回输出批次,但它可以立即计算、过滤、更新哈希表或聚合状态。getOutput 提供后续取出结果的接口,分离的是输入接收与输出交付的节奏,而不是强制把全部计算放到 getOutput。

    10.1 核心原因:解耦"消费输入速率"与"产生输出速率"

    计算可以放在 addInput,且可以只更新内部状态而不立即交付结果。真正需要避免的是把输入到达与下游输出批次强制一一绑定;将接收输入和取输出设计为独立接口,便于表达一对多、多对一、攒批和提前结束。

    算子 输入 → 输出映射 计算若在 addInput 的问题
    FilterProject 1 输入 → 0 或 1 输出(可能全被过滤) 尚可
    Aggregation N 输入 → 攒到 noMoreInput 后才吐 M 个输出 addInput 时根本无输出可推
    Unnest 1 输入行 → 爆炸成很多输出行(可能多批吐出) addInput 得"产出多批并逐个 push",控制流失控
    HashProbe 1 probe 批 → 0..N 个输出批(一行匹配多行) 一次输入要触发多次产出
    Limit N 输入 → 截断后提前结束 需在投喂中途叫停

    getOutput 把"产出"变成一个可被反复询问的动作:"你现在能给一批吗?能就给,给不出返回 null。" 于是:

    • **一次 addInput 可以对应 0..N 次 getOutput**——算子按自己的节奏产出。Aggregation 在 addInput 里只更新哈希表、 到 getOutput 才吐结果;Unnest 可以连续多次 getOutput 把一行的多行逐批吐完;
    • 输入/输出的批大小、批数量彻底解耦,算子内部可自由缓冲、攒批、spill。

    push 执行器也可以用缓冲、回调或协程支持 N:M 映射。Velox 采用独立 addInput/getOutput 接口,使 Driver 在批次边界统一调度;这一具体收益不需要以其他执行模型“做不到”为前提。

    10.2 getOutput 空返回与阻塞、完成分别判断

    getOutput 返回空只表示本次没有可交付批次。Driver 还需检查 isBlocked、isFinished、noMoreInput 传播与下游需求,才能决定继续扫描、等待还是结束;空返回本身不等于让出线程。Stop/pause/yield 检查发生在调度循环的相应边界,addInput 内部同样可能执行计算并影响一次调用的耗时。

    10.3 统一性:source / 中间 / 攒批算子,Driver 一视同仁

    getOutput 让 Driver 无需区分算子类型:

    • TableScan(source)无上游,addInput 永不被调用——数据在 getOutput 里"凭空"从 split 读出;
    • FilterProject(中间)在 getOutput 里消费缓冲输入产出;
    • Aggregation(攒批)在 getOutput 里吐累积结果。

    对 Driver 而言它们完全一样:**"调 getOutput 看有没有货"**。循环里唯一的区分只是 source(i=0)没有 i-1、无需先 addInput。这与 5.1.5 的 isBlocked 统一抽象是同一种品味——用统一窄接口(getOutput 产出 / addInput 投喂 / isBlocked 等待 / isFinished 结束)让一份调度代码驱动所有算子。

    10.4 与 demand-gated push 的关系:pull 触发,push 传递

    下游 needsInput()=true       (需求信号,pull 味道)
            │
            ▼
    Driver 调上游 getOutput()    (按需触发一次产出)
            │
            ▼
    Driver 调下游 addInput()     (把产出 push 下去)
    

    "从 getOutput 发起" = 由下游需求触发的、按批产出。它既不是纯 pull(不是下游递归调上游 next()),也不是纯 push (不是上游来数据就无脑灌),而是 Driver 居中、以 getOutput 为产出节拍器的 demand-gated push(呼应第 4 章)。

    10.5 小结

    独立的 addInput/getOutput 接口允许计算与批次交付解耦:一批输入可产生零到多批输出,多批输入也可在攒批后输出。addInput 可执行计算,getOutput 空返回需结合阻塞与完成状态解释;真正让出执行线程的是相应 StopReason/future 协议,不能由空返回单独推断。


    11. 架构设计品味、代码品味与多线程编程范式

    11.1 架构设计品味

    11.1.1 "协作式调度"而非"抢占式线程"

    整个设计中最高层次的品味。朴素实现是每个 Pipeline 绑定一个 OS 线程、阻塞在 IO/锁上靠操作系统调度;Velox 反其道而行—— **Driver 不拥有线程,而是被调度到线程池上的一段段"执行片"**:

    Driver 不是线程,而是一个可被反复"挂起/恢复"的状态机。
    线程是稀缺资源(CPUThreadPoolExecutor),Driver 是廉价的逻辑单元。
    

    收益:

    • M 个 Driver 可以复用 N 个 executor 线程:普通异步等待让 Driver 退出线程,future 就绪后再入队;实际并发还受状态内存、队列和资源预算约束。同步 I/O 或 suspended 仲裁等待仍可能占住调用线程,不能由这个模型直接保证任意数量的 Pipeline 都能运行。
    • 等 split、join build 或下游消费的普通 future 路径通常让 Driver 返回 kBlock 并归还 executor 线程;发起内存仲裁的 suspended 路径则保留调用栈,不能与这些异步等待混为一谈。
    • CPU 时间片到期主动 kYield,实现公平性。

    本质上是把操作系统的线程调度,下沉为应用层的协程调度。StopReason 枚举就是这套协程调度的"返回码协议"。

    11.1.2 单一终止漏斗(terminate as a funnel)

    无论是正常完成(kFinished)、用户取消(kCanceled)、外部错误(kAborted)、还是自身算子抛异常(kFailed), 所有路径最终都汇聚到 Task::terminate(terminalState) 这一个函数。资源清理逻辑只写一遍,杜绝"正常路径清理了 A 忘了 B,异常路径清理了 B 忘了 A"的经典 bug。

    11.1.3 所有权方向与生命周期的"单调性"

    Task ⊇ Driver ⊇ Operator   (寿命包含关系)
    

    这个不变量靠所有权结构强制:Driver 通过 DriverCtx::task 这个 shared_ptr 死死拽住 Task 直到自己析构; Driver 用 unique_ptr 独占 Operator。于是 Operator 里可以放心地用裸指针访问 DriverCtx、Task, 不需要任何 weak_ptr / 生命周期检查——类型系统已经证明了它们一定还活着。

    11.1.4 循环引用的"受控泄漏"

    Task 和 Driver 在执行期间互相保活,并通过明确的 removeDriver、barrier 清理和终止路径拆开引用。使用 weak_ptr 需要重新设计访问与失效处理,可能增加检查;但 shared_ptr 引用计数也有成本,不能把当前结构称为运行期零开销或把选择原因只归结为 lock() 的性能。

    11.1.5 窄接口统一异构等待——isBlocked 作为唯一的"挂起-唤醒"抽象

    整个执行引擎里,算子要等的外部事件五花八门:等 split(TableScan)、等远程数据(Exchange)、等本地 queue (LocalExchange)、等 join build(HashProbe)、等内存仲裁、等 RPC……但 Driver 调度循环只认识一种东西: isBlocked 返回的 (BlockingReason, ContinueFuture) 二元组(Operator.h:284)。

    多数 BlockingReason 都通过同一 future 续调协议连接,但 reason 还参与统计、诊断以及某些路径的策略判断。统一接口不等于所有枚举值在全部 Driver/Task 代码中完全等价。

    这条与 4.1.1 的"协作式调度"互为表里:5.1.1 讲 Driver 如何让出线程,5.1.5 讲算子如何统一地告诉 Driver"该让出了"。

    11.2 代码品味

    11.2.1 快路径无锁,慢路径加锁

    Task::shouldStop()(Task.cpp:3505)被 Driver 主循环每个算子每轮调用:

    流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

    StopReason Task::shouldStop() {
      if (pauseRequested_) return StopReason::kPause;          // atomic 读,无锁
      if (terminateRequested_) return StopReason::kTerminate;  // atomic 读,无锁
      if (toYield_) {                                          // 只有 yield 才加锁
        std::lock_guard<std::timed_mutex> l(mutex_);
        return shouldStopLocked();
      }
      return StopReason::kNone;                                // 绝大多数走到这里
    }
    

    Task.h:1414 注释明确解释了为何这些字段是 atomic:为了让 ThreadSanitizer 不误报,同时表明"在锁外读 0/false 是安全的"。

    11.2.2 "Locked 后缀"约定——把锁契约编码进函数名

    foo()(自己加锁)与 fooLocked()(调用者必须已持锁)成对出现,把"调用前必须持有 mutex_"这一契约固化进函数签名。

    11.2.3 RAII 兜底一切

    CancelGuard 保证退出执行作用域时处理 Task::leave 等收尾,并依据是否异常及终止原因执行相应 callback。正常 kBlock/kYield 返回不能把 Driver 直接关闭,因为后面还要恢复;notThrown() 用来区分已按协议返回与异常展开,不能把 guard 概括成“每次返回都 close”。

    11.2.4 统计的双保险

    Driver.cpp:533 的 on-thread 时间统计,既在 closeOperators() 手动 finalize,又有 scope guard 兜底,并用 onThreadStartUs_ = 0 哨兵值防止重复计数。两条路径都覆盖且防双计。

    11.2.5 主循环写得极致紧凑

    runInternal 用一个 for 配合 i += 2; continue; 实现数据流调度,没有递归、没有 goto(符合 CODING_STYLE)。 注意:扁平循环不仅是风格,更是协作式调度的前提——详见第 4 章。

    11.2.6 延迟报告模式——阻塞在数据操作里发现,在下一轮 isBlocked 报告

    多数算子并不在 isBlocked 里现场探测阻塞,而是在 addInput/getOutput 真正干活时就拿到了 future 与 reason, 存进成员变量(blockingReason_ / future_),等下一轮 Driver 循环调用 isBlocked 时才报告出来。 CallbackSink(CallbackSink.cpp:20,48,在 addInput 里记录)与 LocalExchange(LocalPartition.cpp:286,275, 在 getOutput 里记录)是对称的两例。

    CallbackSink 在 addInput 中记录背压,LocalExchange 在 getOutput 中记录等待,两者都由后续 isBlocked 交给 Driver。对相邻上游算子的空 getOutput 分支,这次检查就在同一次循环迭代中发生;sink 或 addInput 中产生的阻塞由之后访问该算子时报告。这里分离的是“发现等待”和“交出等待句柄”,不能把两者概括成必须隔一轮调度。

    11.2.7 future 组合语义必须匹配等待语义——collectAny vs collectAll

    当一个算子同时等多个 future 时,用哪种组合不是风格而是正确性:

    • folly::collectAny:任一就绪即唤醒。用于"任何一边推进都值得重试"——Exchange 同时等数据和等 split (Exchange.cpp:171)。
    • folly::collectAll:全部就绪才唤醒。用于"必须所有阻塞条件都解除"——LocalPartition 扇出到多个下游 queue, 必须等所有满的下游都腾出空间(LocalPartition.cpp:633)。

    选错有真实后果:扇出 sink 若用 collectAny 会在仍满的下游上"唤醒→立刻又阻塞"空转;source 若用 collectAll 会在数据已到但 split 未到时白等。详见第 6 章 7.5。

    11.3 涉及的多线程编程范式

    范式 体现 文件/行
    协作式调度(无栈协程) Driver 主动返回 StopReason 让出线程,现场压缩成几个标量后可重入 runInternal
    三种让出语义区分 Block(等外部事件,需 future)/ Yield(时间片,回队尾)/ Suspend(保留栈,原地等) 第 6 章 7.3
    Future/Promise 续延 阻塞→ContinueFuture,兑现时 thenValue 重新入队 BlockingState::setResume Driver.cpp:227
    窄接口统一异步等待 异构等待全部归一成 (BlockingReason, ContinueFuture) 二元组 Operator::isBlocked Operator.h:284
    唤醒句柄所有权 谁掌握"事件就绪"的知情权(queue/exchangeClient/JoinBridge/consumer)谁负责兑现 future 第 6 章 7.4
    future 组合语义 collectAny(任一就绪即重试)vs collectAll(全部解除才继续) Exchange.cpp:171 / LocalPartition.cpp:633
    状态机驱动的异步 isBlocked 推进算子内部状态机,跨状态挂起/恢复 HashProbe::isBlocked HashProbe.cpp:651
    同步屏障 join build 多线程汇合,最后一个 fulfill 其余 promise Task::allPeersFinished Task.h:542
    线程计数 + 等待信号 numThreads_ 归零时 fulfill threadFinishPromises_ allThreadsFinishedLocked Task.cpp:3532
    可重入挂起 numSuspensions 递归计数,保留调用栈与线程,暂时退出 Task 活跃线程计数 enterSuspended/leaveSuspended Task.cpp:3424
    快慢路径分离(lock elision) atomic flag 走无锁快路径 shouldStop Task.cpp:3505
    RAII 资源/信号管理 CancelGuard、makeGuard 兑现 promise 贯穿 Driver.cpp / Task.cpp
    锁内收集、锁外触发副作用 swap 到局部变量出锁再执行 Task::leave、Task::terminate
    协作式取消 terminateRequested_ 标志 + 轮询,并桥接 folly::CancellationToken Task::terminate、getCancellationToken
    单写者所有权消除竞态 Driver unique_ptr 独占 Operator operators_ Driver.h:732

    11.3.1 一条贯穿全局的主线

    让"线程"成为可被任意复用的纯计算资源,让"Driver"成为可被反复挂起/恢复/取消的轻量状态机,并用所有权结构 + 单一终止漏斗 + RAII,把并发正确性从"靠程序员小心"转化为"靠类型与结构保证"。

    Task::mutex_ 串行化相关 Task 状态,Operator、bridge、队列和内存系统还各有自己的同步域。ThreadState 保存的是调度状态而非完整 CPU 寄存器现场;“无栈协程”是理解其保存状态、退出、续调模式的类比,并不表示使用了 C++ coroutine 或一把全局锁。


    12. 从执行进度、等待关系与所有权看调度设计

    这套执行模型同时管理三件事:数据如何前进、线程何时可以执行、状态由谁持有。把它们分开,才能在扫描 I/O、输出反压、Join barrier 和内存回收之间切换,而不要求每条 pipeline 永久占用一个线程。

    选择 解决的约束 需要维护的边界
    Driver 保存进度,Executor 提供线程 等待期间释放线程,多个 Driver 共享线程资源 恢复时不能并发进入同一个 Driver,状态必须在调度间保持有效
    由统一循环协调算子接口 把输入需求、输出、阻塞与结束放到显式控制流中 接口返回值的含义必须一致,空输出不能被当作唯一完成信号
    用 future 表达等待条件 算子说明何时可以继续,调度器负责重新安排执行 正常完成、取消和异常路径都要处理 waiter 与持有的引用
    先更新共享状态,再在锁外通知 避免通知触发同步 continuation 时重新进入临界区 锁内必须先建立可被观察的完整状态,不能只把 fulfill 移到外面
    用进入、离开和暂停协议协调回收 回收执行状态前确定相关线程已经到达安全条件 暂停请求、停稳、恢复是不同阶段,不能用一个布尔开关概括

    工程技巧的价值取决于它维护哪个不变量。RAII guard 让特定退出动作覆盖异常路径;显式状态限制恢复时允许做的事;生命周期清理负责收束仍在飞行的异步工作。它们帮助减少遗漏,但并不能替代对锁顺序、回调可重入性和强引用关系的分析。

    读到一处调度或清理实现时,可以沿三条线核对:它改变了哪份执行进度,唤醒条件由谁兑现,最后一个持有者在什么条件下释放。这样的分析比给整套系统贴上 push、pull 或“无锁”的单一标签更能解释真实行为。前面的调度实例、逐项 future 分析和资源回收代码共同给出了这些问题的答案。

    一个可复用的技巧是 OperatorCtx 返回 peer Operator 时的 aliasing shared_ptr:访问粒度缩到算子,所有权仍绑定拥有它的 Driver,减少调用者重复查找而不扩大 Operator 的独立生命周期。它的前提是同 pipeline 对应 operatorId 和 operatorType 一致;代价是仍需明确清除 barrier 与异步持有者,别把共享指针当作自动打破引用环的机制。

    13. 附录:TaskDriverOperatorLifecycle.md 中文翻译

    以下为 velox/exec/TaskDriverOperatorLifecycle.md 的完整翻译。

    Task、Driver、Operator 生命周期

    摘要:Task 的生命周期长于 Driver 和 Operator。Task 与 Driver 之间存在循环引用,这些循环引用在 Task 释放对 Driver 的引用时被解除。Task::terminate 包含了兜底逻辑,用于在异常或提前终止时清理所有资源。

    13.1 OutputBufferManager

    OutputBuffer 持有对 Task 的引用。这个引用用于在所有输出数据被消费完毕时,通过调用 task->setAllOutputConsumed()(从 OutputBuffer::deleteResults 触发)来通知 Task。

    该引用在 OutputBuffer 被销毁时释放,销毁发生在两个地方:(1) OutputBufferManager::deleteResults(正常路径); (2) OutputBufferManager::removeTask(错误路径)。

    在 Prestissimo 应用中,OutputBufferManager::deleteResults 由 TaskManager::abortResults 调用,后者由 TaskResource::abortResults(挂载在 HTTP DELETE /v1/task/{id}/results/{id} 路由)触发。

    OutputBufferManager::removeTask 则由 Task::terminate 调用。

    13.2 Operator

    Operator 不直接持有对 Task 的引用。它通过 DriverCtx 间接访问 Task。

    13.3 Driver 与 Task

    DriverCtx 存储了一个指向 Task 的 shared_ptr。Driver 独占拥有 DriverCtx,并通过 DriverCtx 间接持有对 Task 的引用。Operator 则通过裸指针引用 DriverCtx。

    流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。

    // Driver 侧:
    std::unique_ptr<DriverCtx> ctx_;
    // DriverCtx 侧:
    std::shared_ptr<Task> task;
    

    Driver 对 Task 的引用只有在 Driver 析构时才会释放。

    Task 共享拥有所有 Driver:

    std::vector<std::shared_ptr<Driver>> drivers_;
    

    Task 在两个地方存储对 Driver 的引用:(1) drivers_ 成员变量,存储所有 Driver;(2) SplitGroupState.barriers,存储包含 join build 的 Driver(用于 Task::allPeersFinished)。

    drivers_ 的引用在以下两处被清理:(1) Task::removeDriver——由 Driver::close 调用;(2) Task::terminate。

    barriers 中存储的 Driver 引用在以下情况下被清理:(1) Task::removeDriver 中,通过 SplitGroupState::clear();(2) Task::allPeersFinished,在所有 join build 流水线完成时;(3) Task::terminate 中,通过 SplitGroupState::clear()。

    Task::removeDriver 和 Task::allPeersFinished 在正常路径(包括 Driver 自身遇到错误并主动关闭的路径)中被调用。 Task::terminate 则在错误路径中被调用。

    Task 与 Driver 之间存在循环引用:Task 引用 Driver,Driver 引用 Task。这些循环引用在 Task 释放对 Driver 的引用时被解除。 Driver 至少通过 DriverCtx::task 持有一个 Task 引用,直到析构为止。因此,Driver 的生命周期不可能超过 Task 的生命周期。 此外,Driver 独占拥有其所有 Operator,所以 Operator 的生命周期也不可能超过 Driver。总结:Task 的寿命 ≥ Driver 的寿命 ≥ Operator 的寿命。

    原生命周期资料称 BlockedState,当前名称为 BlockingState。它持有 Driver,并由 future continuation 捕获;成功或异常完成后释放相应引用,具体析构时机取决于其他持有者。如果关联 promise 与 continuation 持续存活但永不完成,也没有接入取消清理,Driver 与 Task 就可能持续保活。详见当前 BlockingState::setResume。

    Executor 队列中的 lambda 同样持有对 Driver 的引用:

    当前源码摘录(1d1b76567870):velox/exec/Driver.cpp:313。省略外围声明;此片段未作为独立程序编译。

    void Driver::enqueue(std::shared_ptr<Driver> driver) {
      process::ScopedThreadDebugInfo scopedInfo(
          driver->driverCtx()->threadDebugInfo);
      // This is expected to be called inside the Driver's Tasks's mutex.
      driver->enqueueInternal();
      if (driver->closed_) {
        return;
      }
      driver->task()->queryCtx()->executor()->add(
          [driver]() { Driver::run(driver); });
    }
    

    Driver 从以下几处被加入 Executor 队列:

    • Task::start:开始运行;
    • BlockingState::setResume:阻塞 future 兑现后继续运行;
    • Task::ensureSplitGroupsAreBeingProcessedLocked:在 Grouped Execution 中启动新的 split group;
    • Task::resume:此路径目前未被使用; 译注(2026-09-20):这是所引生命周期文档的旧描述;当前 Task::resume 已用于内存回收等暂停恢复路径。
    • Driver::run 被 Task::requestYield 触发:此路径目前未被使用。 译注(2026-09-20):这是所引生命周期文档的旧描述;当前 Task::resume 已用于内存回收等暂停恢复路径。

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