Velox Memory Pool and Arbitrator
Velox 的内存管理沿着两条路径工作:Pool 和 Arbitrator 决定一次申请能否取得配额;Allocator 提供实际存储,Reclaimer 把可回收的执行状态释放出来。 理解它们之间的交接,比只看一棵 Query / Task / Operator 内存池树更重要。
执行流程与相关实现核对于 2026-09-29,Velox 源码版本为 48883e8521b2。下文源码节选、历史资料和宿主集成引用各自注明版本;教学输入用于解释状态变化,未作为性能基准运行。
1. 内存管理中的资源、额度与回收职责
Velox 的内存管理既要提供分配接口,也要控制执行对象怎样使用有限额度。MemoryPool 记录和传播用量,Arbitrator 决定各个参与仲裁的 root 可以使用多少 capacity,Allocator 完成实际存储的分配与释放,Reclaimer 则把资源压力转换为执行对象能够安全完成的回收动作。
| 观察对象 | 回答的问题 | 不能直接推导出的结论 |
|---|---|---|
| allocation / used bytes | 当前分配和使用了多少内存 | 不等于整个 pool 已预留额度,也不等于进程 RSS |
| reservation | pool 路径已经取得多少可用额度 | 不代表同等大小的物理页已经被触碰或驻留 |
| root capacity | 这个 root 当前被允许占用多大额度 | 调整 capacity 本身不会完成数据搬移或 spill |
| reclaimable state | 哪些执行状态可以释放,以及怎样释放 | 有可回收字节不代表此刻可在任意线程、任意位置回收 |
Pool 树同时表达资源归属与用量汇总。Query、Task、PlanNode、Operator 的常见层次帮助定位谁在使用资源,但 aggregate pool 与 leaf pool 的分工、root 是否参加仲裁、算子是否可回收,需要分别看待。不能仅凭 pool 的名字推断全部行为。
一次容量不足的申请,会把局部分配路径接到全局资源协调上。仲裁请求先判断额度能否取得;需要实际回收时,还要等执行对象满足暂停、可回收区和具体 spill 阶段等条件。后面的 reservation、局部/全局仲裁、Reclaimer 与 Task 协议,就是这条交接路径的逐层实现。
本文对照 Velox 1d1b765678702e6b5d1b3c582373c618813c3c85(2026-09-17),于 2026-09-19 复核。源码链接固定到该提交。重点是默认 CPU 内存路径和 SharedArbitrator;自定义内存资源单独说明。线程暂停协议可结合 Velox Task & Driver 阅读。
代码片段分别标明源码节选或流程示意;流程示意省略统计、异常包装与无关分支,不是可独立编译的程序。历史资料与宿主集成保留各自版本,不能据此推断它们组成了经过构建验证的发行版本。
2. 架构总览
2.1 MemoryManager 管什么
| 组件 | 负责的事情 | 不应混淆的边界 |
|---|---|---|
MemoryManager |
创建默认 allocator / arbitrator,创建并跟踪 root pools,提供 system pools | 管理器可以显式创建,也有进程级实例接口 |
MemoryPool |
分配接口、用量统计、reservation 传播、root capacity 检查 | Aggregate 汇总子树;实际分配通过 leaf |
MemoryAllocator |
分配 / 释放存储,实施自身容量限制,可与 AsyncDataCache 协作腾空间 | 分配器也有 capacity;它的失败不等于整台机器耗尽 RAM |
MemoryArbitrator |
决定 root pool 的 capacity 如何分配和转移 | 改 capacity 不等于申请或释放物理页 |
MemoryReclaimer |
按执行对象的语义回收资源,把 pool 树连接到 Task / Operator | 能否回收、能回收多少,由具体实现和当前状态决定 |
MemoryManager::Options 默认选择 MallocAllocator;useMmapAllocator=true 才选 MmapAllocator。前者委托 malloc 等系统接口,后者管理 mmap 相关的页和 arena。不能把 Velox 的 allocator 固定写成 jemalloc,更不能把分配成功理解为所有虚拟页已经驻留在 RSS 中。
默认 arbitrator 的配置容量是:
|
arbitratorKind 为空时创建 NoopArbitrator,并非自动启用 SharedArbitrator。Noop 为新 root 直接设置 maxCapacity,不实施跨查询总配额仲裁;有限的 root 上限和 allocator 自身限制仍然存在。生产环境需要显式选择并注册所需 arbitrator。
源码:Manager 配置、构造 allocator / arbitrator、NoopArbitrator。
2.1.1 System pool 是另一棵树
Manager 创建 __sys_root__,其下包含 spilling、caching、tracing 和 shared leaf pools。这个 root 不注册到 SharedArbitrator,容量设为 kMaxMemory,但默认情况下仍使用同一个 allocator。
MemoryManager::addLeafPool() 创建的是 system root 下的 leaf。查询算子应从 query root 的子树取得 pool,而不是用这个便捷接口绕过查询配额。allocatorCapacity - arbitratorCapacity 为 cache、spill 等系统用途留下配置空间;它不是预先 mmap 出来的独占区域,也不覆盖进程里所有不经该 allocator 的分配。
2.2 组件关系(简洁版)
MemoryManager (显式实例;也提供进程级实例接口)
├── SharedArbitrator
│ capacity 配额管理 · freeNonReservedCapacity_
│ Arb 线程 ×1(全局串行决策)· Rcl 线程池 ×N(per-victim 并行执行)
│
├── MemoryAllocator 物理页 分配 / 归还
│ ├── MmapAllocator 大页 · Slab 管理
│ └── MallocAllocator malloc / aligned_alloc 后端
│
└── Memory Pool 树
├── Root Pool capacity 配额 · 唯一 capacity 检查点
│ └── Task Pool 生命周期边界
│ └── Node Pool 逻辑聚合 · 懒建 · 跨 Driver 共享
│ └── Op Pool 算子使用的 leaf 分配入口
│
│ ↑ growCapacity / shrink ←── SharedArbitrator
│ ↕ allocate / free ←→ MemoryAllocator(leaf pool 调用,含系统用途 leaf)
│
└── MemoryReclaimer 树 (镜像 Pool 树)
├── Root → QueryCtx::MemoryReclaimer 查询回收状态 + children 路由
├── Task → Task::MemoryReclaimer pause / resume 时间窗
├── Node → ParallelMemoryReclaimer fork-join 多 driver
└── Op → Operator::MemoryReclaimer 桥接算子虚函数
│
└──▶ Spiller
SpillState · SpillPartition · SpillFile
2.3 组件关系(详细版:关键字段与职责)
MemoryManager (显式实例;也提供进程级实例接口)
│ systemMemoryBytes_ · arbitrator_ · allocator_
│ pools_ : map<name, weak_ptr<MemoryPool>>
│
├── SharedArbitrator
│ capacity_ 进程总配额(启动时固定,纯逻辑配额)
│ freeNonReservedCapacity_ 当前空闲可划拨配额
│ participants_ 注册的 Root Pool 集合
│ arbitrationQueue_ 全局仲裁等待队列
│ ├── Arb 线程 ×1 串行决策:pickVictim · shrink · grant
│ └── Rcl 线程池 ×N 并行执行:pause → reclaim → resume → shrink
│
├── MemoryAllocator <<interface>>
│ allocateContiguous() · allocateNonContiguous() · free()
│ ├── MmapAllocator mappedPages_ · freeLists_ · 大页预留 · Slab 管理
│ └── MallocAllocator 系统分配接口 · allocator 容量检查
│
└── Memory Pool 树
│
├── Root Pool [kAggregate] ←── Arbitrator 管理 capacity
│ capacity_ [0, maxCapacity_] 动态调整
│ maxCapacity_ query 级硬上限(契约,不可超)
│ reservationBytes_ query 子树的 reservation 汇总(不等同 used)
│ reclaimer ──▶ QueryCtx::MemoryReclaimer(常规查询)
│ │
│ └── Task Pool [kAggregate]
│ reservationBytes_ task 子树的 reservation 汇总(不等同 used)
│ reclaimer ──▶ Task::MemoryReclaimer
│ │ weak_ptr<Task> 避免循环引用
│ │ reclaim() = requestPause → base → resume
│ │
│ └── Node Pool [kAggregate] 懒建 · 同 planNodeId 跨 Driver 共享
│ reservationBytes_
│ reclaimer ──▶ ParallelMemoryReclaimer
│ │ spillExecutor_
│ │ reclaim() = fork each op → join
│ │
│ └── Op Pool [kLeaf] 唯一 allocate 发起点
│ reservationBytes_
│ freeReservationBytes_ 本地预留缓冲(减少跨层加锁)
│ reclaimer ──▶ Operator::MemoryReclaimer
│ op_ : Operator*
│ reclaim() = op_->reclaim()
│ enterArbitration() = enterSuspended()
│ │
│ └──▶ MemoryAllocator.allocate() / free()
│
└── reclaim(targetBytes) 沿 Reclaimer 树向下级联 ──▶ Spiller
SpillConfig
SpillState
SpillPartition
SpillFile · spillExecutor
3. 一次分配如何经过配额、回收,再得到内存
以 HashAggregation 的 operator leaf pool 需要增长缓冲区为例。Task 和 node 的 aggregate pool 负责汇总;真正申请字节的是 leaf。下面假设 SharedArbitrator 已开启,root 的量化 reservation 为 120 MiB,capacity 为 128 MiB,本次 reserve 需要再向 root 提交 16 MiB;这里的 16 MiB 是传播后的 reservation 增量,不等于应用传给 allocate 的原始字节数。
| 阶段 | 调用与责任 | 账本或执行状态 |
|---|---|---|
| 建立 pool 树 | 宿主 / QueryCtx → query root → Task → node → operator leaf | root capacity 控制配额,leaf used 记录实际申请;system pool 是另一棵树。 |
| 检查已有预算 | leaf.allocate → reserve → 父链 → root | 120 + 16 > 128,现有额度不能提交本次 reservation;此时尚未取得新的物理 buffer。 |
| 按本地顺序找额度 | SharedArbitrator::growCapacity | 重查自身余量、检查上限、申请全局空闲 capacity、回收其他 root 的未用额度;这些步骤成功就不必 spill。 |
| 需要回收已用资源 | reclaimer → QueryCtx → Task pause → node / Operator | 先等待相关执行线程到达安全边界,再由算子 spill / 释放资源。used 与 reservation 下降,capacity 此时仍可不变。 |
| 把额度交还再分配 | participant.reclaim → shrink → arbitrator.freeCapacity → 请求方 grow | shrinking 降低被回收 root 的 capacity;请求方增长 capacity 与提交 reservation 在 root 锁内完成。 |
| 取得 buffer 或失败 | MemoryAllocator::allocateBytes | 预算通过仍可能遇到 allocator limit 或底层分配失败;字节路径收到 nullptr 时撤回已记的用量并报错。 |
| 使用完成后释放 | leaf.free → allocator.freeBytes → release | 物理存储先释放,reservation 再按量化与最低预留归还;capacity 不因一次 free 自动转给其他查询。 |
若可用额度仍不足,开启全局仲裁时请求进入全局等待;关闭时尝试请求方自身回收。撞到自身上限也可能在更早的 ensureCapacity 中 self-reclaim,所以不能简化为“local 永不回收、global 才 spill”。调用链入口为 MemoryPoolImpl::allocate,分支顺序见 SharedArbitrator::growCapacity。
3.1 沿同一个请求看两阶段回收
假设被选中的另一个活跃 root 原有 capacity=1 GiB、reservation=768 MiB。算子回收后 reservation 降到 512 MiB,第一阶段释放了执行资源,但 capacity 仍为 1 GiB。第二阶段按默认普通 shrink 的余量策略保留 128 MiB 空闲,在不受 participant minimum 等其他上限约束时,可退回 384 MiB capacity,留下 640 MiB。请求方能获得多少还取决于增长目标和等待者策略,并不是把这 384 MiB 全部直接交给它。
发起 allocate 的 Driver 进入 suspended 仲裁区时保留当前 C++ 调用栈,等待可能占着该线程;同 query 的其他 Driver 在循环检查到 underArbitration 时则可通过 future off-thread。前者等待分配返回,后者等待重新调度,两种等待不能画成同一条“全部让出线程”的路径。后文的全局时序图与两阶段 reclaimer 图分别展开这些线程和账本变化。
4. Presto Native 中的完整调用链
4.1 创建时序
POST /v1/task/{taskId} [TaskResource.cpp:345]
│
├─ QueryContextManager::findOrCreateQueryCtx() [QueryContextManager.cpp:80]
│ └─ MemoryManager::addRootPool(
│ "20240515_abc123_0",
│ queryMaxMemoryPerNode)
│ → ROOT POOL 创建(同一 queryId 只创建一次,有缓存)
│
└─ TaskManager::createOrUpdateTask() [TaskManager.cpp:531]
└─ Task::create() → Task::init() [Task.cpp:381, 511]
└─ Task::initTaskPool() [Task.cpp:710]
└─ rootPool->addAggregateChild("task.xxx")
→ TASK POOL 创建
└─ maybeStartTaskLocked() → Task::start()
└─ createDriversLocked()
└─ 每个 Operator 实例化时:
OperatorCtx() [Operator.cpp:31]
└─ DriverCtx::addOperatorPool()
└─ Task::addOperatorPool() [Task.cpp:781]
├─ getOrAddNodePool()
│ └─ taskPool->addAggregateChild("node.xxx")
│ → NODE POOL 创建(首次懒建)
└─ nodePool->addLeafChild("op.xxx")
→ OPERATOR POOL 创建
4.2 Pool 树示例
以 TableScan → Aggregation (2并发) → PartitionedOutput 为例:
Root Pool: "20240515_abc123_0" kAggregate capacity=10GB
└── Task Pool: "task.20240515_abc123_001.0" kAggregate
├── Node Pool: "node.1" kAggregate (TableScan)
│ └── Op Pool: "op.1.0.0.TableScan" kLeaf
│
├── Node Pool: "node.2" kAggregate (Aggregation)
│ ├── Op Pool: "op.2.1.0.Aggregation" kLeaf (Driver 0)
│ └── Op Pool: "op.2.1.1.Aggregation" kLeaf (Driver 1)
│
└── Node Pool: "node.3" kAggregate (Output)
└── Op Pool: "op.3.2.0.PartitionedOutput" kLeaf
多个 task(stage)共享同一个 Root Pool:
Root Pool: "20240515_abc123_0" capacity=10GB
├── Task Pool: "task.20240515_abc123_001.0" (stage 1)
│ └── ...
└── Task Pool: "task.20240515_abc123_002.0" (stage 2)
└── ...
5. 整体层次结构
默认查询执行路径常见 Query root → Task → Node → Operator 的树形归属,便于按执行对象记账和回收;这不是所有 MemoryPool 必须遵循的固定四层规则。System pool 是另一棵树,connector、exchange、自定义资源也可能增加或采用不同子层次,默认执行树没有单独的 Driver pool 层。
MemoryManager(全局单例)
└── Root Pool(Query 级) kAggregate — QueryCtx 持有
└── Task Pool kAggregate — Task 持有
└── Node Pool kAggregate — Task 管理,按 PlanNodeId 划分
└── Operator Pool kLeaf — Operator 持有,实际分配内存
5.1 核心文件
| 文件 | 职责 |
|---|---|
velox/common/memory/MemoryPool.h |
抽象接口与 MemoryPoolImpl |
velox/common/memory/Memory.h |
MemoryManager 生命周期管理 |
velox/common/memory/MemoryPool.cpp |
reservation / release 核心逻辑 |
velox/common/memory/MemoryArbitrator.h/cpp |
MemoryReclaimer 接口与仲裁基类 |
velox/common/memory/SharedArbitrator.h/cpp |
共享内存仲裁器实现 |
velox/common/memory/ArbitrationParticipant.h/cpp |
每个 root pool 的仲裁代理 |
velox/common/memory/ArbitrationOperation.h/cpp |
单次仲裁请求状态机 |
velox/exec/Task.cpp |
各层 pool 的创建入口 |
velox/exec/Driver.cpp |
DriverCtx::addOperatorPool() |
velox/exec/Operator.cpp |
OperatorCtx 持有 operator pool |
presto_cpp/main/QueryContextManager.cpp |
Presto 侧 root pool 创建 |
presto_cpp/main/TaskManager.cpp |
Presto 侧 task 创建编排 |
6. Pool 类型
| 类型 | 能否分配内存 | 能否创建子 Pool | 用途 |
|---|---|---|---|
kAggregate |
否 | 是 | 层次管理、聚合统计 |
kLeaf |
是 | 否 | 实际分配,追踪用量 |
Leaf pool 可选线程安全模式(threadSafe=true),同一 operator 被多线程访问时开启。
7. 各层 Pool 详解
7.1 Query pool 在哪一层创建
7.1.1 Velox 本身有默认创建路径
当前 QueryCtx::Builder::build() 构造 QueryCtx。若调用方没有传入 pool,QueryCtx::initPool() 会执行:
|
生成的名称形如 query.<queryId>.<sequence>。所以 query pool 可以由 Velox 层创建;嵌入 Velox 的查询引擎也可以先调用 addRootPool(name, queryLimit),再通过 Builder 的 .pool(...) 注入。
这里默认 maxCapacity 是 kMaxMemory。不能据此假设每个 query root 都自动读取 Presto 的 query.max-memory-per-node。该限制如何映射、QueryCtx 如何按 query ID 缓存,属于宿主引擎的集成策略;本文不把旧版 Presto Native 调用链当成当前 Velox 的固定行为。
Root 创建时 capacity 为 0,随后由其 arbitrator 的 addPool 决定初始容量。Builder 在 QueryCtx 已由 shared_ptr 持有后调用 maybeSetReclaimer:若 pool 尚无 reclaimer,就安装 QueryCtx 的 reclaimer;调用方已设置的实现会保留。
源码:QueryCtx::Builder、initPool、默认 reclaimer、创建 root。
7.1.2 默认执行树没有单独的 Driver pool 层
query root aggregate:仲裁单位
└─ task.<taskId> aggregate:Task 暂停 / 恢复边界
└─ node.<planNodeId> aggregate:按计划节点组织回收
├─ op.<node>.<pipe>.0.<type> leaf:Driver 0 的算子实例
└─ op.<node>.<pipe>.1.<type> leaf:Driver 1 的算子实例
Task::initTaskPool 创建 task aggregate;Node pool 懒创建,Task::addOperatorPool 为每个算子实例创建 leaf。pipelineId 和 driverId 出现在 leaf 名称中,不代表额外的一层内存池。
Hash Join 的 node pool 还可能按 split group 区分,使用专门的 HashJoinMemoryReclaimer;connector 使用 aggregate child,exchange client 有自己的 leaf。真实树因此比四层示意更丰富。
7.1.3 自定义资源可以有独立的 allocator / arbitrator
当前代码支持 MemoryManager::addCustomRootPool 和 QueryCtx::Builder::customPool(tag, pool)。CustomMemoryResource 提供 allocator、arbitrator、maxCapacity 和 reclaimer factory;Task 在相应 root 下创建 task / node / operator 子树。
它们不是默认 CPU allocator 账本里的另一个普通 child。若资源使用独立的 allocator / arbitrator,应分别观察其容量与统计;资源提供方还需保证被 pool 借用的 allocator / arbitrator 生命期足够长。
源码:Task 建池、Operator leaf、Join node、custom root、custom task trees。
7.2 Root Pool(Query 级)
流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。
// velox/common/memory/Memory.h:216
std::shared_ptr<MemoryPool> MemoryManager::addRootPool(
const std::string& name, // 如 "20240515_abc123_0"
int64_t maxCapacity, // query_max_memory_per_node
std::unique_ptr<MemoryReclaimer> reclaimer);
- 唯一能设置 capacity 上限的节点,子孙节点不独立设 capacity。
- 在同一个 Presto worker 内,共用同一 QueryCtx 的 task 共享其 query root pool。不同 worker 各有进程级资源;缓存失效后重新创建的 QueryCtx 也可能产生新 root,因此 queryId 不是跨进程或跨生命周期的全局 shared_ptr 身份。
- 宿主在 worker 内缓存 QueryCtx,缓存命中时复用该对象及 root。缓存失效后允许创建新的 QueryCtx,旧 root 又可能因后台回收保活而暂时存在,所以 pool 名还要加入唯一后缀;不能理解成同一 queryId 永远只创建一次。
7.3 Task Pool
当前源码摘录(1d1b76567870):velox/exec/Task.cpp:732。省略外围声明;此片段未作为独立程序编译。
void Task::initTaskPool() {
VELOX_CHECK_NULL(pool_);
pool_ = queryCtx_->pool()->addAggregateChild(
fmt::format("task.{}", taskId_.c_str()), createTaskReclaimer());
initCustomTaskPools();
}
- 名字格式:
task.<taskId>(如task.20240515_abc123_001.0) - 在
Task::init()中同步创建,Task 构造完成时即存在。 - 挂载
TaskReclaimer,支持 task 级别的内存回收(spill)。
7.4 Node Pool
当前源码摘录(1d1b76567870):velox/exec/Task.cpp:757。省略外围声明;此片段未作为独立程序编译。
velox::memory::MemoryPool* Task::getOrAddNodePool(
const core::PlanNodeId& planNodeId) {
if (nodePools_.count(planNodeId) == 1) {
return nodePools_[planNodeId];
}
childPools_.push_back(pool_->addAggregateChild(
fmt::format("node.{}", planNodeId), createNodeReclaimer([&]() {
return exec::ParallelMemoryReclaimer::create(
queryCtx_->spillExecutor());
})));
auto* nodePool = childPools_.back().get();
nodePools_[planNodeId] = nodePool;
for (const auto& [tag, _] : queryCtx_->customPools()) {
getOrAddCustomNodePool(tag, planNodeId);
}
return nodePool;
}
- 名字格式:
node.<planNodeId> - 懒创建:第一个使用该 PlanNodeId 的 Operator 实例化时才创建。
- HashJoin 特殊处理:build/probe 侧共享同一个 join node pool(名字带
[<splitGroupId>]后缀),以便跨 pipeline 聚合内存。 - 同一 plan node 的多个并发 Driver 共享同一个 Node Pool,Node Pool 聚合所有并发的内存用量。
7.5 Operator Pool
当前源码摘录(1d1b76567870):velox/exec/Task.cpp:913。省略外围声明;此片段未作为独立程序编译。
velox::memory::MemoryPool* Task::addOperatorPool(
const core::PlanNodeId& planNodeId,
uint32_t splitGroupId,
int pipelineId,
uint32_t driverId,
const std::string& operatorType) {
velox::memory::MemoryPool* nodePool;
if (isHashJoinOperator(operatorType)) {
nodePool = getOrAddJoinNodePool(planNodeId, splitGroupId);
} else {
nodePool = getOrAddNodePool(planNodeId);
}
childPools_.push_back(nodePool->addLeafChild(
fmt::format(
"op.{}.{}.{}.{}", planNodeId, pipelineId, driverId, operatorType)));
return childPools_.back().get();
}
- 名字格式:
op.<planNodeId>.<pipelineId>.<driverId>.<operatorType> - 每个 Driver 实例的每个 Operator 各自独立一个 Leaf Pool,粒度最细。
- 默认 operator pool 是 leaf,connector 相关 pool 由 Task 的专门接口按其用途创建和保活。不能把所有 TableScan/connector 情况一概画成“operator leaf 变 aggregate 再挂子池”;应区分 operator pool 与 connector pool 的创建入口。
触发点在 OperatorCtx 构造:
// velox/exec/Operator.cpp:31
OperatorCtx::OperatorCtx(DriverCtx* driverCtx, ..., std::string_view operatorType)
: pool_(driverCtx_->addOperatorPool(planNodeId, operatorType_)) {}
// ↑ → DriverCtx::addOperatorPool() → Task::addOperatorPool()
8. 内存分配与 Reservation 机制
8.1 used、reservation 与 capacity 是三个量
| 名称 | 含义与观察位置 |
|---|---|
Leaf usedBytes() |
usedReservationBytes_,已经被分配或外部用量报告消耗的记账字节;包含相应对齐,不等于业务有效载荷或 RSS |
reservedBytes() |
当前向上层占住的 reservation;leaf 有量化余量,aggregate 汇总子树 |
Leaf availableReservation() |
leaf reservation 中尚未使用的部分,可供后续申请直接消费 |
Root capacity() |
Arbitrator 当前授予这个 root 的额度,可以 grow / shrink |
maxCapacity() |
Root 的配置上限;子 pool 的该接口沿父链查询 |
freeBytes() |
有限 root 的 max(0, capacity - reservedBytes);子 pool 调用时也返回 root 口径 |
minReservationBytes_ |
Leaf 通过显式预留保住的 reservation 下限;与 arbitrator 的保留容量不是同一个概念 |
普通稳定状态下可以用 used ≤ reserved ≤ capacity ≤ maxCapacity 理解一棵有限查询树。其中 used / reserved 要用同一子树口径;仲裁上下文中允许 reservation 暂时超过 root capacity,以避免回收自身的申请再次触发嵌套仲裁。Allocator 的限制仍独立存在。
Aggregate 的 usedBytes() 递归汇总 leaf;其 peak / cumulative 则在 reservation 增量上传时更新。因此不能把所有层的 cumulativeBytes 都当成精确相同口径的 malloc 流量,或把每层统计相加得到总用量。
源码:统计接口与量化、usedBytes、freeBytes、reservation 更新。
8.1.1 量化的是 reservation,不是每次都 malloc 1 MiB
MemoryPool::quantizedSize 按总需求向上取整:
| 需求范围 | 粒度 |
|---|---|
| 小于 16 MiB | 1 MiB |
| 16 MiB 到小于 64 MiB | 4 MiB |
| 64 MiB 及以上 | 8 MiB |
先消费 leaf 已有 reservation;余量不足时,才把量化后的增量沿父链传播。比如首次申请对齐后的 100 KiB,leaf 可以先取得 1 MiB reservation,但 allocator 仍只收到这次申请的大小。后续 200 KiB 可以直接消费已有余量。
maybeReserve(increment) 是显式预留接口,先把 increment 按 8 MiB 向上取整,再走 reserve(..., reserveOnly=true),并更新 leaf 的最低 reservation。release() 清除这个下限、释放多余 reservation,不会替调用方 free 仍存活的 buffer。预留失败通常返回 false;若 pool 已被 abort,则继续抛异常,阻止使用可能已失效的执行状态。
MemoryPool 同时追踪实际 used 与预留 reservation:先保证可用预留,再调用 allocator 分配;释放时也更新实际使用。Reservation 是为降低逐次向上协调成本而量化的额度,不是取消 used 计量,也不等于每次都向系统申请一个量化单位。
8.2 分配流程
8.3 一次 allocate 与 free 的完整路径
字节分配的 aligned buffer 与页分配的 Allocation / ContiguousAllocation 是不同接口:allocate 调用 allocateBytes,allocateNonContiguous 才得到 PageRun 列表,连续页还涉及旧 allocation / collateral 的处理。它们共享预算框架,但不能把字节申请的开头与非连续页分配的结尾拼成一条调用链。
以 MemoryPoolImpl::allocate(size) 的字节分配路径为例:
leaf.allocate(size)
1. 按 pool alignment 对齐
2. reserve(alignedSize)
leaf 已有余量 → 更新 leaf used
余量不足 → 增量向父链传播,先在 root 检查
root capacity 足够 → 提交 reservation
不足 → 进入 arbitration → grow / reclaim / 失败
3. allocator.allocateBytes(alignedSize, alignment)
成功 → 返回 buffer
返回 nullptr → release(alignedSize),报告分配失败
传播增量时先递归到 root 检查,成功后再逐层完成记账。线程安全 leaf 需要在锁外做可能很慢的仲裁,再回锁重查余量;并发 reserve / release 可能造成重试,numCollisions 记录额外尝试。非线程安全 leaf 省掉自身计数锁,但上层 aggregate / root 的同步仍在。
grow(growBytes, reservationBytes) 在 root 锁内同时增加 capacity 并提交本次 reservation。这避免了“先取得额度,再被同一 root 的另一条申请抢走”的窗口。仲裁返回后还会检查是否已被 abort,必要时撤回本次提交并抛错。
free(ptr, size) 的顺序是先 allocator_->freeBytes,再 release(alignedSize)。Leaf 根据剩余用量和最低 reservation 重算量化值,把多余 reservation 逐层归还。Root 的 reservation 下降不意味着它的 capacity 自动下降:后者需要 arbitrator 执行 shrink。
连续页、非连续页和 reallocate 还有各自的回调与旧 allocation 处理;不要把上面的字节路径拼接成一个并不存在的统一调用链。外部分配报告 reportExternalAllocation 只参与记账,不通过这个 pool 的 allocator 获取存储。
源码:字节分配与外部报告、字节释放、原子 grow、仲裁返回检查。
8.3.1 Reservation 成功后,Allocator 仍可能失败
MallocAllocator::allocateBytesWithoutRetry 首先检查自身用量是否越过 allocator limit,然后才调用 malloc / aligned_alloc。配置了 cache 时,MemoryAllocator::allocateBytes 还通过 cache()->makeSpace 释放可驱逐空间并尝试分配。
因此失败可能来自 allocator 容量、缓存中不可驱逐的数据、底层分配失败等,不能仅凭 MEM_ALLOC_ERROR 就认定宿主物理内存耗尽。查询 root capacity、arbitrator 总 capacity、allocator capacity 与进程 RSS 是相关但不同的观察量。
operator->allocate(size)
├── reserve(size) 预先在 pool 树占位
│ └── reserveThreadSafe()
│ ├── 检查 reservationBytes_ 是否有余量
│ └── 不足 → incrementReservationThreadSafe()
│ └── 递归向 parent 传播
│ └── Root Pool: 检查 capacity_
│ ├── 够用 → 记账
│ └── 不够 → growCapacity() → Arbitrator
│
└── MemoryAllocator::allocate() 真正的物理内存分配
8.4 释放流程
operator->free(ptr)
└── release(size)
└── 重新计算仍需保留的 reservation 量
└── 多余部分 decrementReservation() 递归向 parent 传播
8.5 Reservation 量化策略
为减少跨层加锁频率,reservation 按以下粒度向上取整:
流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。
// velox/common/memory/MemoryPool.h:518
static uint64_t quantizedSize(uint64_t size) {
if (size < 16 * kMB) return roundUp(size, 1 * kMB); // 1MB 粒度
if (size < 64 * kMB) return roundUp(size, 4 * kMB); // 4MB 粒度
return roundUp(size, 8 * kMB); // 8MB 粒度
}
8.6 关键统计字段(MemoryPoolImpl)
| 字段 | 含义 | 有效层级 |
|---|---|---|
reservationBytes_ |
当前持有的预留字节(含量化膨胀) | 所有 |
usedReservationBytes_ |
实际分配出去的字节 | Leaf |
minReservationBytes_ |
强制保持的最小预留量 | Leaf |
peakBytes_ |
历史峰值 | 所有 |
cumulativeBytes_ |
累计分配总量 | 所有 |
不变量:reservationBytes_ >= usedReservationBytes_ >= 0
Aggregate pool 的 usedBytes() 通过遍历所有子孙的 usedReservationBytes_ 求和得到。
8.7 两条正交路径:capacity 配额 vs 物理内存
8.7.1 单一检查点设计
整条 pool 树上只有 Root Pool 一处做 capacity 检查,中间层(Operator/Node/Task)只负责"记账+传播":
当前源码摘录(1d1b76567870):velox/common/memory/MemoryPool.cpp:1035。省略外围声明;此片段未作为独立程序编译。
bool MemoryPoolImpl::maybeIncrementReservation(uint64_t size) {
std::lock_guard<std::mutex> l(mutex_);
if (isRoot()) {
checkIfAborted();
// NOTE: we allow memory pool to overuse its memory during the memory
// arbitration process. The memory arbitration process itself needs to
// ensure the memory pool usage of the memory pool is within the capacity
// limit after the arbitration operation completes.
if (FOLLY_UNLIKELY(
(reservationBytes_ + size > capacity_) &&
!underMemoryArbitration())) {
return false;
}
}
incrementReservationLocked(size);
return true;
}
| 字段 | Operator Pool | Node Pool | Task Pool | Root Pool |
|---|---|---|---|---|
capacity_ |
— | — | — | ✓(动态可变) |
maxCapacity_ |
— | — | — | ✓(创建时定,不变) |
reservationBytes_ |
✓(记账) | ✓(记账) | ✓(记账) | ✓(参与检查) |
中间 aggregate pool 汇总子树 reservation,支持向上额度检查、回收候选估算与诊断。硬 capacity 约束集中在 root,并不代表中间层统计完全不参与策略;reclaimer 选择子树时会使用相应资源估算。
8.7.2 配额与物理内存的正交性
Root capacity 控制逻辑额度,allocator 提供存储并实施自身容量限制;两条职责可分开分析,但 MemoryManager 会协调默认容量配置,cache/回收也可能参与分配路径。“配额通过”不保证 allocator 成功,改变 capacity 也不会直接 mmap 或释放物理页。
pool->allocate(size)
│
├─① reserveThreadSafe() capacity 配额管理 → SharedArbitrator(纯逻辑记账)
│ 更新 reservationBytes_ / capacity_
│
└─② MemoryAllocator::allocate() 物理内存分配 → malloc / mmap(存储申请;不等于页立即驻留)
两步均成功才算分配完成,失败原因相互独立:
reserve失败 → capacity 配额不足,触发仲裁- MemoryAllocator 分配失败可能来自 allocator 自身容量、cache 无法腾出足够空间、底层分配错误或地址空间等条件;不等价于整台机器 RAM 已耗尽,也不能无依据说这种失败极少见。Reservation 已成功时仍需正确处理 allocator 失败并回滚记账。
完整正常分配路径(capacity 充足时,不触发仲裁):
pool->allocate(size)
→ reserveThreadSafe(size)
→ 按 leaf 当前用量、最低预留与量化策略计算 reservation 增量
→ incrementReservationThreadSafe(root, quantized)
// 沿父链先检查 root,再逐层提交;仲裁在 leaf 计数锁外执行
→ root->maybeIncrementReservation(quantized)
✓ reservationBytes_ + quantized <= capacity_
→ reservationBytes_ += quantized // 记账,无其他副作用
→ return true // 不涉及仲裁器
→ MemoryAllocator::allocateBytes(alignedSize, alignment)
→ allocator 申请存储并维护页分配状态
→ 成功返回 buffer;nullptr 时 release(alignedSize) 并报错
MemoryAllocator 不直接执行每棵 query root 的 capacity 策略,但有自己的容量限制。Pool 先申请 reservation,必要时进入仲裁,再调用 allocator;两层分别可能失败,MemoryManager 还会协调它们的配置。因而“没有拿到 query 额度”与“实际存储分配失败”应分别排查。
9. Capacity 检查与内存仲裁
9.1 Capacity 检查
只有 Root Pool 做硬性 capacity 检查:
当前源码摘录(1d1b76567870):velox/common/memory/MemoryPool.cpp:1035。省略外围声明;此片段未作为独立程序编译。
bool MemoryPoolImpl::maybeIncrementReservation(uint64_t size) {
std::lock_guard<std::mutex> l(mutex_);
if (isRoot()) {
checkIfAborted();
// NOTE: we allow memory pool to overuse its memory during the memory
// arbitration process. The memory arbitration process itself needs to
// ensure the memory pool usage of the memory pool is within the capacity
// limit after the arbitration operation completes.
if (FOLLY_UNLIKELY(
(reservationBytes_ + size > capacity_) &&
!underMemoryArbitration())) {
return false;
}
}
incrementReservationLocked(size);
return true;
}
9.2 仲裁流程
当前源码摘录(1d1b76567870):velox/common/memory/MemoryPool.cpp:1010。省略外围声明;此片段未作为独立程序编译。
void MemoryPoolImpl::growCapacity(MemoryPool* requestor, uint64_t size) {
VELOX_CHECK(requestor->isLeaf());
++numCapacityGrowths_;
try {
MemoryPoolArbitrationSection arbitrationSection(requestor);
arbitrator_->growCapacity(this, size);
} catch (const VeloxRuntimeError& veloxError) {
if (FOLLY_UNLIKELY(
debugEnabled() &&
veloxError.errorCode() == error_code::kMemCapExceeded)) {
std::rethrow_exception(wrapExceptionDbg(veloxError));
}
throw;
}
// The memory pool might have been aborted during the time it leaves the
// arbitration no matter the arbitration succeed or not.
if (FOLLY_UNLIKELY(aborted())) {
// Release the reservation committed by the memory arbitration on success.
decrementReservation(size);
VELOX_CHECK_NOT_NULL(abortError());
std::rethrow_exception(abortError());
}
}
仲裁器可以:
- 触发其他 query/task 的 spill 释放内存
- 收缩其他 pool 的 capacity
- 找不到足够内存时抛出 OOM,终止 query
10. 内存仲裁机制(SharedArbitrator)
10.1 整体架构
MemoryPool (leaf) 申请内存
│ maybeIncrementReservation() 失败
▼
MemoryPoolImpl::growCapacity()
│
▼
SharedArbitrator::growCapacity(pool, size) ← 进程全局唯一仲裁器
│
├── ArbitrationOperation 一次仲裁请求,状态机 kInit→kWaiting→kRunning→kFinished
├── ArbitrationParticipant 每个 root pool 在仲裁器侧的代理
└── MemoryReclaimer 挂在各层 pool 上,执行实际内存回收
触发入口(MemoryPool.cpp:~1128):
当前源码摘录(1d1b76567870):velox/common/memory/MemoryPool.cpp:1035。省略外围声明;此片段未作为独立程序编译。
bool MemoryPoolImpl::maybeIncrementReservation(uint64_t size) {
std::lock_guard<std::mutex> l(mutex_);
if (isRoot()) {
checkIfAborted();
// NOTE: we allow memory pool to overuse its memory during the memory
// arbitration process. The memory arbitration process itself needs to
// ensure the memory pool usage of the memory pool is within the capacity
// limit after the arbitration operation completes.
if (FOLLY_UNLIKELY(
(reservationBytes_ + size > capacity_) &&
!underMemoryArbitration())) {
return false;
}
}
incrementReservationLocked(size);
return true;
}
!underMemoryArbitration() 是防重入关键:仲裁线程自身申请内存时允许临时超限,不会触发嵌套仲裁。
10.2 ArbitrationParticipant 与 ArbitrationOperation
ArbitrationParticipant 是 root pool 在仲裁器侧的代理,每个 root pool 注册时创建一个,拥有单调递增的 id(用于 abort 时选"最年轻"的牺牲者)。
同一 participant 的 operation 串行执行(ArbitrationParticipant.cpp:217):
当前源码摘录(1d1b76567870):velox/common/memory/ArbitrationParticipant.cpp:217。省略外围声明;此片段未作为独立程序编译。
void ArbitrationParticipant::startArbitration(ArbitrationOperation* op) {
ContinueFuture waitPromise{ContinueFuture::makeEmpty()};
{
std::lock_guard<std::mutex> l(stateLock_);
++numRequests_;
if (runningOp_ != nullptr) {
op->setState(ArbitrationOperation::State::kWaiting);
WaitOp waitOp{
op,
ContinuePromise{fmt::format(
"Wait for arbitration on {}", op->participant()->name())}};
waitPromise = waitOp.waitPromise.getSemiFuture();
waitOps_.emplace_back(std::move(waitOp));
} else {
runningOp_ = op;
}
}
if (waitPromise.valid()) {
waitPromise.wait();
}
}
10.3 Capacity 体系:从进程到 pool 的三层 quota
10.3.1 三层 quota 概览
整个 Velox 内存配额体系是一个三层嵌套约束:
| 量 | 字段 | 含义 | Presto 配置 | 典型值 |
|---|---|---|---|---|
| 进程总上限 | SharedArbitrator::capacity_ |
所有 query 持有的 capacity 之和不能超过 | system-memory-gb |
80GB |
| 保留池 | SharedArbitrator::freeReservedCapacity_ |
为 participant 的最低容量需求保留的全局额度;不是按“低优先级/饥饿”自动授予的专属池。 | query-reserved-memory-gb |
4GB |
| 单 query 上限 | MemoryPool::maxCapacity_ |
query 永远不能超过的 capacity(创建时定,不变) | query.max-memory-per-node |
30GB |
| 单 query 当前 | MemoryPool::capacity_ |
该 query 目前持有的动态 capacity | — | 运行时变化 |
| 初始 capacity | initMemoryPoolCapacity_ |
query 创建时同步划拨的起步配额 | memory-pool-initial-capacity |
256MB(当前 Velox 默认目标,实际授予受空闲量与策略限制) |
三个限额的角色分工:
| 限额 | 目的 | 生效层 |
|---|---|---|
arbitrator.capacity_ |
进程稳定性:所有 query 总和不能撑爆 worker | 最外层 hard line |
maxCapacity_ |
公平隔离:单个 query 不能独吞整机 | 每个 root pool 自己的上限 |
initCapacity |
性能优化:新 query 不必冷启动就触发仲裁 | 仅创建时一次性使用 |
10.3.2 Arbitrator.capacity_ 的来源:MemoryManager 初始化
Arbitrator 的总配额只在 MemoryManager 构造时设定一次,运行期不变:
流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。
// velox/common/memory/Memory.cpp
MemoryManager::MemoryManager(const Options& options)
: capacity_(options.allocatorCapacity), ... {
arbitrator_ = MemoryArbitrator::create({
.kind = "SHARED",
.capacity = options.arbitratorCapacity, // ← 进程级配额
.arbitratorReservedCapacity = options.arbitratorReservedCapacity,
...
});
}
在 Presto Native 中,这些值在 worker 启动时从配置文件读入:
// presto_cpp/main/PrestoServer.cpp
memory::MemoryManagerOptions options;
options.allocatorCapacity = systemConfig->systemMemoryGb() * GB;
options.arbitratorCapacity = options.allocatorCapacity;
options.arbitratorReservedCapacity = systemConfig->queryReservedMemoryGb() * GB;
memory::MemoryManager::initialize(options);
关键性质:arbitrator.capacity_ 只是一个"逻辑配额数字",不预占任何物理内存。存储仍由所选 allocator 通过 malloc / mmap 等路径按需申请,arbitrator 只负责账面上不超量。
10.4 Pool capacity 的赋值:addPool 初始分配 vs 运行时扩容
初始授予和运行期增长都通过仲裁器管理的 free capacity 记账,并结合普通空闲容量、保留容量与 participant minimum 策略;不能把全部来源简化成只从 freeNonReservedCapacity_ 扣减。
10.4.1 创建时:addPool 划拨初始 capacity
addPool 在注册 root 时同步计算初始授予量,受初始目标、root 上限、当前空闲量及保留策略约束。相关配置描述的是额度,不会在此一步预先分配等量物理内存。
当前源码摘录(1d1b76567870):velox/common/memory/SharedArbitrator.cpp:452。省略外围声明;此片段未作为独立程序编译。
void SharedArbitrator::addPool(const std::shared_ptr<MemoryPool>& pool) {
checkRunning();
VELOX_CHECK_EQ(pool->capacity(), 0);
auto newParticipant = ArbitrationParticipant::create(
nextParticipantId_++, pool, &participantConfig_);
{
std::unique_lock guard{participantLock_};
VELOX_CHECK_EQ(
participants_.count(pool->name()),
0,
"Memory pool {} already exists",
pool->name());
participants_.emplace(newParticipant->name(), newParticipant);
}
auto scopedParticipant = newParticipant->lock().value();
std::vector<ContinuePromise> arbitrationWaiters;
{
std::lock_guard<std::mutex> l(stateMutex_);
const uint64_t minBytesToReserve = std::min(
scopedParticipant->maxCapacity(), scopedParticipant->minCapacity());
const uint64_t maxBytesToReserve = std::max(
minBytesToReserve,
std::min(
scopedParticipant->maxCapacity(), participantConfig_.initCapacity));
const uint64_t allocatedBytes = allocateCapacityLocked(
scopedParticipant->id(), 0, maxBytesToReserve, minBytesToReserve);
if (allocatedBytes > 0) {
VELOX_CHECK_LE(allocatedBytes, maxBytesToReserve);
try {
checkedGrow(scopedParticipant, allocatedBytes, 0);
} catch (const VeloxRuntimeError& e) {
VELOX_MEM_LOG(ERROR)
<< "Failed to allocate initial capacity "
<< succinctBytes(allocatedBytes)
<< " for memory pool: " << scopedParticipant->name() << "\n"
<< e.what();
freeCapacityLocked(allocatedBytes, arbitrationWaiters);
}
}
}
for (auto& waiter : arbitrationWaiters) {
waiter.setValue();
}
}
只要系统空闲容量足够,这一步同步完成,无需排队。若系统空闲容量不足 initMemoryPoolCapacity_,则按实际可用量给。
10.4.2 运行时:capacity 不足时通过仲裁快速扩容
当 maybeIncrementReservation() 返回 false,触发 growCapacity,进入本地仲裁。Step 1(maybeGrowFromSelf)是非 spill 快速路径:
root->growCapacity(requestor, size)
→ MemoryPoolArbitrationSection ← driver 进入 suspended(enterSuspended)
→ ScopedArbitration ← thread-local = kLocal
→ SharedArbitrator::growCapacity(pool, size)
→ Step 1: maybeGrowFromSelf(op)
→ allocateCapacityLocked(minGrow, maxGrow)
if (freeNonReservedCapacity_ >= minGrow):
freeNonReservedCapacity_ -= allocated
checkedGrow(participant, allocated, reservationBytes)
→ pool->capacity_ += allocated ← 原子增长
→ pool->reservationBytes_ += needed
return true ✓ 无需 spill,直接成功
即使是这种"不需要 spill"的快速扩容,仍然要经过:
MemoryPoolArbitrationSection(driver 短暂进入 suspended 状态)- thread-local kLocal 仲裁上下文标记
SharedArbitrator::freeNonReservedCapacity_扣减
与"完整仲裁"的区别仅在于:Step 1 就成功返回,不需要 Step 2-5(不涉及 shrink/spill/abort 其他 pool)。
10.4.3 capacity 增长策略(非线性)
allocateCapacityLocked 并不总是按请求量分配,而是按增长策略计算目标量,以减少未来频繁扩容:
流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。
// SharedArbitrator.cpp: maxGrowBytes
uint64_t maxGrowBytes = participant->maxGrowBytes(requestBytes);
// ① pool->capacity_ < fastExponentialGrowthCapacityLimit (512MB) → 翻倍
// ② pool->capacity_ >= 512MB → cap * slowCapacityGrowPct (默认 25%)
// 最终取 [minGrowBytes, maxGrowBytes] 区间内能分配的最大量
三种路径对比汇总:
| 路径 | 触发条件 | capacity 操作 | driver 状态 | 耗时级别 |
|---|---|---|---|---|
| 正常路径 | leaf 本地预留足够,或增量传播后 root reservation 不超过 capacity | 无(仅记账) | RUNNING | 依赖路径、锁竞争与构建;此处未测量延迟 |
| 快速扩容(Step 1) | capacity_ 不足,arbFree 足够 | 从 freeNonReservedCapacity_ 划拨 | 短暂 SUSPENDED | 微秒级 |
| 完整仲裁(Step 2-5) | arbFree 不足,需回收其他 pool | spill/shrink/abort 其他 pool | SUSPENDED / WAITING_ARB | 毫秒~秒级 |
10.4.4 addPool 不够分怎么办
freeNonReservedCapacity_ 不足时,addPool 按现有量给而非失败:
流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。
// addPool 的额度决策摘要,不是只从普通空闲池取 min。
minBytesToReserve = min(root.maxCapacity, participant.minCapacity)
maxBytesToReserve = max(minBytesToReserve,
min(root.maxCapacity, config.initCapacity))
allocatedBytes = allocateCapacityLocked(id, 0, maxBytesToReserve,
minBytesToReserve)
if allocatedBytes > 0:
checkedGrow(participant, allocatedBytes, reservationBytes=0)
// grow 抛错时,将本次 allocatedBytes 归还全局空闲额度。
// allocateCapacityLocked 同时遵守普通空闲与 reserved capacity 策略。
初始 capacity 为零时,首次需要增长 reservation 的分配会进入容量获取路径;它仍可能在 self/free-capacity 快速阶段成功,并不必然执行完整的全局回收步骤。初始授予是性能和容量策略,不能用它是否为零判断查询一定失败。
10.4.5 全部 capacity 耗尽时的兜底链
当 freeNonReservedCapacity_ 与 freeReservedCapacity_ 都耗尽:
growCapacity(op)
Step 1: maybeGrowFromSelf → root 自身空闲不足,FAIL(不读取 arbFree)
Step 2: ensureCapacity → 检查是否超 maxCapacity
Step 3: growWithFreeCapacity → arbFree 仍不够, FAIL
Step 4: reclaimUnusedCapacity → 扫所有 participants
把它们持有但没用的 cap(capacity_ - reserved)shrink 回来
不涉及 spill,纯账面操作
Step 5: 全局仲裁 → 后台线程 spill 其他 query 的"已用但可 spill"内存
超时 50% 后转 abort 最年轻的 victim
最终兜底:所有路径都救不了 → 抛 VELOX_MEM_POOL_CAP_EXCEEDED
该 query 被 abort,持有的 cap 全部归还给 arbitrator
仲裁器会在配置和时间预算允许时回收资源、转移 capacity 或 abort participant。但它管理的是所注册 root 的额度,system pool、allocator 自身限制以及不经该 allocator 的进程内存都要另看;这套策略不能保证进程永远不会遇到 OOM。
10.5 本地仲裁流程
10.6 SharedArbitrator:从快速路径到全局等待
每个注册的 root 对应一个 ArbitrationParticipant,每次扩容对应一个 ArbitrationOperation。同一 participant 的 operation 通过 runningOp_ + waitOps_ 串行执行;不同 root 可以并发执行本地路径。
MemoryPoolArbitrationSection 调用请求 leaf 的 enter / leave 回调;ScopedArbitration 安装线程局部仲裁上下文、登记 operation,并在退出时结束 operation。两个 scope 分别负责执行状态桥接与仲裁生命周期。
10.6.1 初始额度与两个空闲容量池
addPool 只尝试分配初始 capacity,不为新 root 直接发起一轮 spill。目标受 memory-pool-initial-capacity、participant minimum 和 root maxCapacity 限制;资源不足时,初始额度可低于目标,甚至为 0。
Arbitrator 有普通空闲容量 freeNonReservedCapacity_ 和保留空闲容量 freeReservedCapacity_。allocateCapacityLocked 优先用普通容量;只有普通容量不足以达到本次 minAllocateBytes 时,才从保留部分补充。它不是“普通池耗尽后任意请求都能继续用”的后备池。
配置 memory-pool-reserved-capacity 表达每个活跃 participant 的最低容量策略,但不会凭空创造资源,也不保证任意数量的查询都拿到该额度;当前 participant 构造还要求这个 minimum 不大于 root maxCapacity。归还容量时先补回全局 reserved 部分,再增加普通空闲容量。
源码:addPool、allocateCapacityLocked、freeCapacityLocked、participant minimum 校验。
10.6.2 扩容按顺序尝试
- **
maybeGrowFromSelf**:执行participant->grow(0, requestBytes)。如果并发释放或前一轮仲裁留下了足够 root 余量,直接提交 reservation;不增加 capacity,也不扣减 arbitrator 的空闲容量。 - **
ensureCapacity**:检查这次增长能否落在 root maxCapacity 和 arbitrator 总 capacity 内。必要时 shrink 自身空闲 capacity、回收自身已用资源,再重查。即使启用了全局仲裁,撞自身上限时仍可能在此 self-reclaim。 - **
growWithFreeCapacity**:计算增长目标,向 arbitrator 取得空闲额度,提交 grow 与 reservation。 - **
reclaimUnusedCapacity**:遍历各 participant,收回可收缩的空闲 capacity,再尝试第 3 步。这不调用算子 spill,也不暂停那些 Task,但会减少它们未来可用的额度。 - 资源仍不足:若关闭 global arbitration,尝试回收请求方自身的已用资源;若开启,进入全局等待。失败、超时或 abort 则沿异常路径返回。
Local 表示请求线程上的仲裁阶段,不能理解为“只访问当前 query”。反过来,Local 中的 self-reclaim 也可能暂停本查询的 Task 和其他 Driver,不是完全无干扰的操作。
源码:growCapacity 主流程、maybeGrowFromSelf、ensureCapacity、收回空闲容量。
10.6.3 增长目标与实际授予量分开
当前增长策略在 2 × currentCapacity <= fastExponentialGrowthCapacityLimit 时以翻倍为目标,否则按 slowCapacityGrowRatio 计算增量;再保证覆盖 request 和 minimum gap,并截断到 root 剩余可增长量。
默认 fast limit 是 512 MiB、slow ratio 是 25%。因此 currentCapacity 为 256 MiB 时满足翻倍条件,384 MiB 时已走比例增长,不能写成“capacity 小于 512 MiB 就翻倍”。目标增长不是承诺授予量;实际还受空闲额度和等待者策略约束。
源码:getGrowTargets。
**SharedArbitrator.cpp:891**,按序尝试,成功即返回:
| 步骤 | 操作 | 说明 |
|---|---|---|
| 1 | maybeGrowFromSelf() |
仲裁器空闲容量够用?直接分配 |
| 2 | ensureCapacity() |
是否超出 pool maxCapacity?shrink 自身→reclaim 自身→再 shrink |
| 3 | growWithFreeCapacity() |
allocateCapacityLocked() 从仲裁器空闲池取容量,checkedGrow() 转给 pool |
| 4 | reclaimUnusedCapacity() |
shrink 所有其他 participants 的空闲容量(不涉及 spill),再重试步骤 3 |
| 5 | 全局仲裁 | 加入 globalArbitrationWaiters_,唤醒后台线程,阻塞等待 |
仲裁器维护两个空闲池:
freeNonReservedCapacity_:普通空闲池freeReservedCapacity_:保留池(per-query 保护容量)
普通空闲与保留容量的使用受 participant minimum 等条件约束;保留池不是任何请求在普通池不足时都可以无条件取用的第二余额。分析授予量要读取 allocateCapacityLocked 的实际分支。
pool capacity 增长策略:
capacity < fastExponentialGrowthCapacityLimit(默认 512MB)→ 翻倍增长- 超过后 → 按
slowCapacityGrowPct(默认 25%)增长
10.6.4 Step 2 ensureCapacity 详解:处理 maxCapacity 边界
maxCapacity_ 是 query 创建时定的硬上限(= query.max-memory-per-node),即使系统有大量空闲内存也不能跨过。Step 2 专门处理"扩容会撞 maxCapacity 上限"的情形:
当前源码摘录(1d1b76567870):velox/common/memory/SharedArbitrator.cpp:1167。省略外围声明;此片段未作为独立程序编译。
bool SharedArbitrator::ensureCapacity(ArbitrationOperation& op) {
if ((op.requestBytes() > capacity_) ||
(op.requestBytes() > op.participant()->maxCapacity())) {
return false;
}
if (checkCapacityGrowth(op)) {
return true;
}
shrink(op.participant(), /*reclaimAll=*/true);
if (checkCapacityGrowth(op)) {
return true;
}
reclaim(
op.participant(),
op.requestBytes(),
op.timeoutNs(),
/*localArbitration=*/true);
// Checks if the requestor has been aborted in reclaim above.
checkIfAborted(op);
if (checkCapacityGrowth(op)) {
return true;
}
shrink(op.participant(), /*reclaimAll=*/true);
return checkCapacityGrowth(op);
}
关键性质:
| 性质 | 解释 |
|---|---|
| maxCapacity 永远是硬约束 | 即使全局仲裁 spill 了 1TB 也救不了——别人 spill 出的容量不能给已到上限的 query |
| Self-reclaim 优先于 global-reclaim | 撞自己上限时,先逼自己 spill;不行才波及其他 query |
| Self-reclaim 仍是本地仲裁 | 全程在 driver 自己的线程同步执行,不需要后台仲裁线程介入 |
| 失败语义不同 | 系统级 OOM = 整机配额耗尽;单 query OOM = 自己撞 maxCap 但系统可能还很空 |
典型场景:query.max-memory-per-node = 10GB,当前 Q1.capacity_ = 9.5GB,Q1.reservationBytes_ = 9.4GB,系统 arbFree = 50GB:
// 容量示例统一使用 GiB;忽略量化与并发变化,只解释边界。
Q1.capacity = 9.5 GiB,reservation = 9.4 GiB,maxCapacity = 10 GiB
本次上传的 reservation 增量 = 0.7 GiB;全局空闲 = 50 GiB
maybeGrowFromSelf: 尝试 grow(0, 0.7 GiB),不是向全局取 1 GiB
当前空闲 0.1 GiB,不足提交 -> 失败
ensureCapacity:
检查 capacity + requestBytes 能否增长,必要时先 shrink 空闲 capacity
若 shrink 后仍不能满足边界,则 self-reclaim,再检查 / shrink
后续 growWithFreeCapacity 才向全局取额度并提交 reservation
目标语义:旧 reservation + 新 increment 不能靠突破 maxCapacity 解决
实际 self-reclaim 目标含 minReclaimBytes / minReclaimPct,
算子回收还可能超过目标,因此不能推导恰好只 spill 300 MiB。
全局空闲很多,仍不能越过该 root 的 maxCapacity。示例说明 self-reclaim 的必要性;回收量不是由两个数相减就精确决定,还受最低回收目标、算子回收粒度和并发变化影响。
10.7 全局仲裁(后台线程)
本地仲裁步骤 4 仍不够时,进入全局仲裁。
后台线程主循环(SharedArbitrator.cpp:1045):
流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。
void SharedArbitrator::globalArbitrationMain() {
while (true) {
wait(globalArbitrationWaiters_ 非空 或 shutdown);
GlobalArbitrationSection section{this};
runGlobalArbitration();
}
}
runGlobalArbitration() 每轮逻辑:
计算 targetBytes = max(所有 waiter maxGrowBytes 之和, capacity × reclaimPct)
if !shouldReclaimByAbort:
reclaimUsedMemoryBySpill() ← 优先 spill 到磁盘
else:
reclaimUsedMemoryByAbort() ← abort 查询
if totalReclaimed >= target: break,否则进入下一轮
Spill → Abort 切换条件(SharedArbitrator.cpp:1064):
bool shouldReclaimByAbort =
globalArbitrationWithoutSpill_ || // 配置强制只 abort
(elapsedTime > maxArbitrationTime * 0.5 && // 超过 50% 超时
(hasReclaimedByAbort ||
(allParticipantsReclaimed && lastReclaimedBytes == 0))); // spill 已无进展
Spill 牺牲者选择(优先级依次):
- 容量分桶(
spillCapacityLimits,默认起始 4GB):优先动大用量的 - 优先级:低优先级(高数值)先 spill
- 可回收量:同优先级内可回收量大的先动
- 跳过可回收量
< minReclaimBytes(默认 128MB)的参与者
Abort 牺牲者选择(优先级依次):
- 优先级:低优先级先 abort
- 容量分桶(
abortCapacityLimits,默认起始 1GB) - participant id:越大(越年轻)越先 abort,保护老查询
对多个 victim 同时提交到 memoryReclaimExecutor_ 线程池并行执行。
10.7.1 Arb 线程与 Rcl 线程池的分工
全局仲裁涉及两类性质完全不同的线程,职责严格分离:
| 维度 | Arb 线程 | Reclaim 线程池 |
|---|---|---|
| 数量 | 1(常驻 globalArbitrationMain) |
N(动态 worker,由 memoryReclaimExecutor_ 管理) |
| 角色 | 决策中心 | 执行单元 |
| 主要职责 | 计算 target、选 victim、分配 cap、唤醒 waiter | requestPause + reclaim + shrink + resume |
| 并发度 | 必须串行(全局视图) | 不同 participant 可提交到 reclaim executor 并行处理;线程池容量、候选数及控制循环也可能让实际工作串行。并行是实现能力,不是每轮必然发生的条件。 |
| 操作性质 | 决策、状态机驱动 | I/O 密集(磁盘 spill) |
| 持有锁时长 | 依具体锁域与工作量变化,无固定微秒保证 | reclaimMutex 覆盖相应回收;含等待及可能的 I/O,无固定秒数 |
| 阻塞类型 | cv.wait(轻量) |
pauseFuture.wait + 磁盘 IO(重量) |
为什么必须分离:
| 原因 | 解释 |
|---|---|
| 决策必须串行 | 多 Arb 线程读不一致的全局状态会做出冲突决策(超分、重复 reclaim、公平性破坏) |
| 执行必须并行 | 线程池允许多个 participant 的回收重叠,控制线程仍会等待本轮结果;并行度受 executor、I/O 与峰值内存限制。 |
| Backpressure 隔离 | 独立 executor 允许多个回收任务重叠,但慢 I/O 或长 pause 仍会影响当前回收轮完成及 waiter 延迟。不能把分工推导成决策线程对执行耗时完全免疫。 |
| 可观测性 | 决策点集中(1 线程)便于调试;执行点分散(N worker)便于追踪 spill |
10.7.2 三层并行结构
Layer 1 — Arb 线程 × 1 全局决策(串行)
│ ├ 计算 target
│ ├ 排序选 victim
│ └ 分配 cap、notify waiter
▼ 提交多 victim
Layer 2 — memoryReclaimExecutor_ × N 多 victim 并行
│ 一个 worker 负责一个 participant 全周期
│ ├ 持 reclaimMutex_
│ ├ pause + reclaim + resume + shrink
▼ 进入单 task 内部
Layer 3 — spillExecutor × M 单 task 内多 driver 并行 spill
(ParallelMemoryReclaimer 调度)
各层有不同的执行与等待职责:
- 全局控制线程会组织本轮回收并取得结果,可能等待 spill I/O 完成;它与 executor 的分工允许任务并行,但不是 fire-and-forget。
- 不同 participant 的回收可以并发,各自受 reclaim 锁与 Task 状态约束,同时仍共享 executor、allocator、CPU 和存储资源,不能保证彼此完全不影响。
- Layer 3 单 query 内多 driver 并发 spill — 同 join 的 build/probe 多 driver 同时写盘
10.7.3 Victim 粒度:query(= root pool = participant)
sortAndGroup(participants_) 显式按 participant 排序——victim 单位就是 root pool,对应一个 QueryCtx,即一个 query:
SharedArbitrator
└── participants_ : Map<id, ArbitrationParticipant>
│
└── 每个 participant 包一个 root pool (weak_ptr)
↑ root pool 一一对应 QueryCtx
↑ QueryCtx 一一对应一个 query
排序依据全部是 query 级指标:
reclaimableUsedCapacity— 该 query 当前可 spill 总量capacity— 该 query 当前持有 cap(用于容量桶)priority— query 优先级(低优先级先被 victim)id— query 单调 id(abort 时选最年轻保护老查询)
为什么以 query 为粒度:
capacity_只在 root pool 有意义(task / node / op 无独立 cap)- 公平性/隔离边界 =
query.max-memory-per-node - abort 的失败语义 = query 整体失败(不是单 task 死)
10.7.4 Rcl worker 完整职责:pause→reclaim→resume→shrink
一个 rcl worker 负责一个 participant 的整个 reclaim 周期,期间持有 reclaimMutex_ 不让出:
rcl_worker (从 executor 接到任务):
① 恢复 thread-local context = kGlobal
② ArbitrationTimedLock(reclaimMutex_, 0) ← 全程持有
③ participant->reclaim()
└─ pool_->reclaim() 沿 Reclaimer 链下沉
└─ TaskReclaimer 对某 task:
├─ task->requestPause() ← 设 pauseRequested_
├─ pauseFuture.wait() ← BLOCK 等 numThreads_=0
├─ ParallelMemoryReclaimer 并发 spill 多 driver
└─ task->resume() ← 清 pause 并重新调度符合条件的 Driver
④ participant->shrink(false) ← 归还 cap 给 arbitrator
⑤ reclaimMutex_ 释放,任务函数返回
关键性质:
- rcl worker 在整个 reclaim 周期内持续持锁,避免并发 reclaim/abort
- pause 是 rcl worker 发起;BLOCK 等的也是 rcl worker;resume 也是 rcl worker 触发
- victim drivers 只是被动响应
pauseRequested_,不参与决策也不执行 spill
详细的 spill 执行机制和算子 spill 状态机参见 §10.13。
10.7.5 同 victim query 内多 task 的串行 reclaim 顺序
同一 query 在一个 worker 上可能有多个 task,例如不同 stage 的任务;grouped execution 则主要是在一个 Task 内组织多个 split group,不能把多 task 与多 group 等同。Root reclaimer 再依据其子树策略遍历 Task。
// memory::MemoryReclaimer 的串行 Task 子树回收摘要。
在 poolMutex_ 保护下取得 child 强引用和估算;然后释放该锁
候选排序:priority 小者在前,同 priority 按 reclaimableBytes 大者在前
for candidate in sorted candidates:
bytes = candidate.pool->reclaim(remainingTarget, maxWaitMs, stats)
accumulated += bytes
若原始 target 为 0:继续处理全部候选
否则 bytes 已满足剩余目标:停止;否则更新 remainingTarget
// 估算与实际回收不同;每个 Task 的暂停 / 恢复由其 reclaimer 负责。
排序与早停规则:
| 规则 | 含义 |
|---|---|
| 贪心最大 | 按 priority 升序、同优先级可回收量降序处理候选,达到目标后早停;实际回收量可偏离估算,这不保证全局最少的 pause 次数。 |
| 早停 | 累计 ≥ target 立即停,不处理剩余 task |
| 不并行 pause 多 task | 默认树内回收依次访问 Task 子树,有助于达到目标后早停。并行是否合适取决于资源共享和策略,task 级标志本身并不禁止不同 Task 并行 pause。 |
| task 内可并行 | 单 task 内多 driver 通过 spillExecutor 并发 spill(Layer 3) |
示例:Q1 在该 worker 上有 3 个 task,target=1024MB,reclaimable=[T2:800, T1:500, T3:200]MB
轮 1: pause T2 → spill 800MB → resume T2 (累计 800, 差 224)
轮 2: pause T1 → spill 224MB → resume T1 (累计 1024 ✓ 退出)
T3 完全不动
少 pause、早收工、影响面最小化是核心目标。
同一 worker 可承载同一 query 的多个 stage/task,不能把一个 root 对应一个 task 当作普遍近似。Root 下的多个 Task 也不要求启用 grouped execution;两者是不同的并行维度。
10.8 Global Arbitration 线程交互时序
10.9 全局仲裁:哪些线程真的在等
10.9.1 请求方保留调用栈
startAndWaitGlobalArbitration 在 stateMutex_ 内最后尝试一次分配;仍不足才登记 waiter,记录 pending growth,通知控制线程,并在 future 上限时等待。
这是同步分配接口内部的等待。若从正常 Operator 的 Driver 线程发起,enterArbitration 已让它进入 Suspended:线程与 C++ 栈仍在,只是从 Task 的活跃线程计数中暂时扣除。这里不是 runInternal 返回 StopReason::kBlock 后把线程还给 executor。非 Driver 线程也可能发起仲裁,此时没有 Driver suspension。
10.9.2 控制线程会执行工作,也会等待结果
globalArbitrationMain 在 condition variable 上等 waiter,随后运行 runGlobalArbitration。每轮根据待满足的增长目标选出多个 victim。
当前 reclaimUsedMemoryBySpill 为 victim 创建 AsyncSource:从第二个任务起提交给 memoryReclaimExecutor_,随后逐个调用 move() 收集结果。第一个任务没有提前提交,会由 move() 在调用线程执行;其余未被 worker 领取的任务也可能由调用方执行,已开始的则需要等待。
所以 ArbitrationMain 并不是“派发后立即服务下一请求、永不等待 I/O”的线程。它参与回收并等待这一批任务结束;回收 worker 在执行期间仍可通过 freeCapacity 给 waiter 分配额度并唤醒请求方,无需等控制线程开始下一轮。
Node 层的 ParallelMemoryReclaimer 使用类似模式,executor 来自 QueryCtx 的 spill executor;没有 executor 时退回基类的逐个子池回收。不能把 reclaim worker 或控制线程都画成只负责派发、不执行任何 spill。
源码:全局登记与等待、控制循环、victim 任务执行、AsyncSource::move、ParallelMemoryReclaimer。
10.9.3 Waiter 顺序与 victim 顺序是两套策略
全局 waiter 按 participant ID 排序;ID 是 root 注册时分配的顺序,不是本次请求的到达时间。allocateCapacityLocked 会限制较年轻 participant 越过老 waiter 获取额度,但 minimum growth 有专门处理。恢复时依次尝试,遇到无法满足的 waiter 就停止本轮。
Spill victim 的选择顺序是:先按可回收已用容量分桶,再在桶内按 reclaimer priority 数值降序、可回收量降序排序;数值较大表示较低的查询优先级,优先承担回收。低于 minReclaimBytes 的候选会跳过。
Abort 的顺序不同:先分查询优先级组,再按 capacity 桶查找,同桶偏向 participant ID 更大、较年轻的查询。判断 capacity 时还计入它正在向全局请求的增长量,避免“小 pool 申请大额度”逃过选择。
每轮全局回收目标为待处理 operation 的 maxGrowBytes 总和与总容量某个比例的较大值;目标不等于最终实际回收字节。默认全局比例是 10%。
切换到 abort 也不只是“过一半 timeout 就 kill”:除 global-arbitration-without-spill 直接启用 abort 外,还要满足时间阈值,以及已进入 abort 阶段、或候选已遍历且上一轮没有回收进展等条件。
10.10 Local vs Global:作用域与线程模型
10.10.1 作用域辨析:local 不等于"只动自己"
| Step | 阶段 | 触及对象 | 是否动用别人 |
|---|---|---|---|
1 maybeGrowFromSelf |
Local | arbitrator 空闲池 | ✗ |
2 ensureCapacity |
Local | 自己(撞 maxCap 时 self-spill) | ✗ |
3 growWithFreeCapacity |
Local | arbitrator 空闲池(同 Step 1,重试) | ✗ |
4 reclaimUnusedCapacity |
Local | **所有 query 的"未用 capacity"**(账面 shrink) | ⚠️ 账面,不 spill |
| 5 进入全局等待 | Global | **victim task 的"已用内存"**(真 spill) | ✓ |
收回其他 root 未使用的 capacity 通常不需要释放其已用数据或暂停 Driver,但会改变该 root 后续增长的可用额度,下一次分配可能更早进入仲裁。因此“只收账面空闲”不等于对后续执行完全无影响。
Local 是请求线程上的容量获取与回收路径,可能收缩其他 root 空闲额度,也可能自身回收;全局路径使用 waiter 与后台控制循环协调更广范围的回收。是否暂停某个 Task 取决于实际进入的 reclaim,而不是单凭 local/global 标签。
10.10.2 线程角色对比
仲裁涉及四种线程角色:
| 角色 | 来源 | 在 Local 中做什么 | 在 Global 中做什么 |
|---|---|---|---|
| Requestor Driver | 触发分配的 driver 线程 | 同步执行 Step 1-4 | 提交 Step 5 后 BLOCK 在 future.wait() |
| Background Arb Thread | globalArbitrationMain() 常驻线程 |
不参与 | 选 victim、分配 cap、唤醒 waiter |
| Reclaim Thread Pool | memoryReclaimExecutor_ 工作线程池 |
不参与 | requestPause + reclaim + shrink |
| Victim's Drivers | victim task 的 N 个 driver 线程 | 不存在 victim 概念 | 自己检测 pauseRequested_ → 协作式进入 PAUSED |
| 维度 | Local | Global |
|---|---|---|
| 决策线程 | Requestor driver 自己 | 后台 arb 线程 |
| 决策时机 | 同步、即刻 | 异步、可能等待 |
| 执行 spill 线程 | Requestor driver 自己(仅 Step 2 self-spill) | Reclaim 线程池 |
| 受影响"别人"的范围 | 别人的空闲 cap(账面) | 别人的已用内存(真 spill) |
| Pause 谁 | 普通 free/shrink 路径不 pause;self-reclaim 可暂停相关 Task | 对选中 participant 下实际需要回收的 Task 请求停稳 |
| Requestor driver 状态 | 仲裁期间 suspended,保留调用栈 | 保留调用栈,同步等待仲裁 future;不能执行另一 Driver |
10.10.3 Pause 粒度:Task,不是 Driver
全局仲裁的 pause 是 整个 victim task 粒度,不能是单 driver:
- HashJoin build table 是同一 join node 下多个 driver 共享的
- 如果只 pause 1 个 driver,其余还在跑,它们可能在写共享数据结构 → spill 时数据不一致
task->requestPause()是 task 级,victim 的 N 个 driver 在各自runInternal()循环里协作式自暂停,等numThreads_ == 0时 reclaim 线程才动手
10.10.4 三种让出机制:SUSPENDED / PAUSED / BLOCKED on future
仲裁系统其实是三种语义完全不同的"让出 CPU"机制协作完成的,每种各管一个场景:
| 机制 | 主语 | 范围 | 触发方 | 唤醒方 | 物理 CPU |
|---|---|---|---|---|---|
| SUSPENDED | 单个 driver | 自己一个 | driver 自己 (enterSuspended) |
driver 自己 (leaveSuspended) |
仍占用(账面让出) |
| PAUSED | 单个 task 的所有 driver | task 内 N 个 driver | 外部 (task->requestPause()) |
外部 (task->resume()) |
普通 Driver 返回 kPause,后续由 Task::resume 重新入队;不是所有 Driver 都栈内 wait |
| 仲裁中的同步 future wait | 单个 driver | 自己一个 | requestor 在仲裁栈内 waitFuture.wait() | 外部 (promise.setValue()) |
保留线程与栈;和普通 Driver kBlock 归还 executor 不同 |
三者实现机制:
流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。
// 流程示意,省略错误检查;三种等待保留不同的执行状态。
// ① 仲裁 / 同步调用需要保留栈:
// Task::enterSuspended(state) 在锁内增加 state.numSuspensions。
// 只有从 0 到 1 的首次进入才 --numThreads_;嵌套进入不重复扣数。
// 仲裁代码继续在当前线程运行,也可以同步等待 op.waitFuture()。
// Task::leaveSuspended(state) 在最后一层退出时恢复 numThreads_。
// 若 Task 仍请求 pause,当前实现在锁外 sleep 10ms 再检查,不持锁等待。
// ② 普通 Driver 响应 Task pause:
// requestPause() 设置 pauseRequested_,返回“全部活跃线程停稳”的 future。
// Driver::runInternal() 在安全检查点返回 kPause,经 leave 退出活跃计数。
// 该 Driver 退出执行器任务;Task::resume() 再将符合条件的 Driver 入队。
// 这里不在 Driver 栈内调用 pauseFuture.wait()。
// ③ 普通 Operator future block:
// isBlocked(&future) -> blockDriver() -> 返回 kBlock -> 离开 Driver
// BlockingState::setResume() 安装 continuation;future ready 后重新入队。
// 这是归还 executor 线程的路径,和①的同步 wait 不同。
关键不变量:Requestor driver 在整个仲裁期间始终 SUSPENDED(numThreads_ 不计自己)。这样如果 requestor 和 victim 是同一个 task(自己被 global 选为 victim),requestPause 等 numThreads_ == 0 时不会因为 requestor 还在"运行"而死锁。
10.10.5 enterArbitration / leaveArbitration / enterSuspended / leaveSuspended 的分层设计
这四个 API 体现了一套"RAII + 桥接 + 多态分发 + 状态复用"的优雅设计,把 pool / reclaimer / driver / task 四层语义桥接起来。
各层设计动机:
| 层 | 设计动机 |
|---|---|
| L1 RAII | 自动配对 enter/leave,避免遗漏导致 numThreads_ 不一致;throw-safe |
| L2 Pool 委托 | Pool 不直接知道 task;通过 reclaimer 间接桥接 |
| L3 多态 virtual | 不同 pool 层级(Root/Task/Node/Operator)对"进入仲裁"语义不同 |
| L4 Task 计数 | 状态机的真正落地:原子计数变化反映 driver 在/不在跑 |
10.10.5.1 Layer 1:MemoryPoolArbitrationSection(RAII)
流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。
// MemoryArbitrator.h (简化)
class MemoryPoolArbitrationSection {
public:
explicit MemoryPoolArbitrationSection(MemoryPool* requestor)
: pool_(requestor) {
pool_->enterArbitration();
}
~MemoryPoolArbitrationSection() noexcept {
pool_->leaveArbitration(); // ★ noexcept 是关键,析构不能 throw
}
private:
MemoryPool* const pool_;
};
// 调用位置示例(原稿行号):MemoryPool.cpp:979
void MemoryPoolImpl::growCapacity(MemoryPool* requestor, uint64_t size) {
MemoryPoolArbitrationSection sec(requestor);
arbitrator_->growCapacity(this, size);
// sec 析构 → leaveArbitration 自动调用
}
RAII 价值:仲裁中途 throw(如 abort)触发栈展开,dtor 自动跑 leaveArbitration,状态不残留;任何 return 路径都不会漏。
10.10.5.2 Layer 2:MemoryPoolImpl 委托
// MemoryPool.cpp:1201
void MemoryPoolImpl::enterArbitration() {
if (reclaimer_ != nullptr) reclaimer_->enterArbitration();
}
void MemoryPoolImpl::leaveArbitration() noexcept {
if (reclaimer_ != nullptr) reclaimer_->leaveArbitration();
}
只做"指针检查 + 转发"。Pool 不应硬编码"task 怎么 suspend"这种特定语义——挂哪个 Reclaimer 决定 enter/leave 实际干什么。
10.10.5.3 Layer 3:MemoryReclaimer 多态分发
class MemoryReclaimer {
public:
// 基类默认 no-op
virtual void enterArbitration() {}
virtual void leaveArbitration() noexcept {}
};
MemoryReclaimer 的 enter/leave 行为由实际挂接的实现决定。当前 exec::MemoryReclaimer 也可借助线程上下文桥接 Driver 的 suspended 协议,Operator reclaimer 另有对应实现;不能声称只有 Operator 一层会响应进入仲裁,或把所有 root/task/node 一概视为 no-op。
当前源码摘录(1d1b76567870):velox/exec/Operator.cpp:682。省略外围声明;此片段未作为独立程序编译。
void Operator::MemoryReclaimer::enterArbitration() {
DriverThreadContext* driverThreadCtx = driverThreadContext();
if (FOLLY_UNLIKELY(driverThreadCtx == nullptr)) {
// Skips the driver suspension handling if this memory arbitration request
// is not issued from a driver thread. For example, async streaming shuffle
// and table scan prefetch execution path might initiate memory arbitration
// request from non-driver thread.
return;
}
Driver* const runningDriver = driverThreadCtx->driverCtx()->driver;
if (!FLAGS_velox_memory_pool_capacity_transfer_across_tasks) {
if (auto opDriver = ensureDriver()) {
// NOTE: the current running driver might not be the driver of the
// operator that requests memory arbitration. The reason is that an
// operator might extend the buffer allocated from the other operator
// either from the same or different drivers. But they must be from the
// same task as the following check. User could set
// FLAGS_transferred_arbitration_allowed=true to bypass this check.
VELOX_CHECK_EQ(
runningDriver->task()->taskId(),
opDriver->task()->taskId(),
"The current running driver and the request driver must be from the same task");
}
}
if (runningDriver->task()->enterSuspended(runningDriver->state()) !=
StopReason::kNone) {
// There is no need for arbitration if the associated task has already
// terminated.
VELOX_FAIL("Terminate detected when entering suspension");
}
}
driver->state() 不是抽象的状态枚举,而是 Driver 的内部状态对象(包含 isOnThread、阻塞原因等),让 Task 精确知道是哪个 driver 进入 suspended。
10.10.5.4 Layer 4:Task::enter/leaveSuspended
// 当前 Task.cpp 原样节选:嵌套 suspended、锁内计数、锁外通知与 pause 重检查。
StopReason Task::enterSuspended(ThreadState& state) {
VELOX_CHECK(!state.hasBlockingFuture);
VELOX_CHECK(state.isOnThread());
std::vector<ContinuePromise> threadFinishPromises;
auto guard = folly::makeGuard([&]() {
for (auto& promise : threadFinishPromises) {
promise.setValue();
}
});
std::lock_guard<std::timed_mutex> l(mutex_);
if (state.isTerminated) {
return StopReason::kAlreadyTerminated;
}
const auto reason = shouldStopLocked();
if (reason == StopReason::kTerminate) {
state.isTerminated = true;
return StopReason::kTerminate;
}
// A pause will not stop entering the suspended section. It will just ack that
// the thread is no longer inside the driver executor pool.
VELOX_CHECK(
reason == StopReason::kNone || reason == StopReason::kPause ||
reason == StopReason::kYield,
"Unexpected stop reason on suspension: {}",
reason);
if (++state.numSuspensions > 1) {
// Only the first suspension request needs to update the running driver
// thread counter in the task.
return StopReason::kNone;
}
if (--numThreads_ == 0) {
threadFinishPromises = allThreadsFinishedLocked();
}
VELOX_CHECK_GE(numThreads_, 0);
return StopReason::kNone;
}
StopReason Task::leaveSuspended(ThreadState& state) {
VELOX_CHECK(!state.hasBlockingFuture);
VELOX_CHECK(state.isOnThread());
TestValue::adjust("facebook::velox::exec::Task::leaveSuspended", this);
for (;;) {
{
std::lock_guard<std::timed_mutex> l(mutex_);
VELOX_CHECK_GT(state.numSuspensions, 0);
auto leaveGuard = folly::makeGuard([&]() {
VELOX_CHECK_GE(numThreads_, 0);
if (--state.numSuspensions == 0) {
// Only the last suspension leave needs to update the running driver
// thread counter in the task
++numThreads_;
}
});
if (state.numSuspensions > 1 || !pauseRequested_) {
if (state.isTerminated) {
return StopReason::kAlreadyTerminated;
}
if (terminateRequested_) {
state.isTerminated = true;
return StopReason::kTerminate;
}
// If we have more than one suspension requests on this driver thread or
// the task has been resumed, then we return here.
return StopReason::kNone;
}
VELOX_CHECK_GT(state.numSuspensions, 0);
VELOX_CHECK_GE(numThreads_, 0);
leaveGuard.dismiss();
}
// If the pause flag is on when trying to reenter, sleep a while outside of
// the mutex and recheck. This is rare and not time critical. Can happen if
// memory interrupt sets pause while already inside a suspended section for
// other reason, like IO.
std::this_thread::sleep_for(std::chrono::milliseconds(10)); // NOLINT
}
}
这一层的关键设计:与 pause 协议原生双向集成——enterSuspended 内部检查"是不是 numThreads_=0 后该唤醒 pauseFuture";leaveSuspended 内部检查"task 是否仍 paused 需要继续等"。不需要外部"check + notify"包装。
10.10.5.5 完整调用链回放
D0 driver 线程:
growCapacity(...) {
MemoryPoolArbitrationSection sec(this); ← Layer 1 ctor
this->enterArbitration(); ← Layer 2
reclaimer_->enterArbitration(); ← Layer 3 多态
// Operator::MemoryReclaimer 重写:
driverCtx->task->enterSuspended(state); ← Layer 4
// state.numSuspensions++ ; numThreads_--
// state.numSuspensions > 0(状态判定,不是赋值)
arbitrator_->growCapacity(this, size); ← 实际仲裁工作
// sec 析构开始 ↓
this->leaveArbitration(); ← Layer 2
reclaimer_->leaveArbitration(); ← Layer 3
driverCtx->task->leaveSuspended(state); ← Layer 4
// state.numSuspensions-- ; numThreads_++
// state.numSuspensions == 0(状态判定,不是赋值)
}
10.10.5.6 enterSuspended 同时支撑两种语义
enterSuspended/leaveSuspended 是 Layer 4 底层原语,被两种上层场景复用:
| 上层场景 | 调用路径 | 是否物理 BLOCK |
|---|---|---|
| Requestor 进入仲裁(SUSPENDED 语义) | enterArbitration → Operator::MemoryReclaimer |
否,继续跑仲裁代码 |
| Victim driver 响应 pause(PAUSED 语义) | Driver::runInternal 检查点返回 kPause | 由 Task::resume 续调;leaveSuspended 配合 pause 是另一个保栈路径 |
// 流程伪代码:普通 pause 与 suspended 不是同一路径。
普通 Driver:
runInternal 检查 shouldStop() -> kPause
CancelGuard 析构 -> Task::leave -> 撤销本次线程登记
Driver::run 收到 kPause -> return
Task::resume -> 按状态选择重新 enqueue
处于 suspended 的 Driver:
调用栈保留,首次 enterSuspended 已从 numThreads_ 扣除
leaveSuspended 若发现 Task 仍 paused,锁外等待后重新检查
最后一层 suspension 退出时恢复计数,再继续原调用栈
Task::resume 不为它重复 enqueue
暂停与 suspended 都参与 Task 活跃线程计数,但进入和恢复路径不同:普通 pause 可让 Driver 返回,suspended 则保留调用栈。共享部分计数与通知机制不代表它们调用完全相同的 API 或使用相同等待方式。
10.10.5.7 设计精髓汇总
| 设计要点 | 价值 |
|---|---|
| RAII guard 在最外层 | 编译期保证 enter/leave 配对,throw-safe |
leaveArbitration 声明 noexcept |
析构安全;leave 不应失败 |
| Pool 委托给 Reclaimer 多态分发 | Pool 不绑定具体执行语义 |
| Task 继承 exec 层 enter/leave;是否进入 suspended 由当前 Driver 线程上下文决定 | 按实际挂接的 reclaimer 分发,exec::MemoryReclaimer 同样能桥接 suspended;不承诺零开销 |
| Operator reclaimer 和 exec::MemoryReclaimer 等实现分别结合其上下文处理进入仲裁;应按创建/挂接路径确认具体多态分发。 | exec::MemoryReclaimer 通过 driverThreadContext() 找到 Driver 并进入/离开 suspended;没有 Driver 上下文时跳过这一桥接。 |
enterSuspended 与 pause 协议双向集成 |
共享 Task 停稳计数语义;普通 pause 与保栈 suspended 使用不同的进入/恢复路径 |
state 对象传引用 |
Task 能精确知道是哪个 driver,避免多对一查找 |
| Layer 4 内部处理唤醒 | 不需要外部 check+notify,原子操作内联完成 |
整体优雅性:通过严格分层让每层只关心自己的事——Pool 不知道 Task 存在、Reclaimer 不知道 Driver 内部、Task 不知道仲裁逻辑——但通过 enter/leave 这条桥连起来形成完整的协调系统。
10.10.6 Local 阶段:谁 suspend/resume
D0 driver 线程:
① reserve(1024 MiB) → maybeIncrementReservation 返回 false
② MemoryPoolArbitrationSection 构造
→ enterArbitration → Task::enterSuspended()
★ D0 标记为 SUSPENDED:numThreads_-- ; state.numSuspensions++
(但 D0 物理上还在 running 仲裁代码,没让出 CPU)
③ Step 1 maybeGrowFromSelf ─┐ 成功即返回,无任何 pause
④ Step 2 ensureCapacity │ D0 持续
⚠️ 若撞 maxCap 需要 self-reclaim │ 跑在
→ 通过 reclaim 链触发 task->requestPause() │ 自己线程
→ D1/D2/D3(同 task 其他 driver)短暂 PAUSED │ 上
→ spill self 后立即 task->resume() │
⑤ Step 3 growWithFreeCapacity │
无 pause │
⑥ Step 4 reclaimUnusedCapacity │
只动其他 query 的"未用 cap"账面 │
★ 不 pause 其他人 ★ │
⑦ 成功或失败 ────────────────────┘
⑧ MemoryPoolArbitrationSection 析构
→ leaveArbitration → Task::leaveSuspended()
★ D0 numThreads_++ → 恢复 RUNNING
Local 中谁受影响:
| 谁 | 状态 | 仅在哪些 Step | 解除时机 |
|---|---|---|---|
| D0(requestor) | SUSPENDED 全程 | 进入仲裁到析构 | leaveArbitration |
| D1/D2/D3(同 task 其他 driver) | 仅 Step 2 self-reclaim 时 PAUSED | Step 2 | task->resume()(reclaim 完立即) |
| 其他 query 的 driver | 通常不暂停其他 query;收回其空闲 capacity 会影响后续增长,CPU / allocator 等共享资源也可能相互影响。 | — | — |
只有实际进入 self-reclaim 等需要停稳算子的路径时,local 阶段才会请求相关 Task pause;free capacity 成功或仅收缩空闲额度的路径通常不需要。具体频率需要运行统计,不能给出没有测量依据的百分比。
10.10.7 Global 阶段:完整线程交互时序
[T0] D0 driver thread:
reserve(1024 MiB) failed at root pool
→ growCapacity
→ MemoryPoolArbitrationSection (enterSuspended)
numThreads_-- ; state.numSuspensions++
→ ScopedArbitration (thread-local = kLocal)
→ Step 1-4 全部同步执行
→ 都失败 → startAndWaitGlobalArbitration
→ 加入 globalArbitrationWaiters_,唤醒后台
→ future.wait(timeout) ← D0 BLOCK 让出 CPU
[T1] background arb thread:
被唤醒
→ GlobalArbitrationSection (thread-local = kGlobal)
→ runGlobalArbitration
→ 计算 target、扫所有 participants
→ 选 Q2/T1 为 victim
→ createAsyncMemoryReclaimTask
└─ 把 kGlobal context 打包给新线程
→ memoryReclaimExecutor_.add(task)
[T2] reclaim thread pool worker:
接到任务 → ScopedMemoryArbitrationContext 恢复 kGlobal
→ ArbitrationTimedLock(reclaimMutex_, 0) ← 0 = 无超时
→ task->requestPause() → pauseRequested_=true
→ future.wait() ← 等 victim 全部 driver pause
[T3] victim's driver threads (Q2/T1 的 D8-D11):
各自在 runInternal 循环检测 pauseRequested_
→ 完成当前算子 step
→ enterSuspended (numThreads_--)
→ 等待 resume ← 自己 PAUSED
[T2] reclaim thread (待 numThreads_=0):
→ pool->reclaim() → HashJoin spill 写盘 (1500MB)
→ participant->shrink(false) ← 归还 cap 给仲裁器
→ task->resume() ← 解除 pause
[T3] victim drivers:
从 pause 恢复,继续执行(数据已在磁盘,按 merge-spill 路径读回)
[T1] background arb thread:
→ allocateCapacityLocked(D0 的 op)
→ checkedGrow(participant=Q1, allocated=1024MB)
→ op.notifyWaitFuture() ← 唤醒 D0
[T0] D0 driver thread:
从 future.wait() 返回
→ ScopedArbitration dtor: finishArbitration
→ MemoryPoolArbitrationSection dtor: leaveSuspended
numThreads_++ ; state.numSuspensions--
→ maybeReserve returns true → 继续 SortBuffer::addInput
10.10.8 完整协调矩阵
按时间阶段汇总各类线程的状态:
| 阶段 | D0 (requestor) | D0 同 task 其他 driver | Victim task drivers | Arb 线程 | Reclaim worker |
|---|---|---|---|---|---|
| Local Step 1/3 成功 | SUSPENDED | RUNNING | — | IDLE | IDLE |
| Local Step 2 self-reclaim | SUSPENDED | PAUSED 短暂 | — | IDLE | — |
| Local Step 4 | SUSPENDED | RUNNING | — | IDLE | IDLE |
| Local → Global 入队 | SUSPENDED | RUNNING | RUNNING | 即将唤醒 | IDLE |
| Global D0 wait | SUSPENDED + BLOCKED | RUNNING | RUNNING | RUNNING (决策) | IDLE |
| Global reclaim 期间 | SUSPENDED + BLOCKED | RUNNING | PAUSED | RUNNING | RUNNING (spill) |
| Global resume victim | SUSPENDED + BLOCKED | RUNNING | RUNNING | RUNNING | 即将完成 |
| Global notify D0 | SUSPENDED + 即将 unblock | RUNNING | RUNNING | RUNNING | IDLE |
| 退出仲裁 | RUNNING | RUNNING | RUNNING | IDLE | IDLE |
两个对比维度:
- Local 可以收缩其他 participant 的未用 capacity,self-reclaim 也会影响共享该 root 的相关 Task;Global 进一步通过后台回收轮选择 victim。应区分额度变化与真正释放已用数据。
- 同 task vs 跨 task: 自己 task 的 driver 在 self-reclaim 时受影响;非自己 task / 非 victim task 的 driver 全程不动
10.10.9 一句话总结
- Local 工作主要在 requestor 线程完成;收缩其他 root 空闲额度通常不暂停其 Driver,但 self-reclaim 需要停稳相应 Task,可能影响同 root 下其他 Driver。是否发生 spill 与 pause 必须按实际分支判断。
- Global:requestor driver 提交请求后 BLOCK 不参与;由"后台 arb 线程 + reclaim 线程池 + victim task 全体 driver 协作自暂停"三方协作完成;pause 粒度是 task,不是单 driver。
10.10.10 三方握手协议的设计(why)
前面的时序和协调矩阵讲的是"发生了什么"。这一节退一步,把全局仲裁看成一个四方握手协议(requestor / arb / rcl / victim),讲"为什么这样设计"。理解了这四点,就理解了整套协调的骨架。
10.10.10.1 (1) 三个交接点,三种不同的同步原语
全局仲裁本质是一条控制权接力链:requestor 把活交给 arb,arb 交给 rcl,rcl 让 victim 停下来。链上三个交接点用了三种完全不同的同步原语,不是偶然,而是每个交接点的耦合需求不同:
① requestor:进入 suspended,加入 waiter,保留调用栈同步等待全局仲裁结果。
② global controller:选择 participant,派发回收任务,再收集/等待本轮结果。
③ reclaim:requestPause 等 Task 安全停稳,执行可回收区内 spill,恢复 Task。
④ participant:依据实际 reservation shrink capacity;仲裁器分配可用额度并完成 waiter。
| 交接点 | 原语 | 等待语义 | 为什么是这个原语 |
|---|---|---|---|
| ① requestor → arb | globalArbitrationWaiters_ 队列 + ContinueFuture |
排队等结果 | ContinueFuture/ContinuePromise 传递的是完成或异常信号;capacity 授予通过仲裁器与 participant/pool 的状态更新体现,不是把容量数字当作 folly::Unit future 的返回值。唤醒后仍需按操作协议检查结果与终止状态。 |
| ② arb → rcl | memoryReclaimExecutor_.add(task) |
派发后收集、等待本轮回收结果 | 多个 victim 可并行;控制线程仍受该轮完成时间影响,并非提交后立刻无限处理下一轮 |
| ③ rcl → victim | task->requestPause() + numThreads_==0 future |
协作式,等对方自己停 | 不能抢占——driver 可能正在改写 RowContainer;只能"请求"暂停,等每个 driver 跑完当前算子 step 后自愿进入 suspended |
三个交接点分别连接 requestor 等仲裁结果、控制线程组织并收齐回收任务、以及 reclaim 等 Task 安全停稳。第二步包含结果等待;第三步依赖协作式停止而非抢占正在修改 RowContainer 的算子。
10.10.10.2 (2) "谁阻塞在谁身上"——依赖图无环是不死锁的根
把四方的阻塞依赖画成有向图(A→B 表示 A 阻塞等待 B):
requestor 等待 global controller 的仲裁结果
global controller 等待 reclaim 任务结果
reclaim 等待 Task 的活跃 Driver 停稳
suspended requestor 已退出活跃计数,普通 paused Driver 已返回
回收完成后 Task 恢复;participant 再 shrink,并把额度归还仲裁器
注意:executor 只是分工方式,并没有删掉 controller → reclaim 的等待边。
关键性质:
- 控制线程会等待本轮回收任务,也可能执行同步工作。避免自死锁依赖 requestor suspended 从活跃线程计数扣除、回收上下文传播、锁与对象保活等条件,不能用“arb 永不等待”作为无环证明。
- 依赖图需要标出 requestor 等仲裁、控制线程等回收结果、reclaim 等 Task 停稳及恢复后的继续条件。只把异步提交画成无等待箭头会漏掉真实依赖;要结合调用栈、锁和 numThreads_ 协议验证可能的环。
- Pause 后的 Driver 确实依赖回收作用域结束并恢复任务;是否形成死锁必须看回收还在等待谁。关键是已暂停或 suspended 的 Driver 不再阻止 Task 停稳,而不是把 resume 称为“推送信号”就自动消除依赖。
文档 §10.16 的 Race2(requestor 与 victim 同 task 自死锁)是这张图上一条特殊的潜在环:requestor 在 numThreads_ 里 → rcl 等 numThreads_=0 → 永远等不到。解法 enterSuspended() 把 requestor 从 numThreads_ 摘掉,等于在图上删掉了那条会成环的边。所以 Race2 不是孤立的 bug fix,而是"维持依赖图无环"这条总不变量的一个具体落点。
10.10.10.3 (3) 为什么是 future/promise,而不是条件变量 / 同步调用
仲裁中同时存在条件变量、future 和同步等待:后台主循环可用条件变量等待工作,requestor 同步等待全局仲裁,Task pause 及异步任务结果也需要协调。不同原语服务于不同等待条件,不能把整个实现概括为只有异步 future。
- ContinueFuture 的就绪状态表明操作可继续或以异常结束;capacity 数量保存在相应仲裁操作、participant 与 pool 状态中。Future 可以承载任意类型是库的通用能力,但这里的 ContinueFuture 是 Unit 信号,不能套用“携带容量返回值”的叙述。
- Future 就绪避免了把某次单一完成事件误当条件变量虚假唤醒,但共享状态仍可能已被取消、abort 或其他操作改变。后续检查必须遵守仲裁操作的实际协议;不能由 future ready 推导所有业务前置条件永远有效。
- 控制线程与回收执行线程可以分工并行,但仍有同步等待点,尤其是本轮回收任务收集、Task pause 和 requestor 全局等待。吞吐与延迟要结合这些等待链和 executor 容量分析,而不是假设整条路径完全解耦。
Future 与 executor 是当前实现组织可并行回收和完成通知的手段,控制线程仍会收齐结果。其他并发设计也能表达分工;源码不能支持“唯一方式”或“同步一次就让后续 waiter 永久饿死”的结论。
10.10.10.4 (4) 握手交接的不只是控制权——所有权与上下文要一起交接
最容易被忽略的一点:当控制权在线程间交接时,两样东西必须跟着一起交接,否则接力链断裂。文档把它们当成 Race6/Race7 分开讲了,这里点明它们其实是"握手协议"不可分割的一部分:
| 跟着交接的东西 | 机制 | 不交接会怎样 |
|---|---|---|
| 对象所有权(pool 保活) | ScopedArbitrationParticipant 把 victim 的 weak_ptr<pool> 升级为 shared_ptr,全程持有(Race6) |
rcl 正在 spill victim 的 pool,query 端却把 pool 析构了 → 悬空指针 |
| 仲裁上下文(thread-local) | createAsyncMemoryReclaimTask 把 kGlobal context 显式打包进 lambda,新线程 ScopedMemoryArbitrationContext 恢复(Race7) |
rcl 线程内 spill 又申请内存 → 不知道"自己在仲裁中" → 嵌套触发仲裁 → 死锁 |
线程切换不会自动继承这段仲裁上下文。当前 createAsyncMemoryReclaimTask 捕获的是提交线程的 context 指针,执行时再构造 ScopedMemoryArbitrationContext;因此提交者的上下文必须活到异步任务完成,调用链中的结果收集与 sync guard 正是需要一起审查的部分。对象保活又是另一件事,候选的强引用必须覆盖实际使用。不能把“提交到 executor”当成隐式状态已经全部按值复制的保证。
10.10.10.5 小结:四个 why 串起来
① waiter/future 表达 requestor 的一次完成或异常信号(Unit,不携带容量数字)
② executor 提供有限并行;global controller 仍需收集任务结果
③ Task pause 提供协作式停稳窗口,suspended 协议避免 requestor 阻挡自身停稳
所有权保活、锁顺序、上下文传播与异常清理共同保证协议,不以“异步”一词代替证明。
四点指向同一个设计内核:让决策(arb)保持轻量串行、让执行(rcl)异步并行、让交接(future + 打包)自包含无环。这正是 §13.2「决策/执行分离」在线程协调层面的具体兑现。
10.11 MemoryReclaimer 层次体系
10.12 Reclaimer 怎样安全地接到算子
| 层级 | 当前默认路径的职责 |
|---|---|
| QueryCtx reclaimer | 标记查询正在回收;调用基类遍历 task pools;退出时通知等待该阶段结束的 Driver |
| Task reclaimer | requestPause,等待活跃线程计数为 0;回收下层;通过 guard 调用 Task::resume |
| Node reclaimer | 普通节点可并行处理子 pool;Hash Join 使用专门实现协调共享表和 build / probe |
| Operator reclaimer | 确认 Driver 状态和 Task pause;检查 canReclaim 与 non-reclaimable section;调用 op_->reclaim |
Task pause 是协作式协议:运行中的 Driver 到检查点后停止推进;发起仲裁的线程若已 suspended,就不计入 Task 正在执行的线程数,避免 self-reclaim 等待自己。暂停并不销毁线程,也不保证任意算子内部状态都可回收,nonReclaimableSection_ 仍需参与保护。
QueryCtx::underArbitration_ 是查询执行协调状态;ScopedMemoryArbitrationContext 是当前线程的仲裁上下文;ArbitrationOperation::State 是一次请求的生命周期。三者不要合并成一个“全局仲裁状态”。跨 executor 执行回收时,createAsyncMemoryReclaimTask 显式传播所需上下文。
源码:QueryCtx 的回收阶段、Task pause / reclaim / resume、Operator 检查与回收、仲裁上下文、异步上下文传播。
10.12.1 同一 reclaimer priority,在不同范围内用途不同
全局策略倾向先牺牲优先级较低的查询;基类 MemoryReclaimer::reclaim 在一个 pool 的孩子中却按 priority 数值升序,再按可回收量降序选择。后者是子资源的回收顺序,例如把昂贵的 Join 回收放到后面。
基类会依次处理多个候选,累计达到 target 后停止,并非永远只选一个最大 child。ParallelMemoryReclaimer 则收集可回收的子池并并行处理,不能机械套用基类的串行早停逻辑。共享数据结构是否允许并行回收,需要具体 reclaimer 协调;Hash Join 就不是“所有 Driver 的表完全独立”的例子。
源码:基类候选排序与迭代、Join node 的专用 reclaimer。
10.12.2 释放已用资源与归还 capacity 分两步
ArbitrationParticipant::reclaim 先调整目标,调用 pool_->reclaim,再 shrink(false)。前一步让算子释放执行资源,带动 used / reservation 下降;后一步才降低 root capacity。SharedArbitrator 将返回的 capacity 加回空闲账本。
普通 shrink 对活跃 pool 保留的空闲余量为:
freeToKeep = min(capacity × minFreeCapacityRatio, minFreeCapacity)
shrinkable = min(max(0, freeBytes - freeToKeep),
max(0, capacity - participantMinCapacity))
例如 capacity=1 GiB、reservation=512 MiB,默认 ratio=25%、固定值=128 MiB:保留的是 128 MiB 空闲,普通 shrink 最多收回 384 MiB,留下 640 MiB capacity;不是用 max 得到 256 MiB。
reservedBytes()==0 && peakBytes()!=0 被视为 inactive pool,可跳过这类保护;shrink(reclaimAll=true) 直接尝试收回全部空闲 capacity。注意全部“空闲”仍不包括被活跃 allocation 占用的 reservation。
另一组 minReclaimBytes / minReclaimPct 控制单次 reclaim 的最低目标,使用的是 max。不要把这组公式与 shrink 的最小空闲余量混用。Reclaimer 报告的释放量、shrink 实际退回的 capacity、文件系统上的 spill bytes 也不是同一指标。
10.12.3 所有权与引用关系
Reclaimer 严格绑定到 pool(pool 拥有 reclaimer),但内部持有对应资源(task / operator)的指针来执行真正的回收逻辑:
MemoryPool
└── reclaimer_ : std::unique_ptr<MemoryReclaimer> ← pool 拥有
│
├── TaskReclaimer → 持 Task* 指针
├── ParallelMemoryReclaimer → 持 spillExecutor_
└── Operator::MemoryReclaimer → 持 Operator* 指针
- 正向:
pool->reclaimer()返回它持有的 reclaimer - 反向:reclaimer 内部存指针指向真正"能 spill"的对象
为什么这样设计:
- pool 是结构单元(树形组织),但不知道算子语义
- operator 是行为单元(HashJoin 怎么 spill),但不参与树形遍历
- reclaimer 是桥梁:挂在 pool 上让"按 pool 树遍历"的逻辑能联通;内部持算子指针让"真正干活"的逻辑能调到算子
10.12.4 基类接口
流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。
// 接口节选;定义见 MemoryArbitrator.h。
virtual void enterArbitration() {}
virtual void leaveArbitration() noexcept {}
virtual int32_t priority() const { return priority_; }
virtual bool reclaimableBytes(
const MemoryPool& pool, uint64_t& reclaimableBytes) const;
virtual uint64_t reclaim(MemoryPool* pool, uint64_t targetBytes,
uint64_t maxWaitMs, Stats& stats);
virtual void abort(MemoryPool* pool, const std::exception_ptr& error);
// reclaimableBytes / reclaim / abort 在 .cpp 中有实际默认实现。
默认基类行为:贪心向下递归
当前源码摘录(1d1b76567870):velox/common/memory/MemoryArbitrator.cpp:234。省略外围声明;此片段未作为独立程序编译。
uint64_t MemoryReclaimer::reclaim(
MemoryPool* pool,
uint64_t targetBytes,
uint64_t maxWaitMs,
Stats& stats) {
if (pool->kind() == MemoryPool::Kind::kLeaf) {
return 0;
}
// Sort the child pools based on their reclaimer priority and reserved memory.
// Reclaim from the child pool with highest priority and most reservation
// first.
struct Candidate {
std::shared_ptr<memory::MemoryPool> pool;
int64_t reclaimableBytes;
};
// NOTE: We hold candidate reference for non-reclaimable pools as well. This
// is to make sure child shared pointer is stored to keep child alive,
// avoiding destruction of child pool within below parents' 'poolMutex_' lock.
// Otherwise a double acquisition of 'poolMutex_' can happen in destructor,
// which creates deadlock.
std::vector<Candidate> nonReclaimableCandidates;
std::vector<Candidate> candidates;
{
std::shared_lock guard{pool->poolMutex_};
candidates.reserve(pool->children_.size());
nonReclaimableCandidates.reserve(pool->children_.size());
for (auto& entry : pool->children_) {
auto child = entry.second.lock();
if (child != nullptr) {
const auto reclaimableBytesOpt = child->reclaimableBytes();
if (!reclaimableBytesOpt.has_value() ||
reclaimableBytesOpt.value() == 0) {
nonReclaimableCandidates.push_back(Candidate{std::move(child), 0});
continue;
}
candidates.push_back(
Candidate{
std::move(child),
static_cast<int64_t>(reclaimableBytesOpt.value())});
}
}
}
std::sort(
candidates.begin(),
candidates.end(),
[](const auto& lhs, const auto& rhs) {
const auto lhsPrio = lhs.pool->reclaimer()->priority();
const auto rhsPrio = rhs.pool->reclaimer()->priority();
if (lhsPrio == rhsPrio) {
return lhs.reclaimableBytes > rhs.reclaimableBytes;
}
return lhsPrio < rhsPrio;
});
uint64_t reclaimedBytes{0};
for (const auto& candidate : candidates) {
VELOX_CHECK_GT(candidate.reclaimableBytes, 0);
const auto bytes = candidate.pool->reclaim(targetBytes, maxWaitMs, stats);
reclaimedBytes += bytes;
if (targetBytes != 0) {
if (bytes >= targetBytes) {
break;
}
targetBytes -= bytes;
}
}
return reclaimedBytes;
}
各层 Reclaimer 子类 overload 这个行为,加入特定动作(pause、并发分发等)。
10.12.5 按 pool 层级的 Reclaimer 类型
| Pool 层 | Reclaimer 类型 | 谁创建 | 引用资源 | 主要职责 |
|---|---|---|---|---|
| Root Pool | 基类 MemoryReclaimer(MemoryReclaimer::create()) |
QueryCtx::createPool() 调 addRootPool 时强制传入 |
— | 贪心路由:按 reclaimableBytes 降序遍历 Task children,逐个调 reclaim;自己不持数据,不 pause |
| Task Pool | Task::MemoryReclaimer (TaskReclaimer) |
Task::initTaskPool() |
Task* |
**触发 task->requestPause()**,等 numThreads_=0,递归向下 |
| Node Pool | ParallelMemoryReclaimer |
Task::getOrAddNodePool() |
spillExecutor_ |
跨多 driver 的 op pool fork-join 并发 spill |
| Operator Pool | Operator::MemoryReclaimer 实例 |
Operator 构造时 pool()->setReclaimer() |
Operator* |
调算子虚函数 op->reclaim(),转接到具体 spill |
10.12.6 TaskReclaimer:触发 pause + 委托基类向下递归
Task::MemoryReclaimer 是 Task 的嵌套类,挂在 Task Pool 上。其核心职责是在 reclaim 子 pool 之前 pause 整个 task,spill 完成后 resume。
流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。
// Task.h 的继承关系与关键成员节选。
class MemoryReclaimer : public exec::MemoryReclaimer {
public:
uint64_t reclaim(memory::MemoryPool* pool, uint64_t targetBytes,
uint64_t maxWaitMs,
memory::MemoryReclaimer::Stats& stats) override;
void abort(memory::MemoryPool* pool,
const std::exception_ptr& error) override;
private:
// 构造函数把 priority 传给 exec::MemoryReclaimer(priority)。
// enter/leave 继承 exec 层;估算继承 memory 层。
std::weak_ptr<Task> task_;
};
**为什么用 weak_ptr
Task 持有 shared_ptr<MemoryPool> pool_ ← Task 拥有 Task Pool
Task Pool 持有 unique_ptr<MemoryReclaimer> ← Pool 拥有 Reclaimer
Reclaimer 内部要回指 Task
如果用 shared_ptr<Task> → 循环引用 → 永不释放
所以用 weak_ptr,需要时 lock() 升级
10.12.6.1 reclaim() 实现
当前源码摘录(1d1b76567870):velox/exec/Task.cpp:3864。省略外围声明;此片段未作为独立程序编译。
uint64_t Task::MemoryReclaimer::reclaimTask(
const std::shared_ptr<Task>& task,
uint64_t targetBytes,
uint64_t maxWaitMs,
memory::MemoryReclaimer::Stats& stats) {
auto resumeGuard = folly::makeGuard([&]() {
try {
Task::resume(task);
} catch (const VeloxRuntimeError& exception) {
LOG(WARNING) << "Failed to resume task " << task->taskId_
<< " after memory reclamation: " << exception.message();
}
});
uint64_t reclaimWaitTimeUs{0};
bool paused{true};
{
MicrosecondWallTimer timer{&reclaimWaitTimeUs};
if (maxWaitMs == 0) {
task->requestPause().wait();
} else {
paused = task->requestPause().wait(std::chrono::milliseconds(maxWaitMs));
}
}
VELOX_CHECK(paused || maxWaitMs != 0);
if (!paused) {
RECORD_METRIC_VALUE(kMetricTaskMemoryReclaimWaitTimeoutCount, 1);
VELOX_FAIL(
"Memory reclaim failed to wait for task {} to pause after {} with max timeout {}",
task->taskId(),
succinctMicros(reclaimWaitTimeUs),
succinctMillis(maxWaitMs));
}
stats.reclaimWaitTimeUs += reclaimWaitTimeUs;
RECORD_METRIC_VALUE(kMetricTaskMemoryReclaimCount);
RECORD_HISTOGRAM_METRIC_VALUE(
kMetricTaskMemoryReclaimWaitTimeMs, reclaimWaitTimeUs / 1'000);
// Don't reclaim from a cancelled task as it will terminate soon.
if (task->isCancelled()) {
return 0;
}
uint64_t reclaimedBytes{0};
try {
uint64_t reclaimExecTimeUs{0};
{
MicrosecondWallTimer timer{&reclaimExecTimeUs};
reclaimedBytes = memory::MemoryReclaimer::reclaim(
task->pool(), targetBytes, maxWaitMs, stats);
}
RECORD_HISTOGRAM_METRIC_VALUE(
kMetricTaskMemoryReclaimExecTimeMs, reclaimExecTimeUs / 1'000);
} catch (...) {
// Set task error before resumes the task execution as the task operator
// might not be in consistent state anymore. This prevents any off thread
// from running again.
task->setError(std::current_exception());
std::rethrow_exception(std::current_exception());
}
return reclaimedBytes;
}
void Task::MemoryReclaimer::abort(
memory::MemoryPool* pool,
const std::exception_ptr& error) {
auto task = ensureTask();
if (FOLLY_UNLIKELY(task == nullptr)) {
return;
}
VELOX_CHECK_EQ(task->pool()->name(), pool->name());
task->setError(error);
// TODO: respect the memory arbitration request timeout later.
const static uint32_t maxTaskAbortWaitUs = 6'000'000; // 60s
if (task->taskCompletionFuture().wait(
std::chrono::microseconds(maxTaskAbortWaitUs))) {
// If task is completed within 60s wait, we can safely propagate abortion.
// Otherwise long running operators might be in the middle of processing,
// making it unsafe to force abort. In this case we let running operators
// finish by hitting operator boundary, and rely on cleanup mechanism to
// release the resource.
memory::MemoryReclaimer::abort(pool, error);
} else {
LOG(WARNING)
<< "Timeout waiting for task to complete during query memory aborting.";
}
}
void Task::DriverBlockingState::setDriverFuture(
ContinueFuture& driverFuture,
Operator* driverOp,
BlockingReason blockingReason) {
VELOX_CHECK(!blocked_);
VELOX_CHECK_NULL(op_);
VELOX_CHECK_EQ(blockingReason_, BlockingReason::kNotBlocked);
{
std::lock_guard<std::mutex> l(mutex_);
VELOX_CHECK(promises_.empty());
VELOX_CHECK_NULL(error_);
blocked_ = true;
op_ = driverOp;
blockingReason_ = blockingReason;
blockStartUs_ = getCurrentTimeMicro();
}
std::move(driverFuture)
.via(&folly::InlineExecutor::instance())
.thenValue(
[&, driverHolder = driver_->shared_from_this()](auto&& /* unused */) {
std::vector<std::unique_ptr<ContinuePromise>> promises;
{
std::lock_guard<std::mutex> l(mutex_);
VELOX_CHECK(blocked_);
VELOX_CHECK_NULL(error_);
promises = std::move(promises_);
if ((op_ != nullptr) && !driver_->state().isTerminated) {
VELOX_CHECK_NE(blockingReason_, BlockingReason::kNotBlocked);
op_->recordBlockingTime(blockStartUs_, blockingReason_);
}
clearLocked();
}
for (auto& promise : promises) {
promise->setValue();
}
})
.thenError(
folly::tag_t<std::exception>{},
[&, driverHolder = driver_->shared_from_this()](
std::exception const& e) {
std::vector<std::unique_ptr<ContinuePromise>> promises;
{
std::lock_guard<std::mutex> l(mutex_);
VELOX_CHECK(blocked_);
VELOX_CHECK_NULL(error_);
promises = std::move(promises_);
try {
VELOX_FAIL(
"A driver future from task {} was realized with error: {}",
driver_->task()->taskId(),
e.what());
} catch (const VeloxException&) {
error_ = std::current_exception();
}
clearLocked();
}
for (auto& promise : promises) {
promise->setValue();
}
});
}
void Task::DriverBlockingState::clearLocked() {
VELOX_CHECK(promises_.empty());
op_ = nullptr;
blockingReason_ = BlockingReason::kNotBlocked;
blockStartUs_ = 0;
blocked_ = false;
}
bool Task::DriverBlockingState::blocked(ContinueFuture* future) {
VELOX_CHECK_NOT_NULL(future);
std::lock_guard<std::mutex> l(mutex_);
if (error_ != nullptr) {
std::rethrow_exception(error_);
}
if (!blocked_) {
VELOX_CHECK(promises_.empty());
return false;
}
auto [blockPromise, blockFuture] = makeVeloxContinuePromiseContract(
fmt::format(
"DriverBlockingState {} from task {}",
driver_->driverCtx()->driverId,
driver_->task()->taskId()));
*future = std::move(blockFuture);
promises_.emplace_back(
std::make_unique<ContinuePromise>(std::move(blockPromise)));
return true;
}
}
关键点:
| 设计要素 | 用意 |
|---|---|
weak_ptr<Task> + lock() 检查 |
task 可能在 reclaim 前已被销毁,安全降级返回 0 |
| RAII resumeGuard + reclaim 异常处理 | 离开作用域时尝试 Task::resume;reclaim 抛异常时先 setError,再展开栈,防止恢复可能不一致的算子状态。当前 guard 会捕获并记录恢复时的 VeloxRuntimeError。 |
future.wait(maxWaitMs) 超时 |
等待的是 Task 活跃线程计数归零。maxWaitMs 为 0 表示该等待不设超时;非零且超时会抛错。nonReclaimableSection_ 是另外一项回收资格检查,并不是 Task pause 计数器。 |
复用基类 reclaim() 实现 |
先锁定 weak Task 引用并校验 pool,再停稳;若 Task 已取消则返回。基类按 priority、可回收量排序并依次回收,达到目标早停,而不是只选一个最大 child。 |
priority_ 字段 |
Task reclaimer priority 影响该子树的回收排序。全局 victim 排序读取 root reclaimer priority,并结合容量分组等条件;不能把 Task 优先级直接当成全局字段。 |
10.12.6.2 reclaim 流程的 "三明治结构"
Task::MemoryReclaimer 继承的是 exec::MemoryReclaimer,后者已经实现 enter/leave:在存在 Driver 线程上下文时调用 Task 的 suspended 协议。Task 子类没有再次重写这两个方法,继承到的也不是 no-op。它另外重写 reclaim/abort,在 victim Task 上建立停稳窗口、处理异常并恢复或关闭 Driver;requestor 进入 suspended 与 victim Task 停稳承担不同职责。
10.12.6.3 abort() 实现
当前源码摘录(1d1b76567870):velox/exec/Task.cpp:3928。省略外围声明;此片段未作为独立程序编译。
void Task::MemoryReclaimer::abort(
memory::MemoryPool* pool,
const std::exception_ptr& error) {
auto task = ensureTask();
if (FOLLY_UNLIKELY(task == nullptr)) {
return;
}
VELOX_CHECK_EQ(task->pool()->name(), pool->name());
task->setError(error);
// TODO: respect the memory arbitration request timeout later.
const static uint32_t maxTaskAbortWaitUs = 6'000'000; // 60s
if (task->taskCompletionFuture().wait(
std::chrono::microseconds(maxTaskAbortWaitUs))) {
// If task is completed within 60s wait, we can safely propagate abortion.
// Otherwise long running operators might be in the middle of processing,
// making it unsafe to force abort. In this case we let running operators
// finish by hitting operator boundary, and rely on cleanup mechanism to
// release the resource.
memory::MemoryReclaimer::abort(pool, error);
} else {
LOG(WARNING)
<< "Timeout waiting for task to complete during query memory aborting.";
}
}
Abort 设置错误/终止状态并进入 Task 的清理协议,不是强制终止正在执行的 OS 线程。On-thread Driver 需要在协作式检查点退出,off-thread Driver 才可按终止资格清理;仍被引用的数据和 pool 对象可能稍后才析构。
10.12.6.4 继承基类的 reclaimableBytes()
当前源码摘录(1d1b76567870):velox/common/memory/MemoryArbitrator.cpp:216。省略外围声明;此片段未作为独立程序编译。
bool MemoryReclaimer::reclaimableBytes(
const MemoryPool& pool,
uint64_t& reclaimableBytes) const {
reclaimableBytes = 0;
if (pool.kind() == MemoryPool::Kind::kLeaf) {
return false;
}
bool reclaimable{false};
pool.visitChildren([&](MemoryPool* pool) {
auto reclaimableBytesOpt = pool->reclaimableBytes();
reclaimable |= reclaimableBytesOpt.has_value();
reclaimableBytes += reclaimableBytesOpt.value_or(0);
return true;
});
VELOX_CHECK(reclaimable || reclaimableBytes == 0);
return reclaimable;
}
canReclaim() 检查 task 是否正在运行(非 terminating / terminated 状态)。已经在退出过程中的 task 不应被选为 victim。
10.12.7 ParallelMemoryReclaimer:fork-join 调度器
ParallelMemoryReclaimer 挂在 Node Pool 上,实现的是 "把 children 的 reclaim 任务 fork 到 spillExecutor,自己 BLOCK 等结果":
当前源码摘录(1d1b76567870):velox/exec/MemoryReclaimer.cpp:85。省略外围声明;此片段未作为独立程序编译。
uint64_t ParallelMemoryReclaimer::reclaim(
memory::MemoryPool* pool,
uint64_t targetBytes,
uint64_t maxWaitMs,
Stats& stats) {
if (executor_ == nullptr) {
return memory::MemoryReclaimer::reclaim(
pool, targetBytes, maxWaitMs, stats);
}
// Sort candidates based on priority.
struct Candidate {
std::shared_ptr<memory::MemoryPool> pool;
int64_t reclaimableBytes;
};
std::vector<Candidate> candidates;
{
std::shared_lock guard{pool->poolMutex_};
candidates.reserve(pool->children_.size());
for (auto& entry : pool->children_) {
auto child = entry.second.lock();
if (child != nullptr) {
const auto reclaimableBytesOpt = child->reclaimableBytes();
if (!reclaimableBytesOpt.has_value()) {
continue;
}
candidates.push_back(
Candidate{
std::move(child),
static_cast<int64_t>(reclaimableBytesOpt.value())});
}
}
}
struct ReclaimResult {
const uint64_t reclaimedBytes{0};
const Stats stats;
const std::exception_ptr error{nullptr};
explicit ReclaimResult(std::exception_ptr _error)
: reclaimedBytes(0), error(std::move(_error)) {}
ReclaimResult(uint64_t _reclaimedBytes, Stats&& _stats)
: reclaimedBytes(_reclaimedBytes),
stats(std::move(_stats)),
error(nullptr) {}
};
std::vector<std::shared_ptr<AsyncSource<ReclaimResult>>> reclaimTasks;
for (const auto& candidate : candidates) {
if (candidate.reclaimableBytes == 0) {
continue;
}
reclaimTasks.push_back(
memory::createAsyncMemoryReclaimTask<ReclaimResult>(
[&, reclaimPool = candidate.pool]() {
try {
Stats reclaimStats;
const auto bytes =
reclaimPool->reclaim(targetBytes, maxWaitMs, reclaimStats);
return std::make_unique<ReclaimResult>(
bytes, std::move(reclaimStats));
} catch (const std::exception& e) {
VELOX_MEM_LOG(ERROR) << "Reclaim from memory pool "
<< pool->name() << " failed: " << e.what();
// The exception is captured and thrown by the caller.
return std::make_unique<ReclaimResult>(
std::current_exception());
}
}));
if (reclaimTasks.size() > 1) {
executor_->add([source = reclaimTasks.back()]() { source->prepare(); });
}
}
auto syncGuard = folly::makeGuard([&]() {
for (auto& reclaimTask : reclaimTasks) {
// We consume the result for the pending tasks. This is a cleanup in the
// guard and must not throw. The first error is already captured before
// this runs.
try {
reclaimTask->move();
} catch (const std::exception&) {
}
}
});
uint64_t reclaimedBytes{0};
for (auto& reclaimTask : reclaimTasks) {
const auto result = reclaimTask->move();
if (result->error) {
std::rethrow_exception(result->error);
}
stats += result->stats;
reclaimedBytes += result->reclaimedBytes;
}
return reclaimedBytes;
}
详细 fork-join 安全前提与示意见 §10.13 "并行 spill 安全的两个必要前提",此处不重复。
10.12.8 Operator::MemoryReclaimer:连接 pool 树和算子的桥梁
这是最关键的一层,把"内存树遍历"对接到"算子 spill 逻辑":
流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。
// Operator.h (简化)
class Operator {
public:
class MemoryReclaimer : public memory::MemoryReclaimer {
public:
static std::unique_ptr<MemoryReclaimer> create(Operator* op) {
return std::unique_ptr<MemoryReclaimer>(new MemoryReclaimer(op));
}
uint64_t reclaim(MemoryPool* pool, uint64_t target,
uint64_t maxWaitMs, Stats& stats) override {
// ① 检查算子是否拒绝
if (op_->nonReclaimableSection_) return 0;
// ② 委托到算子虚函数
return op_->reclaim(target, stats);
}
void enterArbitration() override {
// ③ 进入仲裁:把当前 driver suspend
auto* driverCtx = op_->operatorCtx()->driverCtx();
driverCtx->task->enterSuspended(driverCtx->driver->state());
}
void leaveArbitration() noexcept override {
auto* driverCtx = op_->operatorCtx()->driverCtx();
driverCtx->task->leaveSuspended(driverCtx->driver->state());
}
void abort(MemoryPool* pool, const std::exception_ptr& error) override {
op_->close(); // 关闭算子,丢弃所有状态
}
private:
explicit MemoryReclaimer(Operator* op) : op_(op) {}
Operator* const op_; // ← 反向指针到算子
};
};
注意:Operator::MemoryReclaimer 本身不实现具体 spill,它调用 op_->reclaim(target, stats) 这个 Operator 基类虚函数,由具体算子重写。
10.12.9 算子级 spill 实现(重写 Operator::reclaim())
不同算子各自实现 spill 逻辑,统一通过 Operator::reclaim() 虚函数分派,不需要为每个算子写独立的 MemoryReclaimer 子类:
当前源码摘录(1d1b76567870):velox/exec/HashBuild.cpp:1314。省略外围声明;此片段未作为独立程序编译。
void HashBuild::reclaim(
uint64_t /*unused*/,
memory::MemoryReclaimer::Stats& stats) {
TestValue::adjust("facebook::velox::exec::HashBuild::reclaim", this);
VELOX_CHECK(canSpill());
auto* driver = operatorCtx_->driver();
VELOX_CHECK_NOT_NULL(driver);
VELOX_CHECK(!nonReclaimableSection_);
const auto* config = spillConfig();
VELOX_CHECK_NOT_NULL(config);
if (UNLIKELY(exceededMaxSpillLevelLimit_)) {
// 'canReclaim()' already checks the spill limit is not exceeding max, there
// is only a small chance from the time 'canReclaim()' is checked to the
// actual reclaim happens that the operator has spilled such that the spill
// level exceeds max.
LOG(WARNING)
<< "Can't reclaim from hash build operator, exceeded maximum spill "
"level of "
<< config->maxSpillLevel << ", " << pool()->name()
<< ", root pool: " << pool()->root()->name()
<< ", used: " << succinctBytes(pool()->usedBytes())
<< ", reservation: " << succinctBytes(pool()->reservedBytes())
<< ", root pool reservation: "
<< succinctBytes(pool()->root()->reservedBytes());
return;
}
// NOTE: a hash build operator is reclaimable if it is in the middle of table
// build processing and is not under non-reclaimable execution section.
if (nonReclaimableState()) {
// TODO: reduce the log frequency if it is too verbose.
RECORD_METRIC_VALUE(kMetricMemoryNonReclaimableCount);
++stats.numNonReclaimableAttempts;
LOG(WARNING) << "Can't reclaim from hash build operator, state_["
<< stateName(state_) << "], nonReclaimableSection_["
<< nonReclaimableSection_ << "], spiller_["
<< (stateCleared_ ? "cleared"
: spiller_ == nullptr ? "null"
: spiller_->finalized() ? "finalized"
: "non-finalized")
<< "] " << pool()->name()
<< ", root pool: " << pool()->root()->name()
<< ", used: " << succinctBytes(pool()->usedBytes())
<< ", reservation: " << succinctBytes(pool()->reservedBytes())
<< ", root pool reservation: "
<< succinctBytes(pool()->root()->reservedBytes());
return;
}
const auto& task = driver->task();
VELOX_CHECK(task->pauseRequested());
const std::vector<Operator*> operators =
task->findPeerOperators(operatorCtx_->driverCtx()->pipelineId, this);
for (auto* op : operators) {
HashBuild* buildOp = dynamic_cast<HashBuild*>(op);
VELOX_CHECK_NOT_NULL(buildOp);
VELOX_CHECK(buildOp->canSpill());
if (buildOp->nonReclaimableState()) {
// TODO: reduce the log frequency if it is too verbose.
RECORD_METRIC_VALUE(kMetricMemoryNonReclaimableCount);
++stats.numNonReclaimableAttempts;
LOG(WARNING) << "Can't reclaim from hash build operator, state_["
<< stateName(buildOp->state_) << "], nonReclaimableSection_["
<< buildOp->nonReclaimableSection_ << "], spiller_["
<< (buildOp->stateCleared_ ? "cleared"
: buildOp->spiller_ == nullptr ? "null"
: buildOp->spiller_->finalized() ? "finalized"
: "non-finalized")
<< "], " << buildOp->pool()->name()
<< ", root pool: " << buildOp->pool()->root()->name()
<< ", used: " << succinctBytes(buildOp->pool()->usedBytes())
<< ", reservation: "
<< succinctBytes(buildOp->pool()->reservedBytes())
<< ", root pool reservation: "
<< succinctBytes(buildOp->pool()->root()->reservedBytes());
return;
}
}
std::vector<HashBuildSpiller*> spillers;
for (auto* op : operators) {
HashBuild* buildOp = static_cast<HashBuild*>(op);
spillers.push_back(buildOp->spiller_.get());
}
spillHashJoinTable(spillers, config);
for (auto* op : operators) {
HashBuild* buildOp = static_cast<HashBuild*>(op);
buildOp->table_->clear(true);
buildOp->pool()->release();
}
}
10.12.10 整条 reclaim 调用链(带 reclaimer 类型)
reclaim worker
ArbitrationParticipant::reclaim
调整最低回收目标,取得 reclaimMutex_
root->reclaim
QueryCtx 缺省 reclaimer:设置查询仲裁状态,结束时通知等待者
memory 基类:按 priority / 可回收量排序,遍历 Task 子树并早停
Task reclaimer:保活 Task,安装 resume guard,等待 requestPause
memory 基类:按相同规则遍历 Node 子树
ParallelMemoryReclaimer:有 executor 时用 AsyncSource 组织子池
Operator reclaimer:保活 Driver、检查回收资格与 Task 停稳
op->reclaim:按 HashBuild / OrderBy 等自身状态执行 spill
发生回收异常:先 setError,再让 resume guard 完成清理
shrink(false):降低可归还的 root capacity
SharedArbitrator:将回收所得 capacity 加回全局空闲账本
// 宿主已安装其他 root reclaimer 时,以其实际协议为准。
每层 Reclaimer 的职责一目了然:TaskReclaimer 做 pause/resume,ParallelReclaimer 做并发,Operator::MemoryReclaimer 转接到算子,具体 spill 由算子虚函数实现。
10.12.11 生命周期
| 阶段 | 行为 |
|---|---|
| 创建 | pool 创建时通过 addAggregateChild(name, reclaimer) 传入;或后通过 setReclaimer() 单独设置(Operator 走这条路) |
| 运行 | pool->reclaimer() 返回 raw 指针;reclaim 调用都从这里发起 |
| 销毁 | pool 析构时 reclaimer_ 自动析构(unique_ptr) |
| 保活 | Pool 对象保活与 reclaimer 指向的执行对象存活是不同条件。Task reclaimer 使用 weak_ptr/lock 检查 Task;Operator reclaimer 必须结合算子 close 状态和 Task 回收协议避免访问已销毁对象。childPools_ 保活 pool 本身,不能单独证明裸 Operator* 始终有效。 |
10.13 Reclaim 两步设计:reclaim + shrink
这是容易误解的核心设计:
reclaim() → 释放内存(usedBytes 减少),capacity 不变
shrink() → 归还 capacity 给仲裁器,仲裁器才有空闲容量分配给其他 pool
ArbitrationParticipant::reclaim() 完整流程(ArbitrationParticipant.cpp:277):
当前源码摘录(1d1b76567870):velox/common/memory/ArbitrationParticipant.cpp:273。省略外围声明;此片段未作为独立程序编译。
uint64_t ArbitrationParticipant::reclaim(
uint64_t targetBytes,
uint64_t maxWaitTimeNs,
MemoryReclaimer::Stats& stats) noexcept {
const auto minReclaimBytes = std::max(
config_->minReclaimBytes,
static_cast<uint64_t>(capacity() * config_->minReclaimPct));
targetBytes = std::max(targetBytes, minReclaimBytes);
if (targetBytes == 0) {
return 0;
}
uint64_t reclaimedCapacity{0};
try {
ArbitrationTimedLock l(reclaimMutex_, maxWaitTimeNs);
TestValue::adjust(
"facebook::velox::memory::ArbitrationParticipant::reclaim", this);
++numReclaims_;
VELOX_MEM_LOG(INFO) << "Reclaiming from memory pool " << pool_->name()
<< " with target " << succinctBytes(targetBytes);
auto reclaimedBytes =
pool_->reclaim(targetBytes, maxWaitTimeNs / 1'000'000, stats);
reclaimedCapacity = shrink(/*reclaimAll=*/false);
VELOX_MEM_LOG(INFO) << "Reclaimed from memory pool " << pool_->name()
<< " reserved memory " << succinctBytes(reclaimedBytes)
<< ", capacity " << succinctBytes(reclaimedCapacity);
} catch (const std::exception& e) {
VELOX_MEM_LOG(ERROR) << "Failed to reclaim from memory pool "
<< pool_->name() << ", aborting it: " << e.what();
reclaimedCapacity = abortLocked(std::current_exception());
}
return reclaimedCapacity;
}
capacity shrink 委托到 root pool(MemoryPool.cpp:1213):
当前源码摘录(1d1b76567870):velox/common/memory/MemoryPool.cpp:1244。省略外围声明;此片段未作为独立程序编译。
uint64_t MemoryPoolImpl::shrink(uint64_t targetBytes) {
if (parent_ != nullptr) {
return toImpl(parent_)->shrink(targetBytes);
}
std::lock_guard<std::mutex> l(mutex_);
// We don't expect to shrink a memory pool without capacity limit.
VELOX_CHECK_NE(capacity_, kMaxMemory);
uint64_t freeBytes = std::max<uint64_t>(0, capacity_ - reservationBytes_);
if (targetBytes != 0) {
freeBytes = std::min(targetBytes, freeBytes);
}
capacity_ -= freeBytes;
return freeBytes;
}
普通 shrink 的空闲保留目标采用字节阈值与比例阈值的 min,再结合 participant minimum 和实际 reservation 限制可归还量;单次 reclaim 的最小目标才涉及另一组 max 规则。两个公式不能混写,具体数值例子见两阶段配图。
10.13.1 Spill 执行模型:reclaim 线程操刀,driver 旁观
外部回收线程可在 Task 停稳后执行 spill;local self-reclaim 也可在发起仲裁的 Driver 线程上同步执行,并保持该 Driver suspended。Spiller 还可能使用 spill executor 并行写出。因此要区分当前调用线程、Driver 活跃状态以及后续异步子任务。
reclaim 线程主流程:
① task->requestPause() ← 标记 pauseRequested_
② future.wait() ← 等所有 driver 进入 PAUSED (numThreads_=0)
③ pool->reclaim() ← 在 reclaim 线程内直接执行
→ 走 MemoryReclaimer 调用链
→ 找到可 spill 的算子
→ 调用 HashJoin::reclaim() / OrderBy::reclaim() / ...
↓ 这些函数都跑在 reclaim 线程上
→ spiller_->spill() 写盘 + pool->release() 归还 reservation
④ participant->shrink() ← 归还 cap 给 arbitrator
⑤ task->resume() ← 唤醒 paused drivers
安全与否取决于是否位于可回收状态及共享任务是否停稳,不取决于线程是不是最初的 Driver 线程。算子自身可在安全点触发 spill;外部 reclaim 也必须遵守相同状态不变量,不能仅换一个线程就保证一致。
10.13.2 "挑选谁 spill" 的真实粒度:算子,不是 driver
回收沿实际 Reclaimer 子树下沉;基类按 priority 和可回收量排序,ParallelMemoryReclaimer 另有批量并行路径。执行到 Operator 时,还需检查能否回收及共享 peer 状态,不是对全部算子做一次全局“最大者”选择。
TaskReclaimer::reclaim(task_pool, target)
└── for each node_pool in task.childPools_:
ParallelMemoryReclaimer::reclaim(node_pool, target) // 不按 N 简单均分;子任务可超量回收
└── for each op_pool in node.children_:
Operator::MemoryReclaimer::reclaim(op_pool, share)
└── op->reclaim()
例如:HashJoin::reclaim()
SortBuffer::spill()
Window::spill()
多 driver 同 plan node 的算子可通过 spillExecutor 并行 spill——这正是 ParallelMemoryReclaimer 的用途。
10.13.3 并行 spill 安全的两个必要前提
ParallelMemoryReclaimer 把多 driver 的 spill fork 到 spillExecutor 并发执行,两个前提缺一不可:
| 前提 | 缺失会怎样 |
|---|---|
| ① Task 全部 driver PAUSED | driver 正在写数据结构,spillExecutor 线程读到中间状态 → 数据损坏 / segfault |
| ② Per-driver op pool 数据独立 | 即使 pause 了,多线程访问同一份共享数据仍需加锁,并行收益被锁开销吃掉 |
前提 ① 由 TaskReclaimer 保证:在调用 ParallelMemoryReclaimer 之前 task->requestPause() 并等到 numThreads_=0。
不同 Driver 的算子常有独立局部状态,但 HashJoin 等还共享 bridge 与最终表。并行 reclaim 必须根据具体算子的所有权和 peer 协调协议决定,不能仅凭每个 Driver 有一个独立 pool 就证明其 payload 互不共享。
流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。
// ParallelMemoryReclaimer::reclaim 的当前调度摘要。
executor == nullptr -> 调用 memory 基类串行回收
否则:
锁内取得 child 的强引用与可回收量,锁外处理
为非零候选创建传播仲裁上下文的 AsyncSource
除第一个任务外,提交 executor.prepare()
依次 move() 收集结果;第一个任务可由调用线程实际执行
合并每个任务的统计和回收量,传播捕获的错误
syncGuard 消费剩余任务,保证它们结束后才释放引用
// child 的共享状态能否同时回收,仍由具体算子协议决定。
ParallelMemoryReclaimer 可以提交子任务并等待结果;不同实现和有无 executor 的路径可能让调用线程参与工作。SpillerBase 也明确让调用线程执行第一个任务,因此不能把 rcl worker 定义为永远只调度、不做实际 spill。
10.13.4 Driver 内多 op 是串行 spill 的
默认树内顺序遍历会使许多不同 node 的 reclaim 依次发生,但算子可进一步并行处理 peers,自定义 reclaimer 也能改变策略。不存在仅凭“同一 Driver 的不同 Operator”就得出的普遍互斥定理;必须按实际资源共享和回收调用路径分析。
为什么 driver 内多 op 不并行:
- 默认 reclaimer 先取得子树候选并按其优先级/估算量排序,再依次回收直到达到目标或候选耗尽。它可以处理多个 child,不是每次只挑一个 child 就返回。
- target 通常一个 node 就满足(没必要展开成全树扫描)
- 同 driver 的多 op 并行会争抢同 driver 的 CPU 缓存、磁盘带宽
- 大多数 pipeline 一个 driver 只有 1 个 spillable op,优化罕见 case 不值得
10.13.5 完整并行结构总览
Global controller:选择 participant,提交并收集本轮回收结果
reclaim executor:不同 participant 可并行,单个 participant 由 reclaimMutex_ 协调
root / Task 子树:memory 基类按 priority、估算量排序,依次处理并早停
Node:有 executor 时可并行处理可回收 children,调用线程参与首项
Operator:由算子阶段与 peer 协议决定真实回收范围
// 一次 HashBuild reclaim 可能协调多个 peer,不能机械等于一个 Driver 一个表。
10.13.6 Resume:重新入队、栈内恢复与通知
Task::resume(self) 在锁内清除 pause 标记,再区分 Driver 状态:suspended Driver 自己从保留的调用栈恢复;已经入队的不能重复入队;仍有外部 blocking future 的继续等待;其余符合条件的普通 paused Driver 才重新提交到 executor。另有显式注册的 resumePromises_ 在锁外兑现,它们不是所有 paused Driver 共用的恢复机制。
当前源码摘录(1d1b76567870):velox/exec/Task.cpp:1353。省略外围声明;此片段未作为独立程序编译。
void Task::resume(std::shared_ptr<Task> self) {
std::vector<std::shared_ptr<Driver>> offThreadDrivers;
std::vector<ContinuePromise> resumePromises;
SCOPE_EXIT {
// Get the stats and free the resources of Drivers that were not on
// thread.
for (auto& driver : offThreadDrivers) {
self->driversClosedByTask_.emplace_back(driver);
driver->closeByTask();
}
// Fulfill resume futures.
for (auto& promise : resumePromises) {
promise.setValue();
}
};
std::lock_guard<std::timed_mutex> l(self->mutex_);
// Setting pause requested must be atomic with the resuming so that
// suspended sections do not go back on thread during resume.
self->pauseRequested_ = false;
if (self->isRunningLocked()) {
for (auto& driver : self->drivers_) {
if (driver != nullptr) {
if (driver->state().suspended()) {
// The Driver will come on thread in its own time as long as
// the cancel flag is reset. This check needs to be inside 'mutex_'.
continue;
}
if (driver->state().isEnqueued) {
// A Driver can wait for a thread and there can be a
// pause/resume during the wait. The Driver should not be
// enqueued twice.
continue;
}
VELOX_CHECK(!driver->isOnThread() && !driver->isTerminated());
if (!driver->state().hasBlockingFuture &&
driver->task()->queryCtx()->isExecutorSupplied()) {
if (driver->state().endExecTimeMs != 0) {
driver->state().totalPauseTimeMs +=
getCurrentTimeMs() - driver->state().endExecTimeMs;
}
// Do not continue a Driver that is blocked on external
// event. The Driver gets enqueued by the promise realization.
//
// Do not continue the driver if no executor is supplied,
// This usually happens in serial execution mode.
Driver::enqueue(driver);
}
}
}
} else {
// NOTE: no need to resume task execution if the task has been terminated.
// But we need to close the drivers which are off threads as task
// terminate code path skips closing the off thread drivers if the task
// has been requested pause and leave the task resume path to handle. If
// a task has been paused, then there might be concurrent memory
// arbitration thread to reclaim the memory resource from the off thread
// driver operators.
for (auto& driver : self->drivers_) {
if (driver == nullptr) {
continue;
}
if (driver->isOnThread()) {
VELOX_CHECK(driver->isTerminated());
continue;
}
if (driver->isTerminated()) {
continue;
}
driver->state().isTerminated = true;
driver->state().setThread();
self->driverClosedLocked();
offThreadDrivers.push_back(std::move(driver));
}
}
resumePromises.swap(self->resumePromises_);
}
普通 Driver 在主循环检查到 pause 后返回 kPause,退出本次执行;suspended Driver 则走另一条保留栈的恢复路径。下面用流程伪代码并列两者:
// 流程伪代码:普通 pause 与 suspended 不是同一路径。
普通 Driver:
runInternal 检查 shouldStop() -> kPause
CancelGuard 析构 -> Task::leave -> 撤销本次线程登记
Driver::run 收到 kPause -> return
Task::resume -> 按状态选择重新 enqueue
处于 suspended 的 Driver:
调用栈保留,首次 enterSuspended 已从 numThreads_ 扣除
leaveSuspended 若发现 Task 仍 paused,锁外等待后重新检查
最后一层 suspension 退出时恢复计数,再继续原调用栈
Task::resume 不为它重复 enqueue
10.13.7 算子的 spill 状态机:对 driver 透明
Driver 仍使用常规 Operator 接口推进,但 spill 可能改变阻塞原因、内存状态、输入路由和算子阶段。数据恢复细节封装在算子内,不表示调度与运行统计完全看不到这次回收。
| 算子 | spill 前状态 | spill 后状态 | resume 后输出阶段 |
|---|---|---|---|
| OrderBy | RowContainer 内存累积行 | RowContainer 清空 + 磁盘上 sorted run | getOutput() merge-sort:内存剩余 + 磁盘 run 合并 |
| HashJoin (build) | HashTable 内存构建 | HashTable 释放 + 磁盘 partition 文件 | Probe 阶段按 hash 分区,逐 partition 读回 build 再 join |
| HashAggregation | 内存 hash table | hash table 释放 + 磁盘 partitioned spill | getOutput() merge:磁盘 partition 读回再 agg |
| Window | RowContainer 累积分区 | 磁盘上排好序的 partition runs | getOutput() merge-read partition runs |
Driver 只是按平常方式调 getOutput(),算子内部自己处理 spill 数据的读回——这是 spill 状态机透明设计的核心。
10.14 Abort 流程
10.15 超时、Abort 与并发边界
同一 participant 的扩容请求在 waitOps_ 排队;finishArbitration 移交给队首,出锁兑现 promise。Operation 的时间从创建时开始计算,包含这段排队时间;但当前 startArbitration 的排队等待本身没有超时参数,后续拿到执行资格才再检查 deadline。因此 max arbitration time 不是可强制中断任意执行栈的严格墙钟上限。
全局等待有单独的限时 future wait,并在超时路径移除 waiter。额度分配与移除受 stateMutex_ 保护,醒来后还会检查是否实际分得 capacity,以及 abort / timeout。Future 就绪不能统一理解为仲裁成功。
stateLock_ 保护 participant 的 grow / shrink 与请求移交状态;reclaimMutex_ 用于协调同一 participant 的 reclaim / abort;arbitrator 的 stateMutex_ 保护共享额度与全局 waiter。慢速回收不持续持有共享额度锁。通知方通常锁内搬出 promise、锁外兑现,避免 continuation 重入共享锁。
源码:本地排队与移交、operation 时间与状态、全局等待、通知 waiter。
10.15.1 Abort 是停止执行和清理协议,不是立即 free 全部内存
Participant 回收抛异常时会进入 abort 路径;全局策略也可主动选择 victim。Participant 先记录 aborted,调用 root 的 abort,把错误通过 reclaimer 传到 Task,再尽可能收回空闲 capacity。
Task::MemoryReclaimer::abort 设置 Task error 后等待 completion;等待成功才继续向下传播 abort,超时则依赖正在运行的算子走到安全边界后清理。当前常量是 6'000'000 微秒,即 6 秒,旁边的 60s 注释与数值不一致,应以实际参数为准。
因此 aborted 标志、Task 的逻辑终态、算子资源全部释放、root pool 析构是不同事件。不能把 abort 画成抢占并强制 free 所有正在使用的数据。
**ArbitrationParticipant.cpp:343**:
当前源码摘录(1d1b76567870):velox/common/memory/ArbitrationParticipant.cpp:345。省略外围声明;此片段未作为独立程序编译。
uint64_t ArbitrationParticipant::abortLocked(
const std::exception_ptr& error) noexcept {
TestValue::adjust(
"facebook::velox::memory::ArbitrationParticipant::abortLocked", this);
{
std::lock_guard<std::mutex> l(stateLock_);
if (aborted_) {
return 0;
}
aborted_ = true;
}
try {
VELOX_MEM_LOG(WARNING) << "Memory pool " << pool_->name()
<< " is being aborted";
pool_->abort(error);
} catch (const std::exception& e) {
VELOX_MEM_LOG(WARNING) << "Failed to abort memory pool "
<< pool_->toString() << ", error: " << e.what();
}
VELOX_MEM_LOG(WARNING) << "Memory pool " << pool_->name() << " aborted";
// NOTE: no matter query memory pool abort throws or not, it should have been
// marked as aborted to prevent any new memory arbitration operations.
VELOX_CHECK(pool_->aborted());
std::lock_guard<std::mutex> l(stateLock_);
return shrinkLocked(/*reclaimAll=*/true);
}
Abort 触发 Task 停止与清理协议;on-thread 执行、异步引用以及共享 buffer 可能继续短暂存活。Root 的错误状态会影响后续 reservation,实际 buffer 释放、capacity 回收和 pool 析构需分别追踪。
10.16 线程安全与 Race 防范
整个仲裁/回收流程涉及多线程协作(driver / arb / reclaim / spill executor),由 5 层同步原语 + 3 类状态机不变量 + 上下文跨线程传递 三方面共同保证无 race。
10.16.1 五层同步原语
| 层 | 原语 | 类型 | 保护内容 |
|---|---|---|---|
| 进程级 | SharedArbitrator::stateMutex_ |
std::mutex |
freeNonReservedCapacity_、freeReservedCapacity_、globalArbitrationWaiters_、shutdown_ |
| Participant 注册表 | SharedArbitrator::participantLock_ |
folly::SharedMutex |
participants_ map 的读写 |
| Participant 仲裁 | ArbitrationParticipant::reclaimMutex_ |
std::timed_mutex (via ArbitrationTimedLock) |
单个 root pool 上 reclaim/shrink/abort 串行化;local 阶段有超时,global 无超时 |
| Participant 队列 | ArbitrationParticipant::stateLock_ |
std::mutex |
runningOp_、waitOps_ 队列;同 participant 仲裁严格串行 |
| Pool 层 | MemoryPoolImpl::mutex_ |
folly::SharedMutex |
reservationBytes_、capacity_、children_;按 leaf→root 顺序加锁 |
| Task 级 | Task::mutex_ + numThreads_ |
mutex + atomic<int32_t> |
pauseRequested_、resumePromises_、drivers_、driver 调度 |
10.16.2 关键 Race 场景与保护
Race 1:Driver 跑到一半被 spill → 读到不一致的数据结构
D8 正在 SortBuffer::addInput 中改写 RowContainer
reclaim 线程同时调用 pool->reclaim() 触发 spill
保护:双闸门
- Pause 闸门:reclaim 线程在
pool->reclaim()前必须task->requestPause()并等numThreads_==0 - ReclaimableSection 闸门:算子用
nonReclaimableSection_标记临界区
流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。
// SortBuffer.cpp:204
void SortBuffer::addInput(const RowVectorPtr& input) {
ensureInputFits(input);
// ↑ 内部 ReclaimableSectionGuard 打开窗口,只有窗口期 reclaim 可介入
data_->newRow(); // 关窗口后的写操作
}
// Operator.cpp: 实际 reclaim 前检查
uint64_t Operator::MemoryReclaimer::reclaim(...) {
if (op->nonReclaimableSection_) return 0; // 算子明确拒绝
...
}
nonReclaimableSection_ 与 Task pause、锁和回收协议共同限定安全区;原子或 tsan 包装不提供“任何线程立即看到所有相关状态”的保证,也不能替代复合状态同步。
Race 2:Requestor 与 victim 是同一 task → 自死锁
D0 (Q1/T1) growCapacity → global 选中 Q1/T1 为 victim
reclaim 线程要求 Q1/T1 numThreads_ == 0
但 D0 自己还在 numThreads_ 里
保护:requestor 进入仲裁时 enterSuspended() 把自己从 numThreads_ 中减掉:
// 当前 Task.cpp 原样节选:嵌套 suspended、锁内计数、锁外通知与 pause 重检查。
StopReason Task::enterSuspended(ThreadState& state) {
VELOX_CHECK(!state.hasBlockingFuture);
VELOX_CHECK(state.isOnThread());
std::vector<ContinuePromise> threadFinishPromises;
auto guard = folly::makeGuard([&]() {
for (auto& promise : threadFinishPromises) {
promise.setValue();
}
});
std::lock_guard<std::timed_mutex> l(mutex_);
if (state.isTerminated) {
return StopReason::kAlreadyTerminated;
}
const auto reason = shouldStopLocked();
if (reason == StopReason::kTerminate) {
state.isTerminated = true;
return StopReason::kTerminate;
}
// A pause will not stop entering the suspended section. It will just ack that
// the thread is no longer inside the driver executor pool.
VELOX_CHECK(
reason == StopReason::kNone || reason == StopReason::kPause ||
reason == StopReason::kYield,
"Unexpected stop reason on suspension: {}",
reason);
if (++state.numSuspensions > 1) {
// Only the first suspension request needs to update the running driver
// thread counter in the task.
return StopReason::kNone;
}
if (--numThreads_ == 0) {
threadFinishPromises = allThreadsFinishedLocked();
}
VELOX_CHECK_GE(numThreads_, 0);
return StopReason::kNone;
}
StopReason Task::leaveSuspended(ThreadState& state) {
VELOX_CHECK(!state.hasBlockingFuture);
VELOX_CHECK(state.isOnThread());
TestValue::adjust("facebook::velox::exec::Task::leaveSuspended", this);
for (;;) {
{
std::lock_guard<std::timed_mutex> l(mutex_);
VELOX_CHECK_GT(state.numSuspensions, 0);
auto leaveGuard = folly::makeGuard([&]() {
VELOX_CHECK_GE(numThreads_, 0);
if (--state.numSuspensions == 0) {
// Only the last suspension leave needs to update the running driver
// thread counter in the task
++numThreads_;
}
});
if (state.numSuspensions > 1 || !pauseRequested_) {
if (state.isTerminated) {
return StopReason::kAlreadyTerminated;
}
if (terminateRequested_) {
state.isTerminated = true;
return StopReason::kTerminate;
}
// If we have more than one suspension requests on this driver thread or
// the task has been resumed, then we return here.
return StopReason::kNone;
}
VELOX_CHECK_GT(state.numSuspensions, 0);
VELOX_CHECK_GE(numThreads_, 0);
leaveGuard.dismiss();
}
// If the pause flag is on when trying to reenter, sleep a while outside of
// the mutex and recheck. This is rare and not time critical. Can happen if
// memory interrupt sets pause while already inside a suspended section for
// other reason, like IO.
std::this_thread::sleep_for(std::chrono::milliseconds(10)); // NOLINT
}
}
Race 3:同一 participant 并发仲裁 → 重复回收/记账错乱
保护:stateLock_ 队列严格串行:
当前源码摘录(1d1b76567870):velox/common/memory/ArbitrationParticipant.cpp:217。省略外围声明;此片段未作为独立程序编译。
void ArbitrationParticipant::startArbitration(ArbitrationOperation* op) {
ContinueFuture waitPromise{ContinueFuture::makeEmpty()};
{
std::lock_guard<std::mutex> l(stateLock_);
++numRequests_;
if (runningOp_ != nullptr) {
op->setState(ArbitrationOperation::State::kWaiting);
WaitOp waitOp{
op,
ContinuePromise{fmt::format(
"Wait for arbitration on {}", op->participant()->name())}};
waitPromise = waitOp.waitPromise.getSemiFuture();
waitOps_.emplace_back(std::move(waitOp));
} else {
runningOp_ = op;
}
}
if (waitPromise.valid()) {
waitPromise.wait();
}
}
Race 4:Pool 树 reservation 跨层并发更新
D0 在 op_pool.reserve(1024 MiB) → 向 root 传播
D8 同时在 op_pool.release(500MB) → 向 root 传播
保护:incrementReservationThreadSafe 自下而上递归加锁,每层一把锁,顺序固定:
// MemoryPool.cpp:952
void incrementReservationThreadSafe(MemoryPool* requestor, uint64_t size) {
if (parent_ != nullptr) {
parent_->incrementReservationThreadSafe(requestor, size);
// ↑ 先递归到 root(按 leaf→root 顺序加锁)
}
std::lock_guard l(mutex_); // 锁本层
maybeIncrementReservation(size);
}
Reservation 传播通常沿子到父方向进行,但完整回收路径还会涉及 participant、Task 与其他锁。局部的树方向有助于减少锁顺序冲突,不能单独证明整个系统永远无死锁。
Race 5:Reclaim 期间 pool 被并发 abort
保护:reclaim 与 abort 共享 reclaimMutex_,互斥:
// 两个入口的互斥与异常边界摘要,完整实现见本章源码摘录。
ArbitrationParticipant::reclaim:
取得 reclaimMutex_,调用 pool->reclaim,再 shrink(false)
捕获 std::exception 后调用 abortLocked,返回可归还 capacity
ArbitrationParticipant::abort:
取得 reclaimMutex_,调用 abortLocked
abortLocked:
stateLock_ 下设置 / 检查 aborted_;对 pool 传播 abort;最后 shrink
// Pool::reclaim 本身只负责转调 reclaimer,并无通用 aborted -> return 0。
Race 6:Pool 在 reclaim 期间被 query 端析构 → 悬空指针
保护:ScopedArbitrationParticipant 升级 weak_ptr 为 shared_ptr,保活全程:
当前源码摘录(1d1b76567870):velox/common/memory/ArbitrationParticipant.cpp:118。省略外围声明;此片段未作为独立程序编译。
std::optional<ScopedArbitrationParticipant> ArbitrationParticipant::lock() {
auto sharedPtr = poolWeakPtr_.lock();
if (sharedPtr == nullptr) {
return {};
}
return ScopedArbitrationParticipant(shared_from_this(), std::move(sharedPtr));
}
这也呼应 §11.3 提到的 poolId 单调递增——后台仲裁可能延迟 pool 析构,重名时序问题靠 poolId 解决。
Race 7:Thread-local 仲裁上下文跨线程丢失 → 嵌套仲裁死锁
后台 arb 线程持有 kGlobal context
提交 reclaim 任务到 memoryReclaimExecutor_
工作线程接到任务时 thread-local 已丢
→ 工作线程内 spill 又申请内存 → 不识别"在仲裁中" → 嵌套触发仲裁 → 死锁
保护:createAsyncMemoryReclaimTask 显式打包传递:
当前源码摘录(1d1b76567870):velox/common/memory/MemoryArbitrator.h:525。省略外围声明;此片段未作为独立程序编译。
template <typename Item>
std::shared_ptr<AsyncSource<Item>> createAsyncMemoryReclaimTask(
std::function<std::unique_ptr<Item>()> task) {
auto* arbitrationCtx = memory::memoryArbitrationContext();
return std::make_shared<AsyncSource<Item>>(
[asyncTask = std::move(task), arbitrationCtx]() -> std::unique_ptr<Item> {
std::unique_ptr<ScopedMemoryArbitrationContext> restoreArbitrationCtx;
if (arbitrationCtx != nullptr) {
restoreArbitrationCtx =
std::make_unique<ScopedMemoryArbitrationContext>(arbitrationCtx);
}
return asyncTask();
});
}
工作线程内 underMemoryArbitration() 返回 true,maybeIncrementReservation 跳过 capacity 检查,避免嵌套。
Race 8:Free capacity 计数错乱
A 线程:freeNonReservedCapacity_ -= 100MB (划拨给 Q1)
B 线程:freeNonReservedCapacity_ += 200MB (Q2 shrink 归还)
保护:所有读写 free pool 计数的代码路径必须先持 stateMutex_,无例外。
10.16.3 状态机不变量
三个核心状态机消除"不可能状态":
① Task::numThreads_ 与 driver 状态
不变量: (RUNNING drivers 数) = numThreads_
(SUSPENDED 计数) = state.numSuspensions
requestPause 等价于:等待 (numThreads_) == 0
即所有"在跑"的 driver 都让出来
enterSuspended / leaveSuspended 在 Task::mutex_ 内同步更新两计数。
② ArbitrationOperation 状态机
kInit ──startArbitration──> kRunning ──finishArbitration──> kFinished
│ ↑
↓ │
kWaiting ────────────┘ (在 waitOps_ 队列)
stateLock_ 保护 participant 的 runningOp_ 与等待队列;不能说 Operation 的所有 setState 都在该锁内。ArbitrationOperation::start/finish 在各自调用位置推进 kRunning/kFinished,setState 校验允许的前驱状态。安全性需要同时依赖 operation 的执行归属和 participant 的交接协议。
③ ArbitrationParticipant::aborted_ 单调性
aborted_ : false -> true,不可逆
ArbitrationParticipant::abortLocked:
在 stateLock_ 下检查已 abort 则返回 0,否则置位
调用 pool->abort(error),捕获并记录错误
检查 root 已标记 aborted,再 shrink 可归还的空闲 capacity
后续 reservation / 仲裁入口按其路径检查 abort 并传播错误
不能据此声称任意 reclaim() 调用都在第一行立即返回 0
单调性消除了"abort 完一半又取消"的复杂状态。
10.16.4 死锁预防
| 死锁来源 | 防范手段 |
|---|---|
| 锁循环依赖 | 锁顺序固定:pool leaf→root;reclaimMutex_ 内不再获取上层锁;ABBA 避免 |
| Self-pause 死锁 | requestor 进入仲裁前 enterSuspended 把自己从 numThreads_ 减掉 |
| Local 等不到 reclaim | ArbitrationTimedLock 在 local 阶段带超时(默认 ~100ms),失败回退到下一步 |
| 异步 spill 嵌套仲裁 | thread-local context 跨线程传递,underMemoryArbitration() 跳过嵌套触发 |
| Reclaim 与析构竞争 | ScopedArbitrationParticipant shared_ptr 保活 |
future.wait() 永久阻塞 |
时间预算限制等待和回收尝试,超时后的异常/abort 行为按当前操作路径执行;它不会强行中断任意正在运行的 I/O 或无视生命周期立即释放对象。 |
10.16.5 可见性与原子性
- **
nonReclaimableSection_**:tsan_atomic<bool>,配合 release-acquire 内存序,driver 开/关窗对 reclaim 线程立即可见 - **
numThreads_/state.numSuspensions**:std::atomic<int32_t>+ Task mutex 双重保护(atomic 用于 fast-path 读,mutex 用于复合修改) - **
pool->capacity_/reservationBytes_**:pool mutex 保护,checkedGrow/shrink在锁内作为原子复合操作 - **
aborted_/pauseRequested_**:bool+ 持锁修改 + 持锁读取(避开复杂内存序)
10.16.6 整体收敛性
| 保证 | 机制 |
|---|---|
| 数据一致性(spill 读到完整状态) | task pause + nonReclaimableSection_ 双闸门 |
| 互斥性(一个 pool 同时只一个 reclaim) | reclaimMutex_ 单调持锁 |
| 串行性(同 participant 仲裁队列) | runningOp_ + waitOps_ 队列 |
| 死锁自由 | 锁层次 + timeout + 自暂停 + shared_ptr 保活 |
| 跨线程上下文 | thread-local + ScopedMemoryArbitrationContext 打包传递 |
| 状态单调性 | aborted_、ArbitrationOperation 状态机不可逆 |
| 可见性 | atomic + mutex 组合,保证读到最新值 |
不同同步域各保护自己的状态,但部分调用会嵌套锁或在持锁期间调用回收逻辑。分析并发时需要画出真实锁顺序和调用边界,不能假设职责分层就消除了所有多锁场景。
10.17 关键配置参数
10.18 配置与排查:先确认失败在哪一层
以下为当前 SharedArbitrator 的代码默认值,不代表宿主服务的最终配置:
| 配置 | 默认值 | 作用 |
|---|---|---|
memory-pool-initial-capacity |
256MB | 新 root 初始 capacity 目标 |
reserved-capacity / memory-pool-reserved-capacity |
0B / 0B | 全局保留容量 / participant minimum 策略 |
fast-exponential-growth-capacity-limit / slow-capacity-grow-pct |
512MB / 0.25 | 容量增长目标 |
memory-pool-min-free-capacity / memory-pool-min-free-capacity-pct |
128MB / 0.25 | 普通 shrink 后保留的空闲余量,取 min |
memory-pool-min-reclaim-bytes / memory-pool-min-reclaim-pct |
128MB / 0.25 | 单次 reclaim 最低目标,取 max |
global-arbitration-enabled |
true | 启用跨查询已用资源回收 |
global-arbitration-memory-reclaim-pct |
10 | 每轮全局回收目标的比例下限,单位是百分数 |
memory-pool-spill-capacity-limit / memory-pool-abort-capacity-limit |
4GB / 1GB | victim 搜索的起始容量桶 |
memory-reclaim-threads-hw-multiplier |
0.5 | reclaim executor 线程数相对可用 CPU 并发数的比例,至少一个线程 |
max-memory-arbitration-time |
5m | 仲裁时间预算;当前构造检查要求大于 0 |
global-arbitration-abort-time-ratio |
0.5 | 切换 abort 的时间阈值因子,仍需结合回收进展条件 |
global-arbitration-without-spill |
false | 全局路径是否跳过 spill、使用 abort |
增长参数两项、最小空闲参数两项分别要求同时为 0 或同时非零。max-memory-arbitration-time=0 在当前 SharedArbitrator 构造中被拒绝,尽管 header 还留着“0 表示无超时”的注释;全局容量桶参数也要符合 setup 的实际校验。
源码:默认配置、SharedArbitrator 构造校验、participant 配置校验。
| 现象 | 优先核对 |
|---|---|
| Pool cap exceeded | root maxCapacity、requestBytes、量化后的 reservation、自身可回收资源 |
| Arbitrator 资源不足 | 两类 free capacity、participant 的空闲保留策略、global 是否开启、spillable 候选是否足够 |
| Allocator allocation failure | allocator 用量与上限、cache 可驱逐空间、底层分配错误;同时看 RSS,但不混为一项 |
| Global wait 很长 | 控制线程当前回收轮、reclaim executor、victim Task pause 等待、spill I/O、participant 排序 |
| Reclaim 返回 0 | canReclaim、non-reclaimable section、候选估算是否已过时、Task 是否已退出 |
| used 已降但 capacity 没降 | reservation 是否仍被显式保留,shrink 是否发生,minimum / free headroom 是否限制归还 |
| Task 已结束但 root 仍存在 | Task / QueryCtx 外部引用、child 的 parent 引用、scoped participant 与异步回收任务 |
运行时统计中,memoryArbitrationWallNanos、localArbitrationWaitWallNanos、localArbitrationExecutionWallNanos、globalArbitrationWaitWallNanos 帮助区分排队、请求线程上的工作和全局等待;memoryReclaimWallNanos、reclaimedMemoryBytes 反映 reclaimer 的执行。Spill 文件字节和退回的 capacity 仍需分别看,不能用一个指标替代另一个。
本次阅读了 MemoryPoolTest.memoryManagerGlobalCap、ArbitrationParticipantTest.reclaimableFreeCapacityAndShrink、arbitrationOperationWait 等测试源码,用于核对容量限制、shrink 边界和排队行为;只修改文档,没有重新运行 C++ 测试。
| 参数 | 默认值 | 含义 |
|---|---|---|
memory-pool-initial-capacity |
256MB | 新 pool 初始容量 |
memory-pool-reserved-capacity |
0B | 每 pool 保护容量下限 |
max-memory-arbitration-time |
5min | 仲裁总超时 |
memory-pool-min-free-capacity |
128MB | shrink 后保留的最小空闲字节 |
memory-pool-min-reclaim-bytes |
128MB | 低于此量跳过(不值得 spill) |
fast-exponential-growth-capacity-limit |
512MB | 翻倍增长的上限 |
slow-capacity-grow-pct |
0.25 | 超过上限后的增长比例 |
global-arbitration-enabled |
true | 是否启用后台仲裁线程 |
global-arbitration-memory-reclaim-pct |
10% | 每轮至少回收的容量比例 |
global-arbitration-abort-time-ratio |
0.5 | 超过此比例超时后切换到 abort |
10.19 完整仲裁调用链
11. Pool 的父子关系与生命周期
11.1 Pool 的引用关系与释放时机
Pool 树的所有权方向容易读反:children_ 是 weak_ptr map,parent_ 是 shared_ptr。父亲不会仅因列着孩子就保活它;孩子则保证自己的父链仍在。Manager 的 root 跟踪表也是 weak_ptr。
ArbitrationParticipant 长期保存 root weak_ptr;选中候选时,lock() 生成 ScopedArbitrationParticipant,在本次操作期间强持有 root。由此,即使 QueryCtx 已释放自己的引用,仲裁活动或残留 child 也可能继续保活 root。不能把 QueryCtx 说成 root 唯一的强引用来源。
Task 的 childPools_ 持有普通 node / operator pools,customChildPools_ 持有自定义资源子池,裸指针 map 负责查找。Operator close 释放执行资源,但 pool 对象可继续存活,支持跨 Driver 共享的向量和缓冲区所需的生命周期约束。
Task 析构按显式顺序清理 Driver、共享状态、索引、子 pools、task pool 与 QueryCtx。最后一个 pool 持有者释放引用后,pool 析构会从 parent 注销;root 析构则经 Manager destruction callback 调用其 arbitrator 的 removePool,收回剩余 capacity、移除 participant。
Pool 析构包含泄漏检查,不是替调用者扫描并释放任意 malloc 指针的垃圾回收器。释放 buffer、归还 reservation、收回 capacity、销毁 pool 应分别观察。默认 Manager 以及 custom resource 的提供方都需要满足 allocator / arbitrator 的生命周期契约。
源码:parent / children 字段、participant scoped 引用、Manager 弱引用跟踪、Task 析构、pool 析构、root 注销。
11.2 父子存储方式
流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。
// MemoryPool.h
std::unordered_map<std::string, std::weak_ptr<MemoryPool>> children_;
父 pool 用 weak_ptr 持有子 pool——子 pool 的实际生命周期由各自的强引用持有者控制,父 pool 不延长其生命周期。
子 pool 析构时在 MemoryPoolImpl::~MemoryPoolImpl() 中自动调用 parent->dropChild(this),将自己从父 pool 的 children_ map 中移除:
// velox/common/memory/MemoryPool.cpp:469
MemoryPoolImpl::~MemoryPoolImpl() {
if (parent_ != nullptr) {
toImpl(parent_)->dropChild(this); // 自动从父 pool 移除
}
// 检查内存泄漏
VELOX_DCHECK(
(usedReservationBytes_ == 0) && (reservationBytes_ == 0) &&
(minReservationBytes_ == 0));
}
MemoryPool::~MemoryPool() 也断言 children_ 必须为空,确保析构顺序正确。
11.3 Root Pool 的创建与生命周期(Presto Native)
宿主在 worker 内缓存 QueryCtx,缓存命中时复用该对象及 root。缓存失效后允许创建新的 QueryCtx,旧 root 又可能因后台回收保活而暂时存在,所以 pool 名还要加入唯一后缀;不能理解成同一 queryId 永远只创建一次。
流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。
// presto_cpp/main/QueryContextManager.cpp:107
static std::atomic_uint64_t poolId{0};
auto pool = memory::MemoryManager::getInstance()->addRootPool(
fmt::format("{}_{}", queryId, poolId++), // 追加单调 id,保证名字全局唯一
queryConfig.queryMaxMemoryPerNode(),
nullptr);
// pool shared_ptr 转移进 QueryCtx,QueryCtx 是唯一强引用持有者
**缓存结构——QueryContextCache 只存 weak_ptr**:
// presto_cpp/main/QueryContextManager.h:28,71
using QueryCtxWeakPtr = std::weak_ptr<velox::core::QueryCtx>;
queryCtxs_[queryId] = {
folly::to_weak_ptr(queryCtx), // ← 只存 weak_ptr,不持有强引用
queryIds_.begin(), false};
return queryCtx; // shared_ptr 返回给调用方
cache 不计入 QueryCtx 的引用计数,QueryCtx(连同 Root Pool)的实际生命周期由持有 shared_ptr<QueryCtx> 的 Task 决定。
evict() 不强杀存活的 QueryCtx:
// presto_cpp/main/QueryContextManager.h:91
void QueryContextCache::evict() {
// 找 weak_ptr 已失效(QueryCtx 已自然销毁)的 LRU 条目清掉
for (auto victim = queryIds_.end(); victim != queryIds_.begin();) {
--victim;
if (!queryCtxs_[*victim].queryCtx.lock()) { // lock 失败 = 已销毁
queryCtxs_.erase(*victim);
queryIds_.erase(victim);
return;
}
}
// 所有 query 都还活着 → 不驱逐,直接扩容(初始容量 256)
capacity_ = std::max(kInitialCapacity, capacity_ * 2);
}
pool 名字追加单调 poolId 的原因(代码注释原文):
"In some edge case, we found some background activities such as the long-running memory arbitration process will still hold the query root memory pool even though the query ctx has been evicted out of the cache."
内存仲裁过程在仲裁期间会持有 Root Pool 引用,导致 Root Pool 析构延迟。若此时同一 queryId 重新进来,不加 poolId 后缀会与 MemoryManager 中的旧 pool 名冲突。
Root Pool 引用计数来源:
TaskResource(瞬时,调用完即丢)
Task::queryCtx_ ×N(N = 该 query 下的 task 数量) ← 主要强引用
内存仲裁过程(短暂持有)
QueryContextCache 中只存 weak_ptr,不计入引用计数
11.4 Task 管理的 Pool 成员
Task 用三个成员统一管理所有 pool:
流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。
// velox/exec/Task.h:1205
// Task Pool(自身),定义在 childPools_ 之前,析构在后
std::shared_ptr<memory::MemoryPool> pool_;
// 所有 Node Pool 和 Operator Pool 的强引用持有者
// 注释原文:Keep plan node and operator memory pools alive for the duration
// of the task to allow for sharing vectors across drivers without copy.
std::vector<std::shared_ptr<memory::MemoryPool>> childPools_;
// 裸指针索引,仅用于 getOrAddNodePool() 去重,不持有所有权
std::unordered_map<std::string, memory::MemoryPool*> nodePools_;
Operator 只持有 pool_ 的裸指针(通过 operatorCtx_->pool()),所有权始终在 Task::childPools_ 中。
11.5 各层 Pool 的释放时机
| 层级 | close/release 时机 | 真正销毁时机 |
|---|---|---|
| Operator Pool | Driver::closeOperators() → op->close() → pool()->release()(仅归还多余 reservation,pool 对象存活) |
Task::~Task() 中 childPools_.clear() |
| Node Pool | 无独立 close | Task::~Task() 中 childPools_.clear() |
| Task Pool | 无独立 close | Task::~Task() 中 pool_.reset() |
| Root Pool | 无独立 close | 最后一个持有 QueryCtx 的 Task 析构后,QueryCtx 引用计数归零 |
Task::~Task() 中的析构顺序(顺序至关重要):
流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。
// velox/exec/Task.cpp:483-493
CLEAR(drivers_.clear()); // 1. Driver/Operator 析构,但 pool 仍在 childPools_ 中
CLEAR(nodePools_.clear()); // 2. 清空裸指针索引(不触发析构)
CLEAR(childPools_.clear()); // 3. Operator Pool 析构(从 Node Pool children_ 移除)
// Node Pool 析构(从 Task Pool children_ 移除)
CLEAR(pool_.reset()); // 4. Task Pool 析构(从 Root Pool children_ 移除)
CLEAR(queryCtx_.reset()); // 5. QueryCtx 引用计数 -1,若归零则 Root Pool 析构
必须先 childPools_.clear() 再 pool_.reset():子 pool 析构时调用 parent->dropChild(),此时 parent(Task Pool)必须还存活。
11.6 完整生命周期时序
[Task 收到请求]
findOrCreateQueryCtx()
→ MemoryManager::addRootPool("queryId_N", maxCap)
→ QueryCtx 创建,持有 Root Pool shared_ptr
→ cache 存 weak_ptr,返回 shared_ptr 给调用方
[Task::create() → Task::init()]
→ Task::initTaskPool()
→ rootPool->addAggregateChild("task.xxx") Task Pool 创建
[Task::start() → Driver 实例化 → Operator 实例化]
→ OperatorCtx() → Task::addOperatorPool()
→ taskPool->addAggregateChild("node.xxx") Node Pool 懒创建(首次)
→ nodePool->addLeafChild("op.xxx") Operator Pool 创建
→ 所有权压入 Task::childPools_
[Operator 执行完毕 / Task terminate]
→ Driver::closeOperators() → op->close()
→ pool()->release() ← 只归还 reservation,pool 对象存活
→ drivers_[i] = nullptr ← Operator 析构,但 pool 仍在 childPools_
[Task::~Task()]
① drivers_.clear()
② nodePools_.clear() 裸指针索引清空
③ childPools_.clear() Operator Pool 析构 → Node Pool 析构
④ pool_.reset() Task Pool 析构
⑤ queryCtx_.reset() Root Pool 引用计数 -1
[同 query 所有 Task 都析构后]
→ QueryCtx 引用计数归零 → QueryCtx 析构
→ Root Pool 析构 → MemoryManager 注销
12. 值得学习的工程实践:多线程与代码品味
前面九章解释了"做什么"和"为什么这样做",这一章总结这套代码里值得借鉴到自己工程中的具体手法。分两类:多线程编程的具体技巧,和代码品味层面的取舍。
12.1 多线程协议维护的不变量
12.1.1 (1) 锁分层而非锁总线:五种锁各管一摊
整个内存仲裁子系统有五把主要的锁,每一把都有非常明确的辖区:
| 锁 | 保护对象 | 持锁时长 | 谁来持 |
|---|---|---|---|
stateMutex_ (Arbitrator) |
全局 freeCapacity / arbitration 队列 | 极短(毫秒级) | Arb 线程 |
participantLock_ |
单 participant 的 op 队列与状态 | 短 | Arb / Rcl 线程 |
reclaimMutex_ |
单 participant 上的 reclaim 互斥 | 长(整个 spill 周期) | Rcl 线程 |
stateLock_ |
单 op 状态切换 | 极短 | requestor driver |
Pool mutex_ |
单 pool 的 reservation 字段 | 极短 | 所有 driver |
锁按状态域拆分,可以缩小临界区,但当前实现存在嵌套锁:例如 participant 的 stateLock_ 内会调用 pool 的 grow/shrink;基类收集 children 时也会查询 child 的状态。需要核对真实锁顺序、对象析构点与锁外通知,不能把“职责分层”当作无需分析死锁的证明。
实务启示:当一个模块出现"需要同时持两把锁"时,先不要急着定义"锁顺序",而是先问能不能让这两把锁的辖区互不相交。
12.1.2 (2) 量化记账:用空间换并发
reservation 传播按总需求量化:小于 16 MiB 时按 1 MiB,16 至不足 64 MiB 按 4 MiB,更大时按 8 MiB 向上取整。leaf 已有余量时可以省去向父链增量传播;thread-safe leaf 的本地记账仍持锁,不能称整次 allocate 无锁。
// 量化与锁的流程摘要。
allocate(alignedSize)
leaf 本地 reservation 足够:更新 used;thread-safe leaf 仍取本地 mutex
不足:量化新的总需求,再将 increment 传播到父链 / root
quantizedSize(total):
total < 16 MiB -> roundUp(total, 1 MiB)
total < 64 MiB -> roundUp(total, 4 MiB)
otherwise -> roundUp(total, 8 MiB)
maybeReserve(increment) 另先按 8 MiB 对显式预留请求向上取整。
量化余量按 leaf pool 和当前档位产生,不是“每个 Driver 最多 1 MiB”;一个 Driver 还可能有多个算子池。以单个 leaf、小需求档、连续分配且没有并发释放为例,取得 1 MiB 余量可以覆盖约 1024 次 1 KiB 请求,但本地锁、对齐、跨档位与显式最低预留都会改变实际成本,不能写成固定锁次数保证。
批量记账的收益,是让一部分细粒度分配复用已取得的额度,降低跨层同步频率;代价是暂时多占 reservation。类似思想也见于批量 I/O 与缓存,但锁争用是否近似线性下降仍取决于并发、分配大小和热点分布,应由测量判断。
12.1.3 (3) 单调状态机消除中间态 race
aborted_ 是单向终止标记,一个 ArbitrationOperation 也按 Init → Waiting(可选)→ Running → Finished 推进。但 numThreads_ 是反复增减的计数,pauseRequested_ 也会置位和清除;它们不是单调状态机。应区分一次性生命周期状态与可重复进入的执行状态。
单向终态能减少需要处理的恢复分支,但不会自动消除 race。失败的 allocation 仍需回滚 reservation;重复的 pause/resume 仍需互斥与状态检查。设计状态机时,应明确可逆操作的补偿动作、不可逆终态及各个字段的同步域。
新增状态字段时,先判断它表示一次性事件、可重复阶段还是计数,再决定是否允许回退和由谁修改。收益是缩小需要证明的状态空间;不能给出“复杂度降低十倍”这样的未测量数字。
12.1.4 (4) RAII Guard 替代显式 enter/leave
整个仲裁过程不靠裸 enter()/leave() 调用,而是 RAII guard 自动配对:
当前源码摘录(1d1b76567870):velox/common/memory/MemoryArbitrator.cpp:492。省略外围声明;此片段未作为独立程序编译。
MemoryPoolArbitrationSection::MemoryPoolArbitrationSection(MemoryPool* pool)
: pool_(pool) {
VELOX_CHECK_NOT_NULL(pool_);
pool_->enterArbitration();
}
MemoryPoolArbitrationSection 负责配对 enter/leave,ScopedMemoryArbitrationContext 恢复线程上下文;ScopedArbitrationParticipant 则临时保活 participant 和 pool。三者都利用作用域管理生命周期,但不是同一组回调,也不能把保活对象的构造当成进入仲裁。
对应的代码品味:永远不要写裸 enter/leave、acquire/release 对——总能用 RAII 包一层。
12.1.5 (5) weak_ptr + lock() 打破生命周期循环
TaskReclaimer 持 weak_ptr<Task> 而非 shared_ptr<Task>:
流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。
auto task = task_.lock();
if (!task) return 0; // task 已析构,安全降级
task->requestPause();
...
weak_ptr 避免 Reclaimer 长期持有 Task 而形成环;lock 失败可以安全放弃。lock 成功得到的 shared_ptr 会保活 Task,直至本次使用结束,其最后一次释放也可能触发 Task 析构。应把长期观察与操作期间的强持有区分开。
这种设计让回收者按需取得 Task 的生命周期保证;Task、Driver、executor 和 continuation 的实际强引用共同决定析构时机。weak_ptr 本身既不负责停止执行,也不证明所有权只属于某一个对象。
12.1.6 (6) shared_ptr 临时升级保活
ScopedArbitrationParticipant 进入临界区时把 weak_ptr<ArbitrationParticipant> lock 成 shared_ptr 拿在手上:
当前源码摘录(1d1b76567870):velox/common/memory/ArbitrationParticipant.cpp:118。省略外围声明;此片段未作为独立程序编译。
std::optional<ScopedArbitrationParticipant> ArbitrationParticipant::lock() {
auto sharedPtr = poolWeakPtr_.lock();
if (sharedPtr == nullptr) {
return {};
}
return ScopedArbitrationParticipant(shared_from_this(), std::move(sharedPtr));
}
reclaim 过程中即使外部把 Participant 注销,这个 reference 保证它活到 reclaim 完成。生命周期的"借出"和"归还"完全自动化。
12.1.7 (7) 决策线程与执行线程分离:天然背压
Arb 一个线程做决策、Rcl 多个 worker 做执行,两者用 queue 解耦:
requestor → [Arb queue] → Arb 线程:从队头取 → 决策 → push 到 [Rcl queue]
↓
Rcl workers 并行消费
如果 Rcl 堆积(spill 太慢),Arb 自然检测到 queue 长度,新 requestor 就要等更久——背压无需任何特殊代码就出现了。
实务启示:异步任务流水线天然带背压,比"显式监控 + 主动降级"简单且可靠。
12.1.8 (8) atomic 替代 mutex 的边界判断
numThreads_ 是 std::atomic<int>,requestor 进入仲裁前 --numThreads_、离开后 ++numThreads_。没有锁,因为这个字段不需要和其他字段一起 atomic 更新。
但 state + reclaimableBytes 必须一起更新时,仍然用 mutex——不强行 lock-free。判断标准很清晰:
单个字段独立读写 → atomic 多字段必须事务性更新 → mutex
不少代码会过度追求 lock-free 把所有字段都拆成 atomic,结果 ABA 和不一致状态频出。Velox 在这点上很克制。
12.1.9 (9) thread_local 防重入
Arb 线程自己跑 reclaim 时如果 Reclaimer 内部也调 allocate(),会不会再次触发仲裁形成无限嵌套?
流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。
// 当前实现不是 thread_local bool + ALLOC_DIRECTLY。
underMemoryArbitration() == (memoryArbitrationContext() != nullptr)
ScopedMemoryArbitrationContext 保存旧 context 并安装当前 context,析构时恢复
MemoryPoolImpl::maybeIncrementReservation:
root 在仲裁上下文中允许临时超用 reservation
仍执行记账,不跳过 allocator 的实际分配 / 失败处理
异步任务通过 createAsyncMemoryReclaimTask 显式恢复上下文
一行 thread_local 标记就解决了重入。不是用锁、不是用计数器,是用线程身份。
12.1.10 (10) "宽 future + 窄 mutex"协调多线程
requestor 在等仲裁结果时挂在 future 上,没有自旋、没有持锁:
流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。
ContinueFuture future;
{
std::lock_guard lock(mutex_);
futures_.push_back(future);
}
future.wait(); // 不持锁等待
mutex 只用来保护 future 列表的 push,等待动作发生在锁外。这避免了"持锁等待"这个最常见的死锁源。
12.1.11 (11) "读后验证"模式
仲裁完成 capacity grow 后,requestor 回到原始 allocate 路径**重新调一次 maybeIncrementReservation()**,而不是相信仲裁器已经准备好。原因:并发 driver 可能在仲裁期间又消耗了一些。重新验证总是廉价且安全。
实务启示:跨线程的"对方答应给我"绝不直接相信,任何 promise 都要 verify。
12.2 代码品味:可借鉴的设计取舍
12.2.1 (1) 命名精确传递语义
reclaimvsshrink:前者是释放内存、后者是归还 capacity。两步分开、动词不同,读到shrink就知道是账面操作。enterArbitrationvsenterSuspended:前者是 driver 主动进入仲裁、后者是泛用的"暂停计数"。命名暗示了语义层级。numThreads_而非activeCount_:明确说"线程数"而不是某种抽象计数。freeNonReservedCapacity_长但极其精确:哪种 free、是否被 reserved、是 capacity 还是 bytes 一目了然。
坏命名会让多线程问题加倍难调试,因为读者要先猜语义;好命名让代码自解释。
12.2.2 (2) 错误信息把运行时数据放尾部
流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。
// ❌
VELOX_USER_FAIL("Column '{}' is ambiguous", name);
// ✅
VELOX_USER_FAIL("Column is ambiguous: {}", name);
前者把可变值嵌在句子中间,日志聚合时无法对相同错误 grep;后者错误描述固定、数据在尾部,便于 alerting 系统按"静态前缀"分组统计。这是被 CLAUDE.md 显式要求的纪律。
12.2.3 (3) 不取巧的两阶段提交
reservation 是"先记账后真分配",但记账阶段就可能失败(capacity 不够)。失败时整个 try block 抛异常、RAII 自动回滚账。没有 manual cleanup、没有 if-else 树。
流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。
auto rollback = makeScopedRollback();
reserve(); // throws → rollback 自动执行
allocate(); // throws → rollback 自动执行
rollback.commit(); // 成功才显式 cancel rollback
这是 C++ RAII 用于事务性操作的标准范式。
12.2.4 (4) 默认行为与扩展点的边界
memory::MemoryReclaimer 只对 enter/leave 提供空默认实现;priority 返回配置优先级,reclaimableBytes 汇总子树,reclaim 排序并遍历候选,abort 向子树传播且对 leaf 有限制。子类复用这些实际行为,再补充所在执行层的协议。
- 减少了 boilerplate(Root/Task/Node 都不用写空函数)
- 默认行为要按接口分别判断;不支持回收的 leaf 不应被误当成已经实现安全的 spill
- 易于扩展(新增虚函数不破坏现有子类)
代价是"基类接口语义可能不够清晰"——但配合好的文档和命名可以化解。
12.2.5 (5) 数据/行为/逻辑三体分离
Pool ─ 数据:层次结构与生命周期
Reclaimer ─ 行为:统一接口、被仲裁器调用
Operator ─ 逻辑:算子自己的 spill 实现
Pool 不知道怎么 spill、Reclaimer 不写 spill 代码、Operator 不知道仲裁器存在。每方只暴露其他方需要的最小接口。这是"单一职责"在大型系统的标准落地形态。
12.2.6 (6) 路径正交而非分支嵌套
仲裁有三条路径(normal / local / global),但代码不是 if (case1) ... else if (case2) ... else ... 的大分支,而是沿同一条主路径自然降级:
流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。
if (maybeReserveQuick()) return; // path 1
enterLocalArbitration();
if (maybeGrowLocally()) return; // path 2 → returns
enterGlobalArbitration(); // path 3
每段路径独立完整、提前 return,主线代码不嵌套。比起一个 switch + 多个 helper 容易读得多。
12.2.7 (7) 量化数字避免魔法常量
流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。
// 当前参数所在位置及语义,不是三个固定的全局常量。
SharedArbitrator::ExtraConfig::kDefaultMemoryPoolInitialCapacity = "256MB"
MemoryPoolImpl::maybeReserve: kGrowthQuantum = 8 << 20
ArbitrationParticipant::Config: fastExponentialGrowthCapacityLimit,
slowCapacityGrowRatio
// 初始额度、显式预留量化和运行期增长分别属于不同决策。
没有裸 128 * 1024 * 1024 散落代码里。每个魔数提到常量上有名字,改起来一改全改、读起来名字即文档。
12.2.8 (8) 测试钩子的注入位置与构建条件
流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。
TestValue::adjust("facebook::velox::Pool::reserve", &this);
TestValue 注入点的启用与成本取决于构建配置;应查定义与构建选项,不能把所有 release 构建统称为空宏。生产代码无需为测试做特殊设计、测试也不需要 friend class 破坏封装——局部协作文件的偏好不等于整个上游项目禁止 friend;测试可观察接口应按实际封装代价评估。
12.2.9 (9) 没有 util/helper/common 类
整个内存子系统找不到 MemoryUtil、PoolHelper、ArbitrationCommon 这种垃圾桶类。功能都收敛在有明确语义的类里:ArbitrationParticipant、MemoryReclaimer、ScopedArbitrationParticipant。
这是命名上的纪律,但更深层是对模块边界的尊重——一旦容忍 util,所有不知放哪的代码都会堆进来。
12.2.10 (10) 注释解释"为什么"而非"是什么"
流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。
// 作者给出的注释写法示例,并非源码原文。
// Quantize reservation propagation to reduce parent-chain updates.
// Small totals use 1 MiB; larger totals use 4 or 8 MiB quanta.
// Tradeoff: temporarily retained reservation versus cross-pool synchronization.
有用的注释应说明量化传播减少了哪些跨层更新,以及保留余量的代价。上面是依据实现写出的示例,不应冒充源码原注释;修改粒度时还需检查小分配、并发释放、预留峰值与仲裁频率。
12.3 一句话总结
好的并发设计应把可检查的不变量、锁顺序和生命周期责任直接表达出来;它降低出错面,并不意味着错误已不可能发生。
Velox 在 race 上的纪律性体现在:消除可能性(单调状态机)、限制传染面(锁分层)、自动配对(RAII guard)、安全降级(weak_ptr lock 失败)、自然背压(队列解耦)。这些不是"加注意"能做到的,而是设计上就把出错的可能性 designed away。
代码品味上则是反复体现的"克制"——克制使用 lock-free、克制公开接口、克制中间层职责、克制魔法常量、克制 friend 和 util。抽象的价值在于明确职责与不变量,而不是单纯减少接口或说明。
把这两条内化到自己的工程里,往往比记住任何具体技巧都有用。
13. 设计哲学:Pool / Arbitrator / Reclaimer 三位一体
前面八章讲清楚了"怎么做",这一章退一步看"为什么这样设计"。Velox 内存系统的优雅之处不在某个具体机制,而在于三个抽象的职责切分与协同方式:Pool 提供结构、Arbitrator 提供决策、Reclaimer 提供行为,三者正交又紧密咬合。
13.1 Pool 树 ↔ Query 树:同构但不对称
13.1.1 结构同构
默认 pool 树沿 query/task/node/operator 组织资源,但与执行模型并非严格一一对应:多个 Driver 可共享 node pool,connector/exchange 有专用 pool,默认没有 Driver pool 层。对应关系表达归属与回收范围,不是每个执行对象都复制一个独立层级。
Query 执行模型 Memory Pool 树
═══════════════════ ═══════════════
QueryCtx ←→ Root Pool (配额所有者)
└─ Task ←→ Task Pool (生命周期容器)
└─ PlanNode ←→ Node Pool (逻辑分组)
└─ Operator ←→ Op Pool (分配单元)
(Driver-N)
这种同构带来一个直接好处:任何执行层面的对象都能找到对应的内存归属。Driver 调度时拿到 Operator,就能立刻定位到 Operator Pool;Task 析构时,整棵子树跟着一起释放。不需要额外的 ID 映射,不需要 side-table。
13.1.2 但功能严重不对称
结构对称不代表功能对称——这正是设计的精彩处。沿着树的不同位置,职责其实分得很开:
| 层级 | 唯一职责 | 不做什么 |
|---|---|---|
| Root Pool | 持有 capacity 配额,单点 capacity 检查 | 不做分配 |
| Task Pool | 持有 Reclaimer,提供 pause 时间窗 | 不做检查、不做分配 |
| Node Pool | 跨 Driver 的逻辑聚合点,挂 ParallelMemoryReclaimer | 不做检查、不做分配、可能不存在(懒建) |
| Operator Pool | leaf 提供 allocate/free;Operator leaf 是常见执行入口,connector/system leaf 也能分配。enterArbitration 的行为由实际 reclaimer 决定。 | 不持有 capacity、不做检查 |
**职责"两端化"**:分配请求从最底层(Operator)发起,配额检查在最顶层(Root)做,中间两层只做记账(reservationBytes_ 透传)。这是一个标准的"末端发起、顶端裁决、中间透传"模式。
为什么这样切?因为这三件事的频次和粒度完全不同:
- 分配 → 高频、细粒度 → 必须在算子内(贴近物理代码路径)
- 检查 → reservation 增量经过量化后上传,root 集中检查该树的 capacity;中间层仍维护自己的统计与同步。
- 聚合 → 中间层汇总 reservation、峰值等信息,也参与子树回收估算与诊断。
集中 root capacity 检查避免为每一层另建独立硬额度策略;保留 aggregate 层则用于归属、统计和回收路由。路径仍会访问父链及其同步状态,不能把中间层称为完全惰性,也不能从此证明这条路径在所有设计中最短。
13.1.3 懒构造与共享:树是"逻辑形状"而非物理副本
Node Pool 不是 Task 启动时统一创建,而是 Driver 跑到对应算子时才 getOrAddNodePool(planNodeId)。这意味着:
- 从未执行的 plan node 没有 pool:plan 树常常比执行树大(有些 node 被裁剪),懒建避免无效开销
- 同 planNodeId 多 Driver 共享同一个 Node Pool:跨 Driver 的内存压力天然在 node 粒度聚合可见
- HashJoin Build 和 Probe 共享 join node pool:跨 pipeline 的资源关联也能表达
这让 Pool 树不是 plan 树的死板复刻,而是根据实际执行形状裁剪出的最小骨架。结构同构是为了归属清晰,懒建和共享是为了不浪费。
13.2 Arbitrator:三层 Quota 与"决策/执行分离"
13.2.1 三层 quota 的层次美学
进程级 arbitrator.capacity_ ← 整个进程内存预算(系统稳定性)
────────────────────
查询级 maxCapacity_ ← 单 query 最高可用上限(公平性)
────────────────────
当前级 capacity_ ← 此刻实际拥有的额度(性能/弹性)
三层各自的语义目标不同:
- Arbitrator capacity 是其管理范围内的逻辑总额度,不是整个进程 RSS 的物理硬上限;allocator、system pool、cache 和非池化分配还存在各自约束。
- 查询级是契约性的:用户/调度器在提交查询时和系统约定的上限,是公平调度的依据。
- 当前 capacity 可以随 grow/shrink 改变,并降到初始授予值以下。其变化受 maxCapacity、空闲额度、minimum 策略与当前 reservation 等条件限制;初值是启动目标,不是运行期下界。
对一个受管理 root,当前 capacity 不应超过其 maxCapacity;所有受仲裁器管理的授予额度又受全局容量限制。单 root 配置的 maxCapacity 可以高于全局总额度(例如无限上限配置),因此不能把 maxCapacity≤arbitrator capacity 写成始终成立的不变量。
这种分层的优雅在于:把"硬约束"(系统不能挂)、"软约束"(query 不能超)、"调度变量"(capacity 流动)分到三个抽象里,每一层只解决一个问题,决策时只看相邻一层。Arbitrator 几乎从不直接看进程级,只在 query 级判公平、在当前级做流动。
13.2.2 Local vs Global 的分层兜底
仲裁路径不是单一的——它是三条正交路径在效率递减、代价递增的轴上排开:
路径 1:normal(不仲裁) capacity 足够 → 提交 reservation,仍有本地 / 父链同步
↓
路径 2:local(请求线程;可收缩其他 root 空闲额度) arbFree 足够 → maybeGrowFromSelf
↓
撞 maxCapacity → self-reclaim
↓
arb 全局有空 → claim free(账面操作)
↓
路径 3:global(后台线程) 需要别人 spill → enter queue → 后台 Arb 决策 → Rcl 执行
绝大多数分配走路径 1(无开销);负载偏紧时走路径 2(单 driver 短暂 SUSPENDED);只有真的需要回收别人的物理内存时才走路径 3(跨查询协调)。
快慢路径能在部分压力场景下回收和转移额度,但参数仍需为系统用途、峰值中间态与不可回收资源留空间。Spill 和 abort 都有约束与代价,不能依赖它们保证激进配置一定安全或不会 OOM。
13.2.3 Arb 和 Rcl 两个线程的彻底分离
全局仲裁拆出两个 background 线程,是这套设计里"职责切分"最教科书式的范例:
| 维度 | Arb 线程(1 个) | Rcl 线程池(N 个) |
|---|---|---|
| 任务 | 决策 | 执行 |
| 并发模型 | 串行(避免决策竞争) | 并行(per-participant) |
| 视野 | 全局(看所有 participant) | 局部(只看自己的 participant) |
| 代价 | 轻量(排序 + 选择) | 重量(pause + spill + I/O) |
| 失败影响 | 重排队即可 | 单 participant 失败不影响别人 |
为什么必须分开?因为如果 Arb 自己执行 spill:
- 决策被 I/O 拖慢,新进队的 requestor 全部堵塞
- 多个 victim 必须串行 spill(同一线程跑不动)
- 决策日志和 I/O 日志混杂,可观测性极差
控制循环有排序、回收和等待结果的成本,reclaim executor 的并发度受配置、CPU、I/O 与内存约束。拆分职责提供有限并行机会,既不是固定毫秒级延迟,也不允许无限增加并发。
13.2.4 Victim 粒度 = Query
仲裁选择 victim 时按 participant(= query = root pool)排序,而不是 task、不是 operator、不是 driver。这是因为:
- capacity 只在 root pool 一级:选 task 没意义,task 没自己的额度
- 公平性是 query 间的:用户付费、调度配额都在 query 级
- abort 失败兜底也在 query 级:一个 query spill 失败可以 abort 整个 query
同一 root 的多个 Task 共用 root capacity,所以仲裁器按 participant 决定回收范围,再由 reclaimer 在树内选择实际状态。并非“spill 谁都行”:优先级、可回收量、算子阶段、共享表与暂停代价都限制具体选择。
13.3 Reclaimer:与 Pool 同构、与算子解耦
13.3.1 结构同构、功能差异化
Reclaimer 的类层次几乎是 Pool 树的镜像:
Pool 层次 Reclaimer(当前默认 / 常见路径)
Root Pool QueryCtx::MemoryReclaimer;已有宿主实现则保留
查询仲裁状态 + memory 基类子树遍历
Task Pool Task::MemoryReclaimer(继承 exec::MemoryReclaimer)
停稳、取消 / 异常处理、恢复 / 关闭
Node Pool ParallelMemoryReclaimer
有 executor 时组织 AsyncSource,缺省回退串行
Operator Pool Operator::MemoryReclaimer
Driver 保活、回收资格检查、调用算子
Root 上挂哪个 Reclaimer 由创建路径决定。当前 QueryCtx 在 root 尚无 reclaimer 时安装 QueryCtx::MemoryReclaimer,它在回收期间设置 QueryCtx 仲裁状态、用 weak_ptr 获取 QueryCtx,再调用 memory 基类遍历 Task 子树。宿主预先设置的 reclaimer 会保留,因此不能统一写成“root 总是基类实例”。
和 Pool 一样,结构同构 ≠ 功能同构:
- Root 层:当前默认 QueryCtx reclaimer 包装查询级仲裁状态,并复用基类按 priority、可回收量排序的遍历。宿主传入的 root 若已有 reclaimer,则以该实现为准。
- Task 层:保活 Task、请求并等待停稳,处理取消与回收异常,再沿子树执行回收;退出时尝试 resume,已终止时进入相应清理分支。这是一个执行协议边界,不只是调用两个函数。
- Node 层:ParallelMemoryReclaimer 在有 executor 时通过 AsyncSource 组织子池回收,第一个任务由调用线程参与执行,其余可提交并行;没有 executor 时回退到基类串行遍历。收集结果及异常清理同样属于它的职责。
- Operator 层 Reclaimer:只做转接(trampoline)。
reclaim()→op->reclaim()→ 落到 HashBuild/OrderBy/HashAggregation 各自的 spill 实现。
层次划分让不同执行范围拥有明确入口,但每层还需处理估算、回收、取消和生命周期。评价职责是否集中,应看这些行为是否围绕同一执行对象,而不是要求一个类只能做一个动作。
13.3.2 默认遍历与执行层特化
memory::MemoryReclaimer 只对 enter/leave 提供空默认实现;priority 返回配置优先级,reclaimableBytes 汇总子树,reclaim 排序并遍历候选,abort 向子树传播且对 leaf 有限制。子类复用这些实际行为,再补充所在执行层的协议。
流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。
// memory::MemoryReclaimer 的默认行为摘要。
enterArbitration / leaveArbitration : no-op
priority() : 返回 priority_
reclaimableBytes(pool, bytes) : bool;aggregate 汇总子树,leaf 返回 false
reclaim(pool, target, timeout, stats): 返回实际回收量;按 priority / 估算量排序
abort(pool, error) : 向 child reclaimer 传播;leaf 默认不支持
// exec::MemoryReclaimer 另实现 Driver 线程上下文中的 suspended 配对。
只有需要做事的层才重写:
exec::MemoryReclaimer和Operator::MemoryReclaimer都实现 enter/leave;后者使用当前 Driver 线程上下文,并检查请求归属与配置。不能仅从分配 buffer 所属的 Operator 推断当前执行它的 Driver。Task::MemoryReclaimer在 reclaim 中请求 Task 停稳,并用 resume guard 负责退出窗口;pause/resume 是 Task 方法,不是 MemoryReclaimer 基类中等待重写的虚函数。- 并行调度只在
ParallelMemoryReclaimer重写
默认 QueryCtx、Task、Node 与 Operator reclaimer 通过统一内存接口衔接,但保留各层状态约束。Task 负责停稳与错误传播,Operator 负责是否可回收及实际算子调用;子类之间有协议依赖,不能理解为彼此完全不知道对方。
这是 OO 里"虚函数 + 默认实现"用得最干净的一处:继承层次和职责分布完全对齐。
13.3.3 Pool 拥有 Reclaimer,Reclaimer 反指算子
最巧妙的关系是:
Pool ───持有──→ Reclaimer
│
└──反向指针──→ Task* / Operator*
│
└─ reclaim/spill 真实逻辑
- Pool 是结构:提供层次、提供生命周期。
- Reclaimer 是行为:被仲裁器调用,但行为不在自己身上。
- Task/Operator 是逻辑:知道怎么 pause 自己、怎么 spill 自己的数据。
为什么不让 Pool 自己实现 reclaim?因为 Pool 是内存抽象,不应该知道"算子怎么 spill"。HashBuild 的 spill 涉及 hash table、Operator 的 spill 涉及 sort 缓冲,这些是算子领域知识,不属于 memory 模块。
为什么不让 Operator 直接被仲裁器调用?因为仲裁器需要统一的接口(reclaim(targetBytes)),而 Operator 接口五花八门。Reclaimer 作为适配层把"统一接口"翻译成"具体算子调用"。
这是经典的"有结构无行为、有行为无结构、有逻辑无接口"三体分离,每一方都不耦合到其他两方的细节。
13.3.4 weak_ptr + try/catch RAII:生命周期与异常的双闸门
Reclaimer 设计里两个细节体现了对 corner case 的细致考虑:
weak_ptr 打破循环:
Task ──shared_ptr──→ Pool ──unique_ptr──→ Reclaimer ──weak_ptr──→ Task
↑
避免循环引用
reclaim 前用 weak_ptr.lock() 临时保活 Task;若对象已经失效则返回 0。完成后释放这个 shared_ptr,如果它恰好是最后一个引用,就可能在当前线程触发 Task 析构。因此回收代码既要保证操作期间存活,也要考虑最后一次引用释放发生的位置。
try/catch 三明治:
流程化代码节选:省略外围声明与非主线分支;实现位置以相邻固定版本源码链接为准。
// Task::MemoryReclaimer::reclaimTask 的异常路径摘要。
安装 resumeGuard:作用域退出时尝试 Task::resume(task)
等待 task->requestPause() 返回的 future,确认 Task 停稳
若已取消:return 0
try:
memory::MemoryReclaimer::reclaim(task->pool(), target, maxWaitMs, stats)
catch:
task->setError(current_exception()) // 先终止,不让不一致的算子恢复运行
rethrow // 随后由 guard 执行恢复 / 关闭协议
// guard 捕获恢复过程的 VeloxRuntimeError 并记录日志。
RAII 覆盖正常返回和栈展开时的恢复尝试;但要区分“执行了清理回调”与“业务继续运行”。回收异常会先把 Task 标记失败,resume 对终止 Task 清理 off-thread Driver,而不是让失败算子继续计算。
13.4 Pool、Arbitrator 与 Reclaimer 的交接顺序
把三者放在一次 global arbitration 的时间线上看:
Operator Pool Reclaimer Arbitrator
│ │ │ │
│ allocate(bytes) │ │ │
├────────────────────→│ │ │
│ │ check capacity ✗ │ │
│ │ growCapacity() │ │
│ ├──────────────────────────────────────→│
│ │ │ enterArbitration() │
│ │←─────────────────────────────────────┤ (callback)
│ │ │ → op.enterSuspended│
│ │ │ (首次 --numThreads_) │
│ │ │ │
│ [Arb thread] │ │
│ │ │ pickVictim │
│ │ │ ↓ │
│ [Rcl thread] │ │
│ │ ← reclaim() ──────│ pause + spill + │
│ │ │ resume + shrink │
│ │ │ │
│ │ ← grant capacity ──────────────────────│
│ │ leaveArbitration() │ │
│ │ retry allocate ✓ │ │
│←────────────────────│ │ │
每个箭头都是单向单职责:
- Pool 向 Arbitrator 说"我要 X 字节"
- Arbitrator 向 Reclaimer 说"你提供 X 字节"
- Reclaimer 向 Task/Operator 说"暂停并 spill"
- 反向回包同样清晰
没有一个对象需要知道全链路,每个对象只和相邻层对话。这让任何一层都可以独立替换或扩展:换一个 arbitrator 策略不影响 reclaimer、加一个新算子的 spill 不影响 pool 树。
13.5 设计的"轴心"
如果非要用一句话总结这套设计的精神:
结构对齐执行模型,职责沿生命周期分布,行为通过适配层挂接。
- "结构对齐执行模型" → Pool 树映射 Query 树
- "职责沿生命周期分布" → Capacity 在 Root、Reservation 在 Op、Pause 在 Task、Allocate 在 Operator
- "行为通过适配层挂接" → Reclaimer 适配仲裁器接口与算子 spill 逻辑
这套抽象的代价是初次理解需要爬几层;但一旦理解,新增功能(一种新的 reclaim 策略、一种新的 capacity 分配规则、一种新的算子 spill 方式)几乎都是正交插入、零穿越——这是好抽象的最终标准。
14. 核心设计总结
14.1 附录:建议的源码阅读顺序
- MemoryManager、QueryCtx::initPool、Task::initTaskPool:确认 allocator、arbitrator 和 pool 树由谁创建。
- MemoryPoolImpl::allocate、reserveThreadSafe、growCapacity:跟一次申请看余量、量化、root 检查和实际分配。
- SharedArbitrator::growCapacity、ArbitrationParticipant:看请求串行化、self fast path、额度获取和自身上限。
- startAndWaitGlobalArbitration、runGlobalArbitration、spill 任务执行:对照线程时序图,分清同步等待与并行回收。
- Task reclaimer、Operator reclaimer、participant reclaim:沿资源释放、shrink、freeCapacity 返回完整闭环。
- Task 析构、pool 析构、Manager::dropPool:最后验证引用和额度何时真正归还。
| 关注点 | 设计决策 |
|---|---|
| 内存限额粒度 | 以 query 为单位(Root Pool),单 operator 不设独立上限 |
| 分配计量粒度 | 以 operator 为单位(Leaf Pool),可精确定位内存热点 |
| 中间层(Task/Node) | 纯聚合统计 + Reclaimer 挂载点,不参与实际分配 |
| Node Pool 懒建 | 避免 plan node 多而实际不执行时的无效开销 |
| 同 PlanNodeId 多 Driver | 共享 Node Pool,跨 Driver 的内存压力在节点粒度可见 |
| HashJoin 特殊处理 | Build/Probe 侧共享 join node pool,跨 pipeline 聚合 |
| Reservation 量化 | 减少跨层加锁频次,以少量内存膨胀换取更低竞争 |
| Capacity 检查位置 | 仅在 Root Pool 检查,子孙节点无需持锁向上询问 |
| 仲裁触发点 | Root Pool maybeIncrementReservation() 返回 false 时 |
| QueryCtx cache 存 weak_ptr | cache 不延长 QueryCtx/Root Pool 生命周期,生命周期由 Task 持有的 shared_ptr 决定 |
| Operator Pool 晚于 Task 析构 | childPools_ 统一持有所有子 pool,允许跨 Driver 共享 vector 而不拷贝 |
| pool 名追加单调 poolId | 仲裁过程可能延迟 Root Pool 析构,防止同 queryId 重建时名字冲突 |
| 三层 quota 嵌套约束 | root capacity ≤ root maxCapacity;全局授予总量受 arbitrator capacity 约束。配置的 maxCapacity 可以大于全局总量,因此三者并非恒定嵌套不等式。 |
| arbitrator capacity 启动时固定 | MemoryManager 根据 arbitratorCapacity 与 allocatorCapacity 等 options 初始化;宿主可从配置转换这些 options。它是逻辑额度,不预先分配全部物理内存。 |
| 单一 capacity 检查点 | 整条 pool 树只有 Root Pool 一处做 capacity 判断,中间层 Operator/Node/Task 只做记账与传播 |
| capacity 与物理内存正交 | pool capacity 负责仲裁额度,allocator 仍有自己的容量与分配约束。reserve 可以先失败,reserve 成功后物理分配也可能失败并回滚。 |
| addPool 预分配初始 capacity | Root 注册时依据 participantConfig_.initCapacity(当前默认 256MB)、minimum、root 上限和全局空闲额度计算授予;取得的是逻辑 capacity,不是预分配物理内存,也不保证首次分配免于仲裁。 |
| addPool 不足时按现有量给 | 初始授予受 root 上限、participant minimum、普通空闲与 reserved capacity 共同影响;可不足目标或为 0,此时后续申请继续走增长路径。 |
| maxCapacity 强制 self-reclaim | Step 2 ensureCapacity 在 query 撞自己上限时优先逼其自身 spill,spill 不动才报 query OOM;与系统是否有空闲无关 |
| 快速扩容无需 spill | arbFree 足够时,Step 1(maybeGrowFromSelf)直接从 freeNonReservedCapacity_ 划拨,driver 仅短暂 suspended,无跨查询影响 |
| 非线性增长策略 | 增长受 fast/slow growth 参数与 maxCapacity 等边界约束;数字必须按当前配置读取,见前文当前默认值表,不能将旧 512MB/25% 示例作为固定常量。 |
| reclaim 与 shrink 分离 | reclaim 释放内存,shrink 归还 capacity,两步解耦使仲裁器精确控制容量分配 |
| 本地优先,全局兜底 | local 可收缩其他 root 空闲容量,也可 self-reclaim;global 控制循环组织更广范围回收并等待结果。 |
| Local 边界:是否 spill 别人 | Local 可 shrink 其他 query 的"空闲 cap"(账面操作),但不让别人 spill;只有 global 才会暂停别的 task 真正 spill |
| 三线程协作执行 global | Requestor driver BLOCK + 后台 arb 线程决策 + reclaim 线程池执行 spill,三类线程角色清晰分离 |
| Arb 与 Rcl 严格分离 | 控制循环组织候选与目标,reclaim executor 执行可并行回收;控制线程仍等待该轮结果,慢 I/O 会延长回收轮次。 |
| 三层并行结构 | global 控制循环 → participant 回收任务 → 子树/算子 spill 子任务。存在结果等待与有限 executor 容量,调用线程还可参与第一个任务。 |
| Victim 粒度 = query (root pool) | sortAndGroup 按 participant 排序,对应 QueryCtx;capacity、公平性、abort 失败语义都在 query 级一致 |
| 同 query 多 task 贪心串行 | 按 reclaimable 量从大到小逐 task pause→spill→resume;早停退出,少 pause 减小影响面 |
| Rcl worker 全程持锁 | 一个 rcl worker 持 reclaimMutex_ 走完整个 pause→reclaim→resume→shrink 周期,确保 participant 上 reclaim 操作互斥 |
| Reclaimer 绑定 pool 但持算子指针 | reclaimer 是 pool 的成员(pool 拥有),内部反向持 Task*/Operator* 让 reclaim 调用能转到真正的 spill 逻辑——pool 提供结构,算子提供行为 |
| Operator::MemoryReclaimer 不写 spill 代码 | 它只做"转接":调 op->reclaim() 虚函数;具体 spill 由 HashBuild / OrderBy / HashAggregation 等子类各自实现,避免每算子一个 Reclaimer 子类 |
| ParallelMemoryReclaimer 是调度器 | 自己不 spill,只把多 driver 的 spill fork 到 spillExecutor,BLOCK 等 fork-join 完成;真正 IO 跑在 spillExecutor 线程上 |
| 并行 spill 双前提 | Task 停稳、可回收区与共享数据所有权协议同时成立。独立 op pool 不能单独证明 join table 等状态互不共享。 |
| Driver 内多 op 串行 spill | 默认树遍历有顺序和早停规则,但算子可以并行处理 peers;实际互斥由资源所有权和回收协议保证。 |
| 三种让出机制各管一摊 | Suspended 保留栈且退出 Task 活跃计数;普通 pause 返回并等待 resume;普通 future kBlock 归还执行器。仲裁 wait 是另一种同步保栈等待。 |
| Local 不动其他 query | local 可以改变其他 root 的空闲 capacity;self-reclaim 可暂停本 root 相关 Task。是否暂停/释放已用内存需按具体分支判断。 |
| enter/leaveArbitration 四层桥接 | RAII guard → Pool 委托 → Reclaimer 多态 → Task 计数;每层只管自己关注的事,组合起来实现完整协调 |
| Operator reclaimer 和 exec::MemoryReclaimer 等实现分别结合其上下文处理进入仲裁;应按创建/挂接路径确认具体多态分发。 | exec::MemoryReclaimer 通过 driverThreadContext() 找到 Driver 并进入/离开 suspended;没有 Driver 上下文时跳过这一桥接。 |
| enterSuspended 一个原语两种用 | suspended 与 pause 都参与停稳协议,但普通 pause 返回 kPause,保栈路径在 leaveSuspended 配合恢复,不能把两者统一画成 pauseFuture.wait。 |
| TaskReclaimer 是 pause 时间窗的提供者 | 自身只负责"开窗 + 关窗"(requestPause + resume),向下找 children 复用基类贪心逻辑;try/catch 三明治结构保证 resume 永远被调用 |
| TaskReclaimer 用 weak_ptr |
避免 Task→Pool→Reclaimer→Task 循环引用;reclaim 入口 lock() 失败即安全降级返回 0 |
| Pause 粒度是 task 而非 driver | 共享数据结构(如 HashJoin build table)跨 driver 共享,必须整个 task 暂停后 spill 才安全 |
| SUSPENDED ≠ PAUSED | Requestor 自己 SUSPENDED 进入仲裁(不计 numThreads_);victim 的 drivers PAUSED 协作式让出;两者机制独立 |
| Spill 由 reclaim 线程执行 | 外部回收可在 reclaim worker 执行;local self-reclaim 可在 suspended requestor 线程执行,Spiller 还可使用 executor 子任务。安全性取决于状态与停稳协议。 |
| Spill 状态机对 driver 透明 | 算子内部 spiller_ 自动切换"纯内存"和"merge-spill"模式,driver 不感知 spill 历史 |
| 双闸门保证 spill 数据一致 | task pause(大粒度时间窗)+ nonReclaimableSection_(算子级临界区标记),缺一不可 |
| 同步原语五层分工 | 各锁保护不同状态域,但会有嵌套获取和持锁调用。需检查真实锁顺序,职责分工不能证明从不同时持多锁。 |
| 状态机单调性 | aborted_ 不可逆、ArbitrationOperation transition 固定,消除"半途回滚"的复杂状态 |
| shared_ptr 保活 reclaim 期间的 pool | ScopedArbitrationParticipant 升级 weak_ptr 为 shared_ptr,防止 pool 在 reclaim 中被并发析构 |
| Spill 优先,Abort 兜底 | 是否切换 abort 由当前超时比例、回收进展与配置判断;victim 还受优先级分组、容量阈值和 participant ID 等排序规则影响。 |
| 同 participant 串行仲裁 | ArbitrationParticipant 内部排队,避免同一查询并发仲裁导致过度回收 |
| 防重入 thread-local 标记 | 仲裁线程自身申请内存时跳过 capacity 检查,防止嵌套仲裁死锁 |
15. 从资源控制的反馈过程看分层与回收安全
Pool、Arbitrator 和 Reclaimer 分工的依据,是它们掌握的信息不同。Pool 知道资源归属和用量;Arbitrator 能比较参与者与可用额度;执行对象才知道哪些状态可以丢弃、落盘或重建。把所有判断集中在 allocator 中,既缺少算子语义,也难以解释跨查询的资源选择。
| 决策 | 设计收益 | 代价与适用边界 |
|---|---|---|
| 将额度管理与实际分配分开 | 可以在资源归属层协调使用权,不把每次额度转移等同于系统内存操作 | 需要同时观察 used、reservation、capacity 与 allocator 状态 |
| 以 root 为仲裁参与者,以子树组织记账和回收 | 查询级资源决策可以沿执行对象层次落实 | pool 的资源层次与对象生命周期并不完全相同,不能靠名字代替所有权分析 |
| 先寻找未使用 capacity,再考虑实际回收 | 优先减少搬移数据和中断执行的成本 | 未使用额度可能不足;有效回收量、等待时间和系统压力仍决定后续动作 |
| 通过 Reclaimer 进入算子回收 | 内存系统复用统一接口,算子保留语义控制 | 需要 Task 停稳、可回收区及算子阶段协议,统一接口不会消除这些前提 |
| 用批量 reservation 和有限粒度回收 | 减少频繁跨层协调的成本 | 粒度过大可能占住暂时不用的额度,过小又增加同步和回收开销 |
安全性不来自“加了一把锁”或“等了一个 future”本身,而来自动作之间的顺序:共享状态先稳定,相关执行活动再停稳,然后回收修改状态,最后恢复执行。两阶段的关键在于建立可观察的安全条件,而不是把一个函数机械拆成两个函数。
工程上应分别记录申请等待、回收用时、释放 used bytes、归还 capacity 和失败原因。只有把这些量分开,才能判断瓶颈是在额度策略、执行停稳、spill I/O,还是不可回收工作集。设计品味体现在这些边界能否被说明、检查和观测,而不是接口数量少或层次图整齐。
源码核对(2026-09-20):本轮按 Velox 1d1b76567870 核对关键接口、控制流、默认值与边界条件。当前源码摘录附固定版本链接;流程伪代码用于说明分支,不是可直接编译的程序。未对全文示例做独立编译或性能复测。涉及宿主集成与历史实验的数据,按各节标注的来源理解。
- velox/common/memory/ArbitrationOperation.cpp:60
- velox/common/memory/ArbitrationParticipant.cpp:117
- velox/common/memory/ArbitrationParticipant.cpp:118
- velox/common/memory/ArbitrationParticipant.cpp:217
- velox/common/memory/ArbitrationParticipant.cpp:273
- velox/common/memory/ArbitrationParticipant.cpp:273,338
- velox/common/memory/ArbitrationParticipant.cpp:310
- velox/common/memory/ArbitrationParticipant.cpp:345
- velox/common/memory/MemoryArbitrator.cpp:216
- velox/common/memory/MemoryArbitrator.cpp:234
- velox/common/memory/MemoryArbitrator.cpp:467
- velox/common/memory/MemoryArbitrator.cpp:492
- velox/common/memory/MemoryArbitrator.h:525
- velox/common/memory/MemoryPool.cpp:1010
- velox/common/memory/MemoryPool.cpp:1035
- velox/common/memory/MemoryPool.cpp:1222
- velox/common/memory/MemoryPool.cpp:1244
- velox/common/memory/MemoryPool.cpp:939
- velox/common/memory/MemoryPool.h:550
- velox/common/memory/SharedArbitrator.cpp:1167
- velox/common/memory/SharedArbitrator.cpp:452
- velox/common/memory/SharedArbitrator.cpp:452,1134,1167
- velox/common/memory/SharedArbitrator.h:64
- velox/core/QueryCtx.cpp:126,155
- velox/core/QueryCtx.h:441
- velox/exec/HashBuild.cpp:1314
- velox/exec/MemoryReclaimer.cpp:143
- velox/exec/MemoryReclaimer.cpp:32
- velox/exec/MemoryReclaimer.cpp:85
- velox/exec/Operator.cpp:682
- velox/exec/Task.cpp:1353
- velox/exec/Task.cpp:1353,3522
- velox/exec/Task.cpp:1353,3842
- velox/exec/Task.cpp:3864
- velox/exec/Task.cpp:3928
- velox/exec/Task.cpp:732
- velox/exec/Task.cpp:757
- velox/exec/Task.cpp:913
- velox/exec/Task.h:999
贯穿分配例子的三个边界也给出了诊断顺序:先看 root 是否能提交 reservation,再看 arbitrator 是否能够转移 capacity,最后看 allocator 能否取得存储。若把这三种失败统一写成“内存不足”,就会把配额调整、spill 策略和 allocator / cache 问题混在一起。保留分层计数增加了账本维护和并发重查的成本,却让资源政策不必嵌入每个 malloc 调用,也让算子只需暴露自己能够安全回收的状态。