1. 1. 1. HashJoin 的数据语义与跨 Pipeline 协作
  2. 2. 2. 总体架构
    1. 2.1. 2.1 一个 JoinNode,两组 Operator
  3. 3. 3. 从两侧输入到 Join 输出,再到下一轮恢复
  4. 4. 4. HashJoinBridge
    1. 4.1. 4.1 缓存和 grouped execution 是显式分支
  5. 5. 5. HashBuild — Build 阶段
    1. 5.1. 5.1 状态机(HashBuild.h:43-57)
    2. 5.2. 5.2 Build 与 Probe 的状态机
    3. 5.3. 5.3 核心流程
    4. 5.4. 5.4 Build:输入先成为可管理的行状态
  6. 6. 6. HashProbe — Probe 阶段
    1. 6.1. 6.1 最后的 build 侧输出与 dynamic filter
    2. 6.2. 6.2 状态机(ProbeOperatorState.h:34-53)
    3. 6.3. 6.3 核心流程
    4. 6.4. 6.4 Probe:查到候选还不等于可以输出
  7. 7. 7. Spill 机制
    1. 7.1. 7.1 Spill:两侧必须保留同一分区关系
    2. 7.2. 7.2 触发路径
    3. 7.3. 7.3 整体迭代(以 1 个 spill 分区为例)
    4. 7.4. 7.4 递归 Spill
  8. 8. 8. Promise/Future 生命周期管理:如何不泄漏、不卡死
    1. 8.1. 8.1 Promise/Future 的生命周期
    2. 8.2. 8.2 底层兜底:folly 的 BrokenPromise 把"挂死"转成"报错"
    3. 8.3. 8.3 VeloxPromise:给每个 promise 加"出生地"标签
    4. 8.4. 8.4 提取-加锁 / fulfill-解锁:异常安全的核心 idiom
    5. 8.5. 8.5 Task 终止:把每一个 promise 源头都抽干
    6. 8.6. 8.6 EventCompletionNotifier:RAII 保证"即使终止过程自己抛异常也 fulfill"
    7. 8.7. 8.7 BlockingState:future 在飞行途中,谁保证 driver 不被析构?
    8. 8.8. 8.8 这套设计的分层防御总览
  9. 9. 9. 多线程编程范式
    1. 9.1. 9.1 finishHashBuild:最后到达者负责协调
    2. 9.2. 9.2 协作式非阻塞算子模型(Cooperative Non-Blocking)
    3. 9.3. 9.3 Promise/Future 做跨 Pipeline 握手
    4. 9.4. 9.4 Last-Peer Barrier(最后一人做归约)
    5. 9.5. 9.5 锁的粒度:临界区只搬指针,重活全在锁外
    6. 9.6. 9.6 输入收集与最终建表的两层并行
    7. 9.7. 9.7 所有权转移用 stateCleared_ 标志位防双重处理
    8. 9.8. 9.8 内存仲裁 reclaim 与算子线程的竞争防护
    9. 9.9. 9.9 三个时机的 reclaim
    10. 9.10. 9.10 范式总览
  10. 10. 10. 从 Join 语义推导状态、共享与异步边界
  11. 11. 11. 关键文件索引
    1. 11.1. 11.1 验证和排查入口
Macduan Notes

Velox HashJoin

Velox HashJoin 的核心是两条 pipeline 之间的协作:HashBuild 收集 build 数据并发布可查询的表,HashProbe 消费 probe 数据并产生结果。启用 spill 后,这个一次性的交接变成多轮“恢复 build 分区 → probe 对应分区”的协议。

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

1. HashJoin 的数据语义与跨 Pipeline 协作

HashJoin 将一侧输入组织成可查找的 build 状态,再用另一侧的 key 查找候选行,按 Join 类型与过滤条件产生结果。Velox 中,HashBuild、HashProbe 和 HashJoinBridge 分别负责建表、查表输出与跨 pipeline 协调;它们使用 exec::HashTable,但完整 Join 语义不由哈希索引单独承担。

对象或阶段 接收什么 产生或发布什么
HashBuild build 侧批次 行数据、可查询的表及相关 spill 信息
HashJoinBridge build 完成、probe 完成、恢复轮次等事件 共享结果和继续执行所需的通知
HashProbe probe 侧批次及当前可用 build 表 候选匹配、经过 filter 的输出及相应的补行处理
spill / restore 当前轮暂时不能在内存处理的分区 后续轮需要恢复的 build 与对应 probe 数据

Join 类型决定的不只是输出列。重复 key 会产生多行候选;outer join 需要处理未匹配行;null-aware 语义和 join filter 又会影响一个候选是否形成最终结果。理解这些要求,才能解释重复链、probed 标记、输出迭代器以及 spill 时需要保存的状态。

Build 完成与 Probe 获得可用表之间存在明确交接。启用 spill 后,交接可以重复发生:每一轮准备相应 build 状态,处理对应 probe 数据,再决定是否恢复下一组分区。因此后面的桥接接口和状态机是在表达执行协议,而不只是替某个指针加一层线程安全包装。

图 1:每个 Driver 拥有自己的 Operator,bridge 协调共享表和恢复轮次。
图 1:每个 Driver 拥有自己的 Operator,bridge 协调共享表和恢复轮次。

本文于 2026-09-19 对照 Velox 1d1b76567870 重写,覆盖普通内存 join、状态机、spill、回收、缓存和异步生命周期。哈希表内部的查找与布局另见 HashTable。

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

2. 总体架构

2.1 一个 JoinNode,两组 Operator

LocalPlanner 把 build 和 probe 分配到不同 pipeline。多个 build Driver 各有 HashBuild 和局部状态;多个 probe Driver 各有 HashProbe;它们共享 HashJoinBridge 以及发布后的表。

对象 职责 主要拥有的状态
HashBuild 消费 build 输入,准备最终表 局部表、spiller、等待 future
HashJoinBridge 发布表与 spill 信息,协调下一轮 build result、分区集合、等待 promises
HashProbe 查表、join filter、结果与补行 当前表引用、输入/输出迭代状态
Task peer barrier、调度与取消 Driver、plan node 共享对象

HashBuild 是 build pipeline 的 sink,getOutput() 返回空;表通过 bridge 发布。HashJoinBridge 当前接受 shared_ptr 表,适配普通 join 和缓存共享,而不是始终进行唯一所有权交接。

HashJoin 由三个核心组件构成:

Build Pipeline (N drivers)          Probe Pipeline (M drivers)
   HashBuild × N  ──────bridge──────  HashProbe × M
                     HashJoinBridge

Bridge 负责 build/probe 交接表、结果就绪通知与 spill 轮次协调;Task 的 peer barrier、Driver 恢复、取消以及内存回收协议也参与同步。阅读调用链时需要同时找到共享状态归属与各等待源的完成责任。


