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 数据,再决定是否恢复下一组分区。因此后面的桥接接口和状态机是在表达执行协议,而不只是替某个指针加一层线程安全包装。
本文于 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。
当前 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 端 spillermaybeSetupSpillInputReader()— 若当前是 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 行匹配的数据也必须保留,等对应表恢复后再处理。
当前 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();
}
}
这个写法防了三件事:
- 死锁 ——
setValue()会触发下游 driver 重新入队调度,可能回调进其他持锁路径。在锁外 fulfill 避免锁重入/锁顺序死锁。 - promise 泄漏 —— 一旦
std::move(promises_)执行,这批 promise 的所有权就完全转移到局部promises变量。无论后面发生什么(包括异常),这个局部 vector 析构时,要么已经setValue过,要么触发 BrokenPromise 唤醒等待方。**绝不会"留在 promises_ 里被遗忘"**。 - 重复 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。普通路径中:
- 非最后到达者保存 future,进入
kWaitForBuild。 - 最后到达者收集 peer Driver,检查并移动局部表与 spiller。
- 估算最终表所需内存,完成 spill 文件的收尾。
- 调用
prepareJoinTable合并局部状态并构造查询结构。 - 发布 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 行再次当作未匹配。
Spill 可能由另一个线程(内存仲裁器)发起,它会回调算子的 reclaim()。这与算子自己的执行线程存在天然竞争。Velox 用两层防护:
- 进入算子 reclaim 前,TaskReclaimer 请求相关 Task 停稳,等待正在执行的 Driver 离开可修改算子状态的区域。发起仲裁的 Driver 可能处于 suspended,仍保留调用栈;这与普通 future blocked 完全退出线程的机制不同。
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 核对关键接口、控制流、默认值与边界条件。当前源码摘录附固定版本链接;流程伪代码用于说明分支,不是可直接编译的程序。未对全文示例做独立编译或性能复测。涉及宿主集成与历史实验的数据,按各节标注的来源理解。