3. 从两侧输入到 Join 输出,再到下一轮恢复

以 inner equi-join 为例,build 输入为 (k, name)=[(7,A), (7,B), (9,C)],probe 输入为 (k, value)=[(7,100), (8,200)],没有额外 join filter。两个 build Driver 可以分别接收前两行和最后一行;两个 probe Driver 分别接收两条 probe 输入。相同 key 的 build 行必须保留,最终结果包含 (7,100,A)、(7,100,B),key=8 没有输出。

阶段发生的交接结果与等待条件
计划与对象建立LocalPlanner 切出 build / probe pipeline;Task 建立共享 HashJoinBridge每个 Driver 有自己的 Operator;bridge 是 Task 中对应 Join 的共享协调对象。
收集 build 输入HashBuild::addInput 解码 key 并把 key / dependent columns 存入 RowContainer各 Driver 收集局部状态;此时还不能把部分表当成完整 Join 结果发布。
汇合与发布noMoreInput → finishHashBuild → OperatorCtx::allPeersFinished最后到达者收集 peer 状态、准备最终表、setHashTable;其他 build Driver 等待通知。
等待 build 的 probe 恢复tableOrFuture → ContinueFuture → Driver 退出 / 重新调度拿到表之后才查找候选,重复 key 链给出两个 build 行;匹配语义再决定输出。
普通轮结束probe 输入耗尽,必要时汇合 peers 并补输出 build 行inner join 无额外补行;right/full/semi/anti 与 null-aware 分支各有不同条件。
存在 spill 时保存对应两侧分区 → 本轮 probe 完成 → bridge 发布恢复分片重新读取 build 分区建表,再读同一 probe 分区;过大时递归分区,不能只恢复其中一侧。

下面的协调时序同时包含普通建表和“有 spill 分区”的扩展路径。普通无 spill 执行在本轮完成后结束;有待恢复分区时,builder 才进入等待 probe / 领取 shard 的循环,不能先 kFinish 再继续读分区。

Build Driver 0  Build Driver 1  HashJoinBridge  Probe Driver 0  Probe Driver 1
     │                │               │                │               │
  addInput()       addInput()         │                │               │
     │                │               │           tableOrFuture()  tableOrFuture()
     │                │               │←─ wait ───────┤               │
     │                │               │←─ wait ────────────────────────┤
  noMoreInput()    noMoreInput()      │                │               │
     │                │               │                │               │
  allPeersFinished=false  allPeersFinished=true         │               │
  → kWaitForBuild  merge tables       │                │               │
     │             setHashTable() ───→│                │               │
     │                │           notify() ────────────→               │
     │                │               │       ─────────────────────────→
     │      无待恢复分区:kFinish;有分区:进入以下 spill 轮次             │
     │             spillInputOrFuture()→ kWaitForProbe │               │
     │                │               │           addInput()       addInput()
     │                │               │           getOutput()      getOutput()
     │                │               │           noMoreInput()    noMoreInput()
     │                │               │         allPeersFinished=false
     │                │               │       allPeersFinished=true
     │                │               │←─ probeFinished() ─────────────┤
     │                │           set shard     │               │
     │             spillInputOrFuture()→shard   │               │
     │             setupSpillInput()     │      │               │
     │             processSpillInput()   │      │               │
     │             setHashTable() ──────→│      │               │
     │                │               notify()──→               │
     │                │               │      addSpillInput()     │
     │                │               │      getOutput()         │
     │                │               │      probeFinished()─────→
     │                │           spillPartitionSet empty        │
     │                │           → notify builders              │
     │             spillInputOrFuture() → empty → kFinish       kFinish

4. HashJoinBridge

4.1 缓存和 grouped execution 是显式分支

setupCachedHashTable 支持 query 范围内多个 task 复用 hash table。没有显式 cache key 时使用 queryId 与 planNodeId 组合;一个 task 构建,其余可能等待或直接复用。缓存表有独立的 pool 生命周期,不应被普通 Task Operator 的销毁提前释放。

当前 canSpill 明确禁用 cached hash table 和 counting join 的 spill。mixed grouped join 又有自己的配置与 concurrentSplitGroups 条件。因此“开启 join_spill_enabled 就一定可 spill”是错误的。

Mixed grouped 的 probeFinished(restart=true) 还能复用此前保留的 spilled partitions,重置轮次而不重新通知 build;它不是普通 ungrouped 执行每轮恢复的同一分支。阅读普通流程时先固定执行模式,再扩展这些分支。

文件: velox/exec/HashJoinBridge.h / HashJoinBridge.cpp

HashJoinBridge 是两组 pipeline 交接表、等待就绪和推进 spill 恢复的主要对象。Task 的 peer barrier、共享表、内存回收和取消协议也参与协调,因此不能将 bridge 说成所有通信与同步的唯一通道。

关键数据成员(HashJoinBridge.h:159-203):

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

uint32_t numBuilders_;                              // build driver 数量
std::optional<HashBuildResult> buildResult_;        // 合并后的 hash table
std::optional<SpillPartitionId> restoringSpillPartitionId_; // 当前在 restore 的 spill 分区
std::vector<SpillPartition*> restoringSpillShards_; // 按 builder 数量拆分后的 shards
IterableSpillPartitionSet spillPartitionSet_;        // 待处理的 spill 分区栈
std::vector<ContinuePromise> promises_;             // 等待 table 的 probe 协程
bool probeStarted_;                                 // probe 是否已启动

核心接口:

方法 调用方 作用
setHashTable() Last build driver 提交合并好的 hash table,唤醒所有等待的 probe
tableOrFuture() Probe driver 获取 hash table,不就绪则挂起等
probeFinished() Last probe driver 通知 probe 结束,Bridge 调度下一个 spill 分区
spillInputOrFuture() Build driver (restore) 获取分配的 spill shard,不就绪则挂起等
appendSpilledHashTablePartitions() Probe driver 向 Bridge 追加新 spill 分区

5. HashBuild — Build 阶段

文件: velox/exec/HashBuild.h / HashBuild.cpp

5.1 状态机(HashBuild.h:43-57)

5.2 Build 与 Probe 的状态机

Build 的状态定义和主要进入原因:

状态 原因 后续
kRunning 接收输入或恢复 spill 输入 处理、等待或结束
kWaitForBuild 等 peer 合并;也用于 cache waiter 就绪后回 running
kWaitForProbe 已发布本轮结果,等待 probe 结束 领取下一分区或 finish
kYield 恢复文件输入时主动让出 CPU 下一次继续恢复
kFinish 全部工作完成 终态

kYield 在 processSpillInput 中直接赋值并设置 ready future;不能仅看 checkStateTransition 的分支就漏掉它。参见 恢复输入与 yield、isBlocked。

Probe 使用 kRunning / kWaitForBuild / kWaitForPeers / kFinish。等待 future 就绪后先回 running,再决定下一动作;不能把示意图中的 WaitForPeers → WaitForBuild 理解成代码直接跳过 running。

图 2:两侧等待的是不同事件;spill 让发布与探测形成重复轮次。
图 2:两侧等待的是不同事件;spill 让发布与探测形成重复轮次。

当前 noMoreInputInternal 的 peer 等待不仅用于 spill:配置允许并行输出 build 侧行时,也需要等所有 probe 完成。只照抄 ProbeOperatorState.h 中较窄的注释会漏掉这个路径。

kRunning → kWaitForBuild → kWaitForProbe → (kRunning 循环) → kFinish
                                               (spill restore)
状态 含义
kRunning 正在接收并处理输入行
kYield 处理完 spill 数据后主动让出 CPU
kWaitForBuild 非 last driver,等待 last driver 完成 table 合并
kWaitForProbe 等待 probe 完成,才能开始 restore 下一个 spill 分区
kFinish 完成

5.3 核心流程

5.4 Build:输入先成为可管理的行状态

HashBuild 从输入列解析 join key 和 dependent columns,使用 VectorHasher、RowContainer 等结构维护数据。哈希表中的列顺序通常是 keys 在前、dependents 在后,不保证与输入 RowVector 的列顺序相同。hashJoinTableType 构造这一布局。

重复值能否丢弃由 join 语义决定。某些无额外 filter 的 semi/anti join 只需知道存在性,不需要所有 payload;普通 inner join 则要保留一对多匹配。null key 的处理也依赖 null-aware 及 join 类型,不能统一“删掉所有 null”。

addInput 在维护表前检查内存和 spill 状态;已经进入 spill 的输入可能直接分区写出。最后一个输入并不立即让每个 build Driver 独立发布自己的表,必须先经过 peer barrier。

(1)initialize() (HashBuild.cpp:123)

setupCachedHashTable() → 若命中缓存直接返回
setupTable()           → 创建 HashTable<IgnoreNull>
setupSpiller()         → 若可 spill 则创建 HashBuildSpiller

(2)addInput() (HashBuild.cpp:442)

  • 调用 ensureInputFits() 做内存预留(不够则触发 spill)
  • 调用 computeSpillPartitions() 计算每行属于哪个 spill 分区
  • 被标记为 spilling 的分区行 → 调 spillInput() 写盘
  • 尚未转入 spill 的行先进入可管理的 RowContainer 状态。特定去重/counting 路径可在输入阶段建立相应索引;常规 join 会收齐各 build peer 的行后,通过 prepareJoinTable 构造最终查找结构。lookup_ 是输入与查找的临时载体,不是持久保存 join 行的容器。

(3)finishHashBuild() — 多 driver 汇聚(HashBuild.cpp:808)

这里是多 driver 并行 build 的核心:

allPeersFinished()?
   ├── false(非 last driver)→ kWaitForBuild,挂起
   └── true(last driver)→
         ① 遍历所有 peer build driver,收集各自的部分 table
         ② ensureTableFits(totalNumRows)  // 预留内存做 merge
         ③ table_->prepareJoinTable(otherTables, executor)  // 并行 merge
         ④ joinBridge_->setHashTable(mergedTable, spillPartitions)
         ⑤ 唤醒所有等待的 probe driver

关键设计: 每个 build driver 各自持有一个 HashTable,last driver 负责 merge 所有 partial table,然后一次性交给 bridge。非 last driver 在 merge 完成后由 last driver 通过 stateCleared_ 标记释放其数据所有权。

(4)postHashBuildProcess() → Spill Restore 循环(HashBuild.cpp:1029)

while (true):
  shard = joinBridge_->spillInputOrFuture()
  if (empty)   → kFinish
  if (waiting) → kWaitForProbe

  setupSpillInput(shard)   // 重置 table + spiller,创建 spillInputReader
  processSpillInput()      // 逐 batch 读取 → addInput() → 重新 build
  noMoreInput()            // → 回到 finishHashBuild(),提交新 table

6. HashProbe — Probe 阶段

6.1 最后的 build 侧输出与 dynamic filter

动态过滤还受逻辑类型与 connector 存储表示的限制。对以 VARCHAR / VARBINARY 为 backing 的注册自定义类型,当前 HashProbe 会跳过这条字符串动态过滤生成路径:内存中的 StringView 字节不一定等于扫描器面对的物理字节,贸然下推可能误删真正匹配的行。回退只是少一次过滤优化,Join 本身仍按完整类型语义匹配。见 pushdownDynamicFilters 的 custom type 分支。

Right/Full 等 join 需要在 probe 结束后查看 build 行的匹配标志。当前实现可在配置允许时让多个 probe Driver 分工输出,通过 bridge 的原子计数领取尚未处理的 RowContainer;不能一概写成“只有最后一个 probe 输出所有 build 行”。参见 RowContainer 领取接口、getBuildSideOutput。

Dynamic filter 也有严格条件。当前 asyncWaitForHashTable 只在适用 join 类型、非通用 kHash 模式、不是恢复输入、没有待恢复 spill 数据且配置开启时尝试下推。若只用当前内存子表的 keys 过滤 probe 全量输入,会错误丢掉属于磁盘分区的匹配,因此 spill 条件是正确性要求。参见 dynamic filter 条件。

文件: velox/exec/HashProbe.h / HashProbe.cpp

6.2 状态机(ProbeOperatorState.h:34-53)

kWaitForBuild → kRunning → kWaitForPeers → kFinish
                   ↑ spill restore 时循环回来

6.3 核心流程

6.4 Probe:查到候选还不等于可以输出

asyncWaitForHashTable 从 bridge 取表,尚未就绪时得到 future。拿到表后设置结果迭代器、恢复分区和 input spiller,再进入正常 probe。

主要步骤是解码 join key、确定参与查找的行、计算 hash/lookup、遍历候选、应用额外 join filter,最后按输出布局组合 probe 和 build 列。一个 probe 行可能产生多个结果;一次 getOutput 也可能只返回其中一段,迭代器要跨调用保留进度。

Join 类别 关键语义
Inner 每个满足 key 和 filter 的匹配都可能输出
Left / Full probe 行无有效匹配时需要补 null
Right / Full / right-side variants 记录 build 行是否被匹配,最后处理 build 侧输出
Semi filter 按存在性保留一侧行
Semi project 产生匹配标记,null-aware 时可能涉及三值逻辑
Anti 输出无匹配者;null-aware 不能简单等同于 NOT EXISTS

只有“null-aware anti、build 存在 null key、无额外 join filter”等明确条件下,才能使用整次查询无输出的早退。当前实现对带 filter 的情况有单独处理。参见 hasNullKeys 分支 和 null-aware anti with filter 测试。

(1)asyncWaitForHashTable() (HashProbe.cpp:434)

  • 调 joinBridge_->tableOrFuture() 取 hash table
  • 若未就绪:挂起,等 setHashTable() 触发
  • 就绪后:
    • maybeSetupInputSpiller() — 若 build 有 spill 分区,则建 probe 端 spiller
    • maybeSetupSpillInputReader() — 若当前是 restore 模式,建 spillInputReader
    • 若 table 为空且 join 类型允许短路 → 直接调 noMoreInput()

(2)addInput() (HashProbe.cpp:703)

if (needToSpillInput()):
  spillInput()  // 将属于 spill 分区的 probe 行写盘

table_->prepareForJoinProbe()
table_->joinProbe(*lookup_)      // 核心:hash 查找,填充 lookup_->hits
resultIter_->reset(*lookup_)     // 初始化结果迭代器

(3)getOutput() / getOutputInternal()

listJoinResults()     // 从 resultIter_ 批量取命中行
evalFilter()          // 若有 join filter 则过滤
fillOutput()          // 组装输出 batch(probe 列 + build 列)
// RIGHT/FULL join:标记已匹配的 build 行(setProbedFlag)

(4)Build 端输出(RIGHT/FULL join)— getBuildSideOutput() (HashProbe.cpp:897)

Right/full 等 join 需要在确定全部 probe 的匹配情况后输出 build 侧行。当前实现支持由最后到达者协调后并行扫描 build 侧输出;是否只有一个 probe 负责输出取决于配置和路径,不能无条件写成单线程。还要区分未匹配输出与 right semi 等使用 probed 标志的语义。

getAndIncrementUnclaimedRowContainerId()  // 认领 row container
table_->listNotProbedRows()  // RIGHT/FULL: 未被探测的 build 行
// 或 listProbedRows()       // RIGHT SEMI FILTER
// 或 listAllRows()          // RIGHT SEMI PROJECT

(5)noMoreInputInternal() — Peer 同步(HashProbe.cpp:1822)

allPeersFinished()?
  ├── false(非 last)→ kWaitForPeers,等 last 唤醒
  └── true(last)→
       ① 设 lastProber_ = true
       ② getBuildSideOutput() 输出 build 端行(RIGHT/FULL)
       ③ joinBridge_->probeFinished()   // 通知 bridge
       ④ wakeupPeerOperators()          // 唤醒等待的 peer probe driver

joinBridge_->probeFinished() 做了什么:

  • 从 spillPartitionSet_ 弹出下一个 spill 分区
  • 调 partition->split(numBuilders_) 拆成 N 个 shard
  • 清空 buildResult_(通知 build driver 可以开始 restore)
  • 唤醒所有等待的 build driver

7. Spill 机制

7.1 Spill:两侧必须保留同一分区关系

Hash join 的 spill 不是独立保存一份 build 文件就完成了。build 写出后,probe 中可能与这些 build 行匹配的数据也必须保留,等对应表恢复后再处理。

图 3:每轮先恢复一个 build 分区,再消费对应 probe 分区;放不下时生成下一层分区。
图 3:每轮先恢复一个 build 分区,再消费对应 probe 分区;放不下时生成下一层分区。

当前 bridge 的 setHashTable 要求“非空内存表”和“本轮 spilled partition 集合”不能同时非空。它采用整表回收的路径,不能按另一个系统的混合 hash join 模型描绘成任意挑几个 bucket spill、其余 bucket 始终留在当前表。参见 setHashTable 不变量。

一轮 probe 完成后,probeFinished 选择下一个分区,把文件分成若干 shard 给 build peers。split(numBuilders) 主要按文件数分配,不承诺各 shard 字节数、行数或恢复耗时相等。

Build 用 unordered reader 重建表;Probe 只读与 restoredPartitionId 对应的输入。若该分区仍放不下,setupSpiller 推进 hash bit offset 形成子分区。SpillPartitionId 记录层级路径,不是简单“bit offset 加一个整数”。

7.2 触发路径

内存仲裁器 → HashBuild::reclaim()
             └── spiller_->spill(partitionId)  // 把某分区的 rows 写盘

addInput() 时检测:spiller_->spillTriggered(partitionId) → 后续该分区的行直接 spill,不入 table。

7.3 整体迭代(以 1 个 spill 分区为例)

Round 1 (正常):
  Build:  addInput → [分区0写盘] → setHashTable(table, spillPartitions={0})
  Probe:  取 table → 探测 → [分区0的 probe 行写盘] → probeFinished()

Round 2 (Restore):
  Bridge: 弹出分区0,split → 2 shards(假设 2 build drivers)
  Build:  各取 1 shard → setupSpillInput → processSpillInput
         → 重新 build → setHashTable(newTable, spillPartitions={})
  Probe:  取 newTable → 从 probe 端 spill 读取分区0的 probe 行 → 探测
         → probeFinished()(无更多分区)→ 结束

7.4 递归 Spill

如果 restore 时内存仍不足,会在更深的 bit range 再次 spill。SpillPartitionId 包含父分区信息,通过 maxSpillLevel 限制递归深度。


8. Promise/Future 生命周期管理:如何不泄漏、不卡死

8.1 Promise/Future 的生命周期

异步等待要同时回答:谁创建 promise、谁持有 future、正常与异常时谁唤醒、等待期间谁保持对象存活。

Bridge 把未就绪的等待者放入 promises;发布结果时在锁内取走这些 promises,锁外 notify。JoinBridge::cancel 同样设置取消标记并通知等待者,后续访问由状态检查报错。唤醒只表示“重新检查状态”,不一定表示 join 成功。

finishHashBuild 的 SCOPE_EXIT 保证 peer promises 在该作用域退出时得到处理;异步 spill 用 guard 消费已发出的任务,即使其中一个报错也不能让其他任务悬空引用即将销毁的对象。Folly BrokenPromise 只是未完成 promise 被销毁时的异常兜底,不是正常取消协议。

mutex 保护共享状态,future 表达继续执行条件,shared_ptr 保持生命周期;三者解决不同问题。这里不存在“用了 shared_ptr 就不需要同步”或“返回 future 就不占任何线程”的通用结论。

Promise/Future 使普通等待可以退出 executor 线程,但生命周期必须闭合:成功、异常和取消都要有完成责任。未完成 promise 被销毁通常让 Folly future 以 BrokenPromise 就绪;更危险的是 promise 与 continuation 仍被引用、却再无完成事件,这会持续保活 Driver。以下机制降低这些风险,并不构成任意扩展代码永不挂起的保证。

先明确两个故障模式:

  • 故障 A:尚未完成的 promise 被销毁。对于这里使用的 folly future,析构通常通过 BrokenPromise 让等待方以异常就绪;真正可能无限保活的是 promise 与 continuation 都仍被某个引用链持有、却再也没有完成或取消事件的情况。
  • 故障 B:future 无人满足 —— 通常发生在 task 异常/取消路径,正常的 fulfill 点被跳过。

8.2 底层兜底:folly 的 BrokenPromise 把"挂死"转成"报错"

这是整套设计的基石。ContinuePromise/ContinueFuture 基于 folly:

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

using ContinuePromise = VeloxPromise<folly::Unit>;
using ContinueFuture = folly::SemiFuture<folly::Unit>;

folly 的语义是:一个 Promise 析构时若尚未 fulfill,会自动给关联的 future 投递一个 BrokenPromise 异常。也就是说——

即使代码因为 bug 把一个 promise 直接丢弃了,等待方也不会无限挂起,而是会"带着错误被唤醒"。挂死(最难排查)被降级成了报错(有栈、有日志、能恢复)。

这是一条至关重要的安全网:它把"故障 A"从"静默死锁"变成了"显式异常"。Velox 的上层设计都建立在这条兜底之上。

8.3 VeloxPromise:给每个 promise 加"出生地"标签

Velox 没有裸用 folly::Promise,而是包了一层 VeloxPromise(VeloxPromise.h:41):

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

~VeloxPromise() {
  if (!this->isFulfilled()) {
    LOG(WARNING) << "PROMISE: Unfulfilled promise is being deleted. Context: "
                 << context_;
  }
}

构造时强制要求传入 context 字符串(如 "HashJoinBridge::tableOrFuture"、"Task::allPeersFinished {taskId}")。

品味点: BrokenPromise 兜底了"不挂死",但还有个问题——哪个 promise 漏了?分布式执行里几百个 promise,光知道"有人 broken"无济于事。VeloxPromise 在析构时打印 context,等于给每个 promise 烙上"出生地"。一旦出现未 fulfill 的析构,日志直接告诉你是哪条协调路径漏了。**这是把"难以复现的并发泄漏"变成"一行可 grep 的 WARNING"**——可观测性优先于事后调试。注意它只是 LOG(WARNING) 而非 VELOX_CHECK:因为在 task 取消路径下,promise 经 cancel() fulfill 后再析构是正常的,未 fulfill 析构是"可疑"而非"必错",用 warning 而非 fail 是恰当的分寸。

8.4 提取-加锁 / fulfill-解锁:异常安全的核心 idiom

JoinBridge::cancel 等已发布 future 的通知点采用“锁内取走 promise、锁外 fulfill”的结构,避免 continuation 重入受保护状态。不能把它扩大成整个引擎所有 setValue 点都如此;尚未交付且没有回调的 future 可以在受控构造路径中直接兑现。

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

void JoinBridge::cancel() {
  std::vector<ContinuePromise> promises;
  {
    std::lock_guard<std::mutex> l(mutex_);
    cancelled_ = true;
    promises = std::move(promises_);   // ① 锁内:把 promise 搬到局部变量
  }                                     //   (promises_ 被清空)
  notify(std::move(promises));          // ② 锁外:逐个 setValue
}

// JoinBridge.cpp:20
static void JoinBridge::notify(std::vector<ContinuePromise> promises) {
  for (auto& promise : promises) {
    promise.setValue();
  }
}

这个写法防了三件事:

  1. 死锁 —— setValue() 会触发下游 driver 重新入队调度,可能回调进其他持锁路径。在锁外 fulfill 避免锁重入/锁顺序死锁。
  2. promise 泄漏 —— 一旦 std::move(promises_) 执行,这批 promise 的所有权就完全转移到局部 promises 变量。无论后面发生什么(包括异常),这个局部 vector 析构时,要么已经 setValue 过,要么触发 BrokenPromise 唤醒等待方。**绝不会"留在 promises_ 里被遗忘"**。
  3. 重复 fulfill —— 搬走后 promises_ 已空,后续再调 cancel() 是幂等的 no-op。

品味点: "把待 fulfill 的 promise 从共享状态搬进栈上局部变量"是关键一招。共享状态里的 promise 容易被遗忘(要靠记得清空),而栈上局部变量一定会析构——C++ 的栈展开机制成了泄漏的最后防线。即使 notify() 中途抛异常,剩余 promise 随栈展开析构 → BrokenPromise → 等待方带错误恢复。把生命周期管理"挂靠"到语言保证的 RAII 上,而不是靠程序员的纪律。

8.5 Task 终止:把每一个 promise 源头都抽干

Task::terminate 处理 Task 自己拥有的等待源,并调用 bridge、split、buffer 等组件的取消或清理接口。自定义 Operator、connector 或外部 I/O 创建的 future 仍由产生方负责在成功、异常和取消路径完成;Task 不能枚举和强行兑现任意用户 future。

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

// Task::terminate() 内
taskCompletionNotifier.activate(std::move(taskCompletionPromises_), ...); // 任务完成
stateChangeNotifier.activate(std::move(stateChangePromises_));            // 状态变更
for (auto& splitGroupState : splitGroupStates) splitGroupState.clear();   // split
for (auto& promise : splitPromises)      promise.setValue();              // split 等待
for (auto& bridge : oldBridges)          bridge->cancel();                // join/其他 bridge
for (auto& barrierPromise : barrierPromises) barrierPromise.setValue();   // barrier 等待

每一类 promise——任务完成、状态变更、split 输入、JoinBridge 的 build/probe 等待、allPeersFinished 的 barrier 等待——都有对应的抽干语句。bridge->cancel() 内部又走 8.3 的 idiom 抽干 bridge 自己的 promises_。

新增等待源时,需要明确其所有者、成功完成点、异常传播和取消入口。Task 拥有的等待源应接入 Task 清理;组件私有的等待源则由该组件响应 cancel/close。context 日志有助于发现遗漏,但日志本身不会执行清理或自动证明覆盖。

8.6 EventCompletionNotifier:RAII 保证"即使终止过程自己抛异常也 fulfill"

终止过程本身也可能抛异常。如果 fulfill 写成裸循环,异常会跳过剩余 promise。Velox 用一个 RAII 包装兜底(Task.cpp:51):

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

class EventCompletionNotifier {
 public:
  ~EventCompletionNotifier() { notify(); }   // 析构兜底
  void activate(std::vector<ContinuePromise> promises, std::function<void()> cb = nullptr) {
    active_ = true; callback_ = std::move(cb); promises_ = std::move(promises);
  }
  void notify() {
    if (active_) {
      for (auto& p : promises_) p.setValue();
      promises_.clear();
      if (callback_) callback_();
      active_ = false;                        // 幂等:标记已通知
    }
  }
};

正常路径显式调 notify()(置 active_=false,析构不再重复);若异常跳过了显式 notify(),析构函数会补上。

品味点: 这是 8.3 idiom 的升级版——连"忘记调用 notify"或"notify 前抛异常"都用析构函数兜住。active_ 标志让它幂等(显式调过就不再调)。把"必须执行的收尾动作"绑定到对象生命周期,是 C++ 异常安全的标准答案,这里用在了 promise fulfill 上。

8.7 BlockingState:future 在飞行途中,谁保证 driver 不被析构?

还有一个隐蔽的泄漏/悬垂问题:driver 返回 future 后,执行线程离开了,是什么阻止 Driver 对象(及其内存池)在 future resolve 前被销毁? 答案是 BlockingState 持有一个 shared_ptr<Driver>(Driver.h:184):

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

class BlockingState {
  std::shared_ptr<Driver> driver_;   // keepalive:把 driver 续命
  ContinueFuture future_;
  // ...
};
// Driver.cpp:227  setResume()
void BlockingState::setResume(std::shared_ptr<BlockingState> state) {
  std::move(state->future_)
      .via(&exec)
      .thenValue([state](auto&&) {           // lambda 捕获 state → 续命整条链
        // ... 未终止则重新入队 ...
        Driver::enqueue(state->driver_);
      })
      .thenError(folly::tag_t<std::exception>{}, [state](std::exception const& e) {
        // future 带 BrokenPromise/其他错误 resolve → 设置 task error
        state->driver_->task()->setError(std::current_exception());
      });
}

续命链条:

future 回调  →  捕获 [state]  →  state (shared_ptr<BlockingState>)
                                   └─ driver_ (shared_ptr<Driver>)

只要 future 没 resolve,回调 lambda 就活着,lambda 捕获的 state 就活着,state->driver_ 就活着——driver 不会被析构。future 一旦 resolve(无论成功 thenValue 还是失败 thenError),lambda 执行完毕、state 析构、driver 引用释放。

品味点: 两点品味。其一,用 shared_ptr 的引用计数把"对象生命周期"和"异步操作生命周期"绑定——异步操作没结束,对象就不会死,从根上消除悬垂。其二,thenValue + thenError 成对出现:这正是对 8.1 BrokenPromise 的接应。promise 被丢弃 → future 带 BrokenPromise resolve → 走 thenError → setError → 整个 task 优雅失败并连锁抽干其余 promise(回到 8.4)。一条 driver 的 future 异常,会触发全 task 的有序拆除,而不是孤立地挂死。错误是会传染的,而这正是设计想要的——宁可让整个 task 带错退出并清理干净,也不要留一个静默挂起的 driver。

8.8 这套设计的分层防御总览

层次 机制 防的是
L0 语言/库兜底 folly BrokenPromise:未 fulfill 的 promise 析构 → future 收到异常 把"挂死"降级为"报错"
L1 可观测性 VeloxPromise 析构打印 context 定位哪条路径漏了 promise
L2 fulfill idiom 锁内 std::move 提取 → 锁外 notify,靠栈上局部变量 RAII 死锁 / 遗忘 / 重复 fulfill
L3 终止排水 Task::terminate() 逐项抽干每类 promise 源 异常/取消路径下的故障 B
L4 异常安全 EventCompletionNotifier 析构兜底 fulfill 终止过程自身抛异常
L5 生命周期 BlockingState 持 shared_ptr<Driver> keepalive future 飞行途中 driver 被析构(悬垂)
L6 错误传染 thenError → setError → 连锁 terminate 单点 future 异常 → 全 task 有序拆除

一句话总结 promise 生命周期哲学:

这些层次能将常见未完成 promise、异常退出和对象保活问题变成可观察的清理协议,但不能保证任意扩展代码都不会挂起。仍被引用而永不完成的外部 future 可以持续保活 Driver;排查时要沿等待源的所有权、取消回调和完成责任逐层验证。

9. 多线程编程范式

9.1 finishHashBuild:最后到达者负责协调

当前直接调用入口为 OperatorCtx::allPeersFinished。Task 仍维护 Driver barrier;OperatorCtx 把 peer Driver 转成同 operatorId 的 aliasing shared_ptr<Operator>,再交给 HashBuild。指针指向算子,引用计数保活 Driver,因此不能将它理解为 Operator 从 Driver 脱离、独立拥有一份生命周期。

finishHashBuild 先释放未用预留,再调用 Task::allPeersFinished。普通路径中:

  1. 非最后到达者保存 future,进入 kWaitForBuild。
  2. 最后到达者收集 peer Driver,检查并移动局部表与 spiller。
  3. 估算最终表所需内存,完成 spill 文件的收尾。
  4. 调用 prepareJoinTable 合并局部状态并构造查询结构。
  5. 发布 HashBuildResult,唤醒 probe;作用域退出时唤醒等待的 build peers。

最后到达者是协调者,不意味着所有合并工作必定单线程。当前代码在存在其他局部表且没有 spill 分区时允许把 query executor 传给 prepareJoinTable;表内部仍根据自身条件决定是否真正并行构建。

peer 的 table/spiller 移动在各自锁保护下完成;较重的 finishSpill 和建表不都放在这些局部锁内。stateCleared_ 标记状态已被取走,配合生命周期检查避免另一路重复处理。

HashJoin 的并发设计是整个 Velox 执行引擎并发哲学的缩影。它几乎不使用传统的"锁 + 条件变量阻塞线程"模型,而是建立在一套协作式、future 驱动的范式上。下面逐个拆解。

9.2 协作式非阻塞算子模型(Cooperative Non-Blocking)

常规异步等待通过 isBlocked 返回 future,让 Driver 离开 executor 线程,再在就绪后重新入队。但同步文件/回收路径、锁和 AsyncSource 结果收集仍可能占用并等待调用线程;这里的“非阻塞”描述的是 Driver 的异步调度协议,不是整个 HashJoin 所有调用都绝不等待。

算子通过 isBlocked() 返回一个"阻塞原因 + future"来表达"我要等":

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

// Operator.h:284
virtual BlockingReason isBlocked(ContinueFuture* future) = 0;
// BlockingReason.h:32
kWaitForJoinBuild,   // probe 等 build table
kWaitForJoinProbe,   // build 等 probe 完成(spill restore)
kWaitForMemory,
kWaitForArbitration,

当算子返回非 kNotBlocked 时,Driver 不会 park 线程,而是 return StopReason::kBlock(Driver.cpp:1379 blockDriver()),把这个 OS 线程还给 executor 线程池去跑别的 driver。等 future 被 fulfill 时,BlockingState::setResume 再把这个 driver 重新调度回线程池。

多个 Driver 可复用宿主 QueryCtx 提供的 executor 线程,线程数由宿主配置及线程池策略决定,并不固定等于 CPU 核数。Build 等表或 peer 的普通 future 路径可归还线程;仲裁中的 suspended 与同步 reclaim 路径则保留调用栈,必须单独看待。

9.3 Promise/Future 做跨 Pipeline 握手

Build 与 probe 由不同 pipeline 的 Driver 执行,普通等待不共享一条递归执行栈。Bridge 用 promise/future 表达表就绪和轮次事件;它还持有共享表与状态,Task 的 peer barrier 和资源协议也连接两侧,所以耦合并不限于几枚 future。

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

using ContinuePromise = VeloxPromise<folly::Unit>;
using ContinueFuture = folly::SemiFuture<folly::Unit>;

消费者(probe)登记 future:

// HashJoinBridge.cpp:301  tableOrFuture()
std::lock_guard<std::mutex> l(mutex_);
if (buildResult_.has_value()) return buildResult_.value();  // 已就绪,直接拿
promises_.emplace_back("HashJoinBridge::tableOrFuture");
*future = promises_.back().getSemiFuture();                 // 没就绪,登记后挂起
return std::nullopt;

生产者(build)fulfill:

// HashJoinBridge.cpp:218  setHashTable()
std::vector<ContinuePromise> promises;
{
  std::lock_guard<std::mutex> l(mutex_);
  buildResult_ = HashBuildResult(...);
  promises = std::move(promises_);   // 锁内只搬走 promise
}
notify(std::move(promises));          // 锁外 fulfill —— 关键!

品味点: 注意 notify() 在锁外执行。fulfill promise 会触发下游 driver 重新入队调度,如果在持锁状态下做,可能引发锁顺序问题甚至死锁。"锁内只搬运状态、锁外做唤醒"是贯穿全代码的纪律。folly::Unit 作为 payload 表示"这是一个纯信号,不传数据"——数据本身(table)已经放在 buildResult_ 里,future 只负责"通知就绪",信号与数据分离得很干净。

9.4 Last-Peer Barrier(最后一人做归约)

多个 build driver 并行建表后,需要一个 barrier 把它们汇聚。Velox 没有用 std::barrier 或计数信号量,而是用一个**"最后到达者胜出并负责归约"**的模式:

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

// Task.cpp:2304  allPeersFinished()
std::lock_guard<std::timed_mutex> l(mutex_);
auto& state = barriers[planNodeId];
if (++state.numRequested == numPeers) {       // 我是最后一个
  peers = std::move(state.drivers);           // 接管所有 peer 的句柄
  promises = std::move(state.allPeersFinishedPromises);
  barriers.erase(planNodeId);
  return true;                                 // 唯一返回 true 的人
}
// 不是最后一个:登记 promise 后挂起
state.drivers.push_back(callerShared);
state.allPeersFinishedPromises.emplace_back(...);
*future = state.allPeersFinishedPromises.back().getSemiFuture();
return false;

最后到达者动态承担归约协调工作,无需预设专用协调线程。它只是最后完成该 barrier 前置工作的一方,不能据此推断它最空闲或数据最少;倾斜时它反而可能是工作最多的一方。

9.5 锁的粒度:临界区只搬指针,重活全在锁外

HashBuild 在取走 peer 状态时尽量缩短局部锁范围,再执行较重的建表工作。但这不是全 join 的“绝不持锁做重活”规则:例如 bridge reclaim 的 callback 路径有自己的持锁协议,需要按具体锁与函数分析。

最关键的 table merge:

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

// HashBuild.cpp:865  锁内:仅"偷"走 peer 的 unique_ptr
{
  std::lock_guard<std::mutex> l(build->mutex_);
  otherTables.push_back(std::move(build->table_));  // O(1) 指针搬移
  spiller = std::move(build->spiller_);
}
// HashBuild.cpp:933  锁外:真正的 merge(可能耗时数秒)
table_->prepareJoinTable(std::move(otherTables), ..., executor);

peer 状态转移的关键临界区主要搬动 table/spiller 所有权,较重的 prepareJoinTable 在相应局部锁外执行。这减少了该路径的串行范围;锁持有时间仍含检查与其他状态操作,且 bridge reclaim 等路径可能在锁内执行 callback,因此不能推导整个实现几乎无锁竞争。

9.6 输入收集与最终建表的两层并行

各 build Driver 在输入阶段主要操作自己的行状态,降低共享索引争用。最后到达者负责协调所有 peer,并不意味着归约内部始终串行;prepareJoinTable 可以在符合条件时使用 executor 建立共享索引的互不重叠区间,之后再收尾。

而且 merge 本身内部还能再并行:

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

// HashBuild.cpp:933
table_->prepareJoinTable(
    std::move(otherTables), ...,
    allowParallelJoinBuild ? queryCtx()->executor() : nullptr);  // 传入 executor 做 fork/join

两层并行分别覆盖输入行收集和最终索引构建。它缩小了必须串行执行的工作,但仍有 barrier、合并、overflow 回填以及任务调度成本;并行收益应按表规模、分区和 spill 状态测量,不能说串行瓶颈已经完全消失。

9.7 所有权转移用 stateCleared_ 标志位防双重处理

last driver "偷"走 peer 的 table 时,用一个 bool 标志保证幂等:

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

// HashBuild.cpp:900
std::lock_guard<std::mutex> l(build->mutex_);
VELOX_CHECK(!build->stateCleared_, "peer state empty...");  // 断言没被偷过
build->stateCleared_ = true;                                // 标记已偷
otherTables.push_back(std::move(build->table_));            // 转移所有权

品味点: unique_ptr 的 move 语义 + stateCleared_ 标志位 + VELOX_CHECK 断言,三者组合让"所有权转移"既安全又自我验证。被偷后的 peer 在 close() 时再次置位是幂等的。用类型系统(unique_ptr 不可复制)来表达"所有权只能有一个归属",而不是靠注释约定——让编译器帮你查并发 bug。

9.8 内存仲裁 reclaim 与算子线程的竞争防护

9.9 三个时机的 reclaim

时机 回收者 关键条件
build 尚在收集输入 HashBuild::reclaim Task 已暂停,所有 build peers 可回收
表已发布、probe 尚未开始 HashJoinBridge::reclaim tableSpillFunc 有效、probeStarted 为 false
probe 已持有表 HashProbe::reclaim 保存剩余输出/输入关系,协同 probe peers

Build reclaim 协调 peer spillers,完成写出后清表并 release 未用预留。它不是只回收某一个 operator 的任意 targetBytes;当前 targetBytes 参数未用于做精确字节裁剪。

Bridge 在 tableOrFuture 被访问时标记 probeStarted,并清除 tableSpillFunc;之后回收职责转到 probe 侧。Bridge::reclaim 当前持锁调用 spill function,因此不能把“所有重活都在锁外”当成全局规则。参见 Bridge reclaim。

Probe reclaim 先处理待输出结果;若还有 probe 输入,再写出表并设置剩余输入的 spill 路由。最后清理 buffers 和表,登记恢复分区。right-side join 还要保存 probed flag,避免恢复后把已匹配 build 行再次当作未匹配。

图 4:表在生命周期的不同阶段由不同对象负责回收;发布与首次 probe 之间也有回收窗口。
图 4:表在生命周期的不同阶段由不同对象负责回收;发布与首次 probe 之间也有回收窗口。

Spill 可能由另一个线程(内存仲裁器)发起,它会回调算子的 reclaim()。这与算子自己的执行线程存在天然竞争。Velox 用两层防护:

  1. 进入算子 reclaim 前,TaskReclaimer 请求相关 Task 停稳,等待正在执行的 Driver 离开可修改算子状态的区域。发起仲裁的 Driver 可能处于 suspended,仍保留调用栈;这与普通 future blocked 完全退出线程的机制不同。
  2. NonReclaimableSectionGuard 临界区守卫:

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

// Operator.h:489
class NonReclaimableSectionGuard {
  NonReclaimableSectionGuard(Operator* op) {
    op_->nonReclaimableSection_ = true;   // 进入不可回收区
  }
  ~NonReclaimableSectionGuard() { ... }   // RAII 自动退出
};
tsan_atomic<bool> nonReclaimableSection_{false};   // tsan 可见,配合 ThreadSanitizer 检测

reclaim() 入口先检查状态,不安全就直接放弃本次回收:

// HashBuild.cpp:1310
VELOX_CHECK(!nonReclaimableSection_);
if (nonReclaimableState()) return;   // 任一 peer 处于临界态则放弃

nonReclaimableSection_ 与 Task pause/reclaim 协议配合,表达当前算子状态是否允许回收。tsan_atomic 的具体实现取决于 folly/构建配置,不应简单归纳为所有非 TSan 构建都是裸 bool;无论包装形式如何,它也不能替代 Task 停稳及对象生命周期协议。

9.10 范式总览

范式 替代了什么朴素做法 收益
协作式非阻塞算子 一算子一线程 + cv.wait() M:N 调度,线程数 = 核数,无上下文切换风暴
Promise/Future 握手 共享标志位 + 轮询/条件变量 跨 pipeline 解耦,信号与数据分离
Last-Peer Barrier std::barrier + 静态协调者 同步点与归约指派合一,动态负载均衡
关键交接路径在锁内移动所有权 持锁做重计算 减少关键临界区中的重计算;锁、barrier 与回收 callback 仍需分析
输入并行 + 条件满足时并行建表 全局锁保护共享 table 两级并行减少部分串行工作;仍有协调与 overflow 收尾
unique_ptr + 标志位转移所有权 裸指针 + 注释约定 类型与断言约束所有权;并发安全仍依赖状态和同步协议
RAII guard + tsan_atomic 手动加解锁 + 无验证 异常安全 + TSan 可验证、零生产开销

一句话总结整体设计哲学:

把"等待"变成 future、把"线程"变成可调度的 driver、把"锁"压缩成指针搬运的瞬间、把"串行归约"二次并行化、把"并发正确性"交给类型系统和 sanitizer 去验证。HashJoin 不是"用锁把并发问题摁住",而是"通过结构设计让大部分并发冲突根本不存在"。


10. 从 Join 语义推导状态、共享与异步边界

HashJoin 的主要复杂度来自“数据还没全部到齐”和“结果已经输出一部分”同时存在。实现既需要并行处理输入,也必须知道哪个阶段可以发布表、输出未匹配行、释放状态和开始恢复下一轮。

设计选择 对应的语义或资源要求 必须保留的约束
Build / Probe 分成不同 pipeline build 状态需要独立收集和准备,probe 等待可用结果 多个 Driver 的局部完成不能等同于全局建表完成
Bridge 统一发布共享结果与轮次事件 把跨 pipeline 协作集中到明确的接口 发布、取消、恢复和缓存共享都要有一致的生命周期处理
spill 以分区和轮次组织 避免把所有外部数据同时恢复进内存 两侧数据必须使用对应的分区语义,倾斜和恢复内存仍限制成功范围
保留重复链、输出位置与匹配状态 有界输出批次不能丢失尚未展开的结果;outer join 需要正确补行 这些状态跨调用、跨等待,必要时还跨 spill 保持一致
锁内提取状态、锁外通知 减少 continuation 可重入导致的锁依赖问题 通知前共享状态应已可见,异常和取消不能遗留等待者

值得学习的抽象,是 Bridge 对“哪个事件发生后谁可以继续”的集中表达,以及 build/probe 内部对执行阶段的显式保存。它们使普通内存路径和恢复路径可以复用较多处理逻辑,但不能把所有异常情形压成一个“完成”状态。

评价并行设计时,应同时看数据所有权、共享表的有效期和等待链。共享指针保证对象存活,协议保证此刻允许怎样使用对象,两者承担不同职责。前面保留的 Promise/Future、取消清理与 spill 协调细节,是这份协议成立的具体条件。

前面的重复 key 例子说明,索引命中只是候选定位:多重匹配、join filter、null 与未匹配行输出都由 Join 语义约束。下推动态过滤必须比完整 Join 更保守,允许多留候选,不能漏掉正确结果。自定义字符串类型的回退正体现了这个边界;保留正确性要优先于少读取几行。

11. 关键文件索引

11.1 验证和排查入口

当前 HashJoinTest 包含取消、spill 下取消、空 build/probe、并行构建、不同 schema 和 null-aware/filter 等用例。阅读这些测试可以验证状态转换背后的必要条件;本文没有运行完整 HashJoinTest,也不报告性能提升。

排查顺序建议是:固定 joinType / nullAware / filter / executionMode / cache / spill 配置 → 找到当前 build/probe 状态 → 检查 bridge result 与 partition id → 找等待 future 的创建者和完成者 → 对照内存与 spill 统计。

继续阅读 finishHashBuild、asyncWaitForHashTable、probeFinished 和 Spiller,可以把表的数据结构、异步协作与内存压力连成一条完整路径。

文件 位置 作用
velox/exec/HashBuild.h :43 Build 状态机定义
velox/exec/HashBuild.cpp :808 finishHashBuild() — 多 driver 汇聚 + merge
velox/exec/HashBuild.cpp :1029 postHashBuildProcess() — spill restore 循环
velox/exec/HashProbe.cpp :434 asyncWaitForHashTable() — 等待 table
velox/exec/HashProbe.cpp :1822 noMoreInputInternal() — peer 同步
velox/exec/HashJoinBridge.cpp :218 setHashTable() — build 提交 table
velox/exec/HashJoinBridge.cpp :320 probeFinished() — probe 完成,调度 restore
velox/exec/HashJoinBridge.cpp :368 spillInputOrFuture() — build 取 shard

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