Macduan Notes

向量化 - A Velox View

Velox 是一个可组合的 C++ 数据执行库。它用 Type 描述数据语义,用 Vector 在内存中组织一批列式数据,用 Function 和 Expression 完成值计算,再由 Operator、Driver、Task 把这些计算接成可调度的流水线。MemoryPool、Arbitrator、Allocator 和 Spiller 负责让这些流水线在有限资源下运行。

理解 Velox,可以沿一批数据的路径展开:数据源怎样生成 Vector,函数怎样访问编码后的值,表达式怎样减少重复工作,算子怎样交换 RowVector,Driver 怎样等待和恢复,以及状态超出内存预算后怎样落盘与重建。下面从这些对象的职责出发,再进入实现与设计取舍。

1. Velox 的职责与边界

Velox 所采用的列式表示与批量执行,来自数据库执行技术长期的演进。KDB+、MonetDB、VectorWise,以及后来的 BigQuery、ClickHouse、Redshift、Snowflake、Photon、DuckDB、DataFusion,都可以放在这个技术背景下理解。但这些系统的产品形态、优化器、执行模型和实现语言并不相同,不能因为都处理列式数据,就把它们归为同一个架构。

其中需要区分 整列执行向量批量执行:MonetDB 的列式思路推动了整列运算,MonetDB/X100 与后来的 VectorWise 则强调适当批量大小,在函数调用开销、CPU cache 和中间结果物化之间取得平衡。Photon 采用原生向量化执行来服务 Spark 工作负载;DuckDB 是嵌入式分析数据库;Velox 则把可复用的执行组件做成 C++ 库。

原分享用 CWI、MonetDB、VectorWise 与这些系统的联系帮助建立技术脉络。这里保留这层背景,但不把人员经历、商业竞争猜测或示意箭头当成代码继承关系;也不使用“第一个向量化数据库”这样的绝对断言。

列式执行的技术脉络
图 1 · 列式执行的技术脉络 打开原图

Velox 是一个面向数据处理的 C++ 执行组件库。它提供类型系统、列式向量、表达式求值、函数库、关系算子、Connector、序列化和资源管理,供不同数据系统组成自己的执行引擎。它既可以执行一个本地计划片段,也可以只被使用其中的表达式或向量组件。

Presto、Spark 等集成场景能复用这些执行机制。原文提到的 PyTorch、XStream、F3、FBETL 等,是技术分享及早期介绍中的应用背景,不代表所有这些集成的实现都在当前开源仓库里。开源代码能直接核对的是组件接口、实现和对应测试。

Velox 不负责解析 SQL、DataFrame 等语言,也不承担跨机器的全局查询优化和任务放置。宿主通常把已经选择好物理执行策略的计划转换为 Velox 的 core::PlanNode / PlanFragment,并提供 Connector、函数注册、执行线程池与查询配置。Velox 再把该片段组织成本机的 Task、Pipeline 和 Driver。

“宿主负责控制面,Velox 负责数据面”是有用的职责划分,但不是绝对隔离:Velox 内部仍然包含本地调度、I/O 等待、背压、内存仲裁和自适应策略。它不做全局优化,并不意味着执行过程只能机械地照着计划运行。

源码:PlanNode.hPlanFragment.hLocalPlanner.cppTask.h

Velox 位于数据系统的执行层
图 2 · Velox 位于数据系统的执行层 打开原图

2. 执行引擎由哪些组件组成

  • Type:表示标量、复杂和嵌套类型,例如 BIGINT、VARCHAR、ARRAY、MAP、ROW;参数化类型还携带精度、小数位、字段名和子类型。Tensor 不是这里可以直接列为通用内建 TypeKind 的基本类型;张量场景需要结合数组、扩展类型或上层约定表达。

  • Vector:在内存中表示列式数据,常见实现包括 Flat、Dictionary、Constant、Array、Map、Row、Lazy;编码枚举还包含 BIASED、SEQUENCE、FLAT_MAP、FUNCTION 等。枚举存在不等于所有执行路径都同等支持。它与 Arrow 的列式思想兼容,并提供转换接口,但不能把所有 Velox buffer 布局视为 Arrow 的逐字节同义物。

  • Expression Eval,基于向量化编码(Vector-encoded)数据构建的,完全向量化的表达式求值引擎(expression evaluation engine),它借助了common subexpression elimination, constant folding, effective null propagation,encoding-aware evaluation和dictionary memoization等技术。
  • Functions,Velox函数库提供了与流行的 SQL 方言兼容的函数包(目前,用于 Presto 和 Spark),提供了simple function interface方便开发人员快速增加新的函数。
  • Operator:实现扫描、过滤投影、聚合、HashJoin、排序等关系计算。计划中的 Filter / Project 常由 FilterProject 组合执行;计划节点和物理算子类不要求一一同名。

  • Connector 与 DWIO:Connector 负责数据源接口和列 / 谓词适配,DWIO 和各格式 reader / writer 负责文件编码与解码,文件系统适配负责存储访问。常见格式包括 ORC / DWRF、Parquet、Nimble,存储适配可覆盖 S3、HDFS 等。具体可用能力取决于编译选项、注册和配置。

  • Serializer:把向量转换为传输或持久化需要的字节表示,例如 PrestoPage、Spark UnsafeRow 相关实现。序列化接口不是完整的网络传输系统;连接、重试和跨节点调度仍由相应运行时或适配层承担。

  • Resource Management,用于处理计算资源的原语集合,例如内存区域(memory arenas),缓冲区管理(buffer managment),tasks,drivers(这里的driver指velox的pipeline的driver,非spark driver),线程池,spilling和cache。
  • Fuzzer / Replayer:围绕表达式、聚合、算子和执行追踪构建随机测试、差分验证与复现工具。优化需要在这些语义与异常测试下成立,不能只验证一个 happy path。

2.1 类型、向量编码与 C++ 类型分别负责什么

从 I/O 到计算,完整的 Velox Type 一直保留。计划节点的输出 schema、表达式的结果类型、Operator 的输入输出类型,以及 Vector 的 type(),都描述逻辑语义与结构;它们不会在开始计算后被全局替换成一个 C++ 类型。

具体实现再通过 TypeKind、TypeTraits、模板分派和 Reader / Writer 把这种语义映射到可运算的数据。例如 BIGINT 对应 int64_t;DATE 有自己的逻辑语义,但 DateType 继承整数实现,复用 INTEGER / int32_t 的底层表示。ARRAY、MAP、ROW 也是 Velox Type,同时各自携带子类型结构。

因此,逻辑类型、物理实现分类、C++ native type 并不是互斥的三份类型清单。BIGINT 这样的基础类型同时承担逻辑类型和实现分派分类;DATE、DECIMAL 等则更能体现语义与存储的分离。Vector encoding 又是另一条维度:一个 VARCHAR 列可以是 Flat、Constant 或 Dictionary。

从逻辑类型到 C++ 实现:类型不会因编码而消失
图 3 · 从逻辑类型到 C++ 实现:类型不会因编码而消失 打开原图

源码:Type.hType.hType.hType.hVectorEncoding.h

3. Vector 如何表示与访问数据

Vector 是各组件交换列式数据的基本对象。BaseVector 保存逻辑行数、完整 Type、encoding、可选 null bitmap 和内存池等元信息,并提供复制、调整大小、hash、比较和打印接口。列式布局把同一字段的数据聚在一起,让批量访问、压缩表示和 SIMD 更容易发挥作用。

“Arrow-compatible”应理解为可以在相近的列式模型之间互操作,而非所有对象可直接互换。例如 Velox ArrayVector 分别保存 offsets 与 sizes,StringView 还携带内部指针;导入 / 导出是否零拷贝必须看具体类型、编码、ownership 和转换实现。

LazyVector 把加载推迟到调用方真正需要值的时候,并允许指定要加载的行集合。对于过滤、Join 条件和投影,如果某列只有少量行最终需要,延迟加载可以减少解码和物化;如果根本不需要该列,有时还能避免相应数据读取。

这个收益不能直接等同于“每少选一行就少一次 S3 / HDFS I/O”:文件页、压缩块、预读和缓存决定真实 I/O 粒度。LazyVector 与产生它的 reader 还有生命周期约束,底层 reader 推进后不能随意加载旧批次;Driver 的背压逻辑也需要维护这个边界。

DecodedVector 提供统一的只读访问视图:沿 Dictionary / Constant 等包装找到底层向量,组合行映射与 null 信息,让调用方用 index(row)isNullAt(row)valueAt<T>(row) 访问逻辑值。它通常不会复制成一个与原向量等长的 FlatVector。

Flat、Constant 和单层 Dictionary 的常见路径能复用输入数据与索引;多层 Dictionary 可能需要分配组合后的索引,null 合并、显式访问 indices() 等操作也可能物化内部数组。复杂值的 base 可以仍是 ArrayVector / MapVector / RowVector,不能概括为“任何编码都转成 FlatVector”。

默认 decode 还可能触发 LazyVector 加载;loadLazy=false 用于特定延迟访问路径,例如聚合通过 ValueHook 在加载时直接更新状态。DecodedVector 的内部临时状态使用系统分配器并可复用,它不是一份自动计入查询 MemoryPool 的独立输出列。

源码:BaseVector.hLazyVector.hDecodedVector.hDecodedVector.h

3.1 FlatVector:值与 null 位图

普通定宽标量的 FlatVector 保存 values buffer 和可选 null buffer。null bitmap 按位存储在机器字中:0 表示 NULL,1 表示非空;没有 null buffer 表示全非空。INTEGER 的值可视为连续的 int32_t,但 BOOLEAN 的 values 本身也采用 bit-packed 形式,不能把所有 FlatVector 都解释为普通 T[]

当某行是 NULL 时,values 槽可能包含旧字节或任意占位值;语义上必须先处理 null,不能因为读取到了某个整数就把它当成该行的值。下图保留原例的 12 行数据和三个 NULL 槽。

FlatVector:值缓冲区与 null 位图
图 4 · FlatVector:值缓冲区与 null 位图 打开原图

3.2 ArrayVector:偏移、长度与元素

ArrayVector 保存 offsets、sizes、可选 null bitmap,以及一个存放元素的 child Vector。第 i 行引用 child 中的 [offsets[i], offsets[i] + sizes[i]) 范围。child 可以是 Flat、Dictionary、Constant,或另一个复杂 Vector,不要求元素先全部展开。

原例四个数组分别是 [1,2,3][5,8][10,12,-1,0][123,9];对应 offsets 为 [0,3,5,9],sizes 为 [3,2,4,2]。offsets 不要求单调,所以可以重排和引用已有元素区间。父数组是否 NULL、数组长度是否为零、数组内某元素是否 NULL,是三个不同问题。

ArrayVector:行范围引用共享 elements
图 5 · ArrayVector:行范围引用共享 elements 打开原图

3.3 DictionaryVector 与 ConstantVector

DictionaryVector 通过 indices 引用底层 values Vector,并可带自己的 null bitmap。它既能表示重复值,也能表示过滤后的选择和排序后的重排:创建一组新索引即可改变逻辑顺序,无需复制底层 payload。底层不必已经去重,因此“dictionary”是编码结构的名称,不是对基数的保证。

ConstantVector 表示逻辑上重复 N 次的同一个值。对常见标量路径,它可以直接保存 value_ 和 null 状态;对 ARRAY 等复杂值,则通过 valueVector_ + index_ 引用底层的某一行。原例的复杂常量引用 ArrayVector 第 2 行,重复表示 [10,12,-1,0]

DictionaryVector:用索引重排、筛选与复用值
图 6 · DictionaryVector:用索引重排、筛选与复用值 打开原图
ConstantVector:标量与复杂值两条表示路径
图 7 · ConstantVector:标量与复杂值两条表示路径 打开原图

源码:DictionaryVector.hConstantVector.hConstantVector.h

3.4 MapVector 与 RowVector

MapVector 和 ArrayVector 都使用 offsets / sizes 表示每行范围,但 MapVector 有平行的 keys 与 values 两个 child Vector:相同位置的 key 和 value 组成一对。原例仍使用 11 个 entry,分成四个 map;它们的父 null 与 value null 需要分别处理。

RowVector 则按字段保存 children,父层的第 i 行由每个 child 的第 i 行共同组成。ROW(INTEGER,VARCHAR,VARCHAR) 可以表示数量、颜色和形状三列:整行 NULL 不等于“每个 child 都被写成 NULL”,非空行内部也允许单个字段为 NULL。Operator 输出一个 RowVectorPtr,就是按这种结构传递一批关系行。

MapVector:同一范围索引两组 children
图 8 · MapVector:同一范围索引两组 children 打开原图
RowVector:列式 children 组合一批关系行
图 9 · RowVector:列式 children 组合一批关系行 打开原图

3.5 StringView:固定描述符与变长字符

Velox 的 StringView 属于常说的 German-style string 表示:固定长度的描述符保留长度与短前缀,短字符串内联,长字符串引用外部字符数据。DataFusion / Arrow 的 StringViewArray 也采用相近方向,但具体布局和所有权协议不能不加区分地互换。原文推荐的两篇 Using StringView / German Style Strings to Make Queries Faster 分别讨论 Parquet 读取和字符串操作,是理解这种布局收益的补充资料。

结构示意,不是重新定义 Velox StringView

// StringView 的字段布局示意;接口、构造和比较逻辑见源码。
struct StringViewLayout {
  uint32_t size;
  char prefix[4];
  union {
    char inlined[8];
    const char* data;
  } value;
};

源码:StringView.hStringView.h

StringView:16 字节描述符与独立字符串缓冲区
图 10 · StringView:16 字节描述符与独立字符串缓冲区 打开原图

FlatVector<StringView> 的 values buffer 由 16 字节描述符组成。长串保存 4 字节长度、4 字节前缀和 8 字节指针,指针指向完整字符串;长度不超过 12 字节的字符串直接占用 prefix 与 union 的空间。这里计算的是字节数,不是 Unicode 字符数。

向量的 stringBuffers 维持外部字符串内存的所有权,而 StringView 本身不是 owning string。复制一个长串 StringView 只复制了描述符;如果没有同时维持字符缓冲区生命周期,指针会悬空。下图把原例中的 orange peel 长度修正为 11 字节,其余短串与长串例子仍保留。

  1. 长度不超过 12 个字节的短字符串完全内联在StringView的value buffer中,这样可以增加内存局部性,cache locality更好。
  2. 比较长度和前缀可以快速排除不同字符串,减少访问长串 buffer 的机会;长串长度和前缀相同后仍需比较剩余内容。相等比较与字典序比较的短路条件不同,实际比较次序以 StringView 的操作符实现为准。

  3. trim / substr 等操作在满足实现条件时,可以让输出 StringView 引用输入字符的一段,避免复制长串字符。Adapter 会根据函数声明获取相应输入 stringBuffers 的共享所有权;短串可能直接内联。零拷贝是特定路径的优化,不是这些函数对所有输入的无条件保证。

  4. 由于 StringView 具有固定大小(16 个字节),因此可以乱序写入 StringView Vector(例如,先在位置 2 写入 StringView,然后是 0 和 1)。

源码:StringView.hStringView.hSimpleFunctionAdapter.h

3.6 SelectivityVector:计算哪些行

当前类名是 SelectivityVector,不是 SelectiveVector。它保存选中行的 bitmap 和有效边界,在过滤阶段间传递活动行集合。后续谓词只处理仍被选中的行,无需每一步都重新整理物理数据;全选时还有连续区间的快速遍历路径。它表达“哪些行要计算”,并不表达“这些行是否为 NULL”。

另一种常见表示是 position list,直接保存活动 row number。原文以 Photon 的 position-list 设计作比较,并推荐论文 Filter Representation in Vectorized Query Execution。稀疏选择时 list 紧凑、遍历直接;bitmap 易于按机器字组合过滤条件,并避免高选择率下保存大量 row IDs。两者都有场景,不能脱离选择率、批量大小和后续访问模式判断优劣。

下图保留原文的两次过滤和 position-list 例子,并显式区分位极性:Velox null bitmap 的 0 表示 NULL;原 position-list 示意图里的 isNull 标志则用 1 表示 NULL。position list 选中某行,也不意味着该行的所有列都非空。

SelectivityVector:选中行与 null 是两件事
图 11 · SelectivityVector:选中行与 null 是两件事 打开原图

源码摘录:全选 / 稀疏遍历 · SelectivityVector.h

template <typename Callable>
inline void SelectivityVector::applyToSelected(Callable func) const {
  if (isAllSelected()) {
    // IMPORTANT: Do not remove this line.  Without it the compiler would not be
    // able to vectorize or unroll the loop.
    const auto end = end_;
    for (vector_size_t row = begin_; row < end; ++row) {
      func(row);
    }
  } else {
    bits::forEachSetBit(bits_.data(), begin_, end_, func);
  }
}

4. Function 如何接入批量执行

Velox 提供两种主要的函数实现方式。VectorFunction 直接接收行集合、输入 Vector、输出 Type 和 EvalCtx,适合需要控制批量布局、共享结果或特殊执行路径的实现;Simple Function Interface 让作者主要编写逐值计算,由 Adapter 屏蔽大量编码、null、读写视图和结果管理细节。

框架能够选择 null-free 路径、ASCII 路径,预处理常量参数,并让支持该语义的字符串函数复用输入 buffer。这里的“能够”仍有明确条件:函数必须正确声明 null 行为、确定性和字符串复用来源,输入也必须满足相应路径的约束。不能因为某个函数提供了 callNullFree,就对含嵌套 NULL 的复杂输入无条件使用它。

Simple Function 如何进入向量化执行框架
图 12 · Simple Function 如何进入向量化执行框架 打开原图

以两个 DOUBLE 相加为例,手写 VectorFunction 至少需要处理结果可写性、选中行、编码和 null。原文的 PlusDouble 片段存在未声明变量、大小写错误、分支重复落入通用路径等问题,不能作为可直接使用的实现。下面保留“手写批量接口”的教学目的,给出一个更明确的通用路径;假设函数已注册为两个 DOUBLE 输入、一个 DOUBLE 输出。

教学示例:通用编码访问与 null 处理

// 接口演示;注册与 include 省略,未作为完整 Velox 目标编译。
class PlusDouble : public exec::VectorFunction {
 public:
  void apply(
      const SelectivityVector& rows,
      std::vector<VectorPtr>& args,
      const TypePtr& outputType,
      exec::EvalCtx& context,
      VectorPtr& result) const override {
    // 先建立只读输入视图,再准备输出,避免假定输入必定为 Flat。
    exec::DecodedArgs decoded(rows, args, context);
    const auto* left = decoded.at(0);
    const auto* right = decoded.at(1);
    context.ensureWritable(rows, outputType, result);
    auto* output = result->asFlatVector<double>();
    auto* values = output->mutableRawValues();

    rows.applyToSelected([&](vector_size_t row) {
      const bool isNull = left->isNullAt(row) || right->isNullAt(row);
      output->setNull(row, isNull);
      if (!isNull) {
        values[row] = left->valueAt<double>(row) +
            right->valueAt<double>(row);
      }
    });
  }
};

这里刻意没有手工接管输入 buffer 的 ownership。result 非空不代表可写:它可能是只读编码,也可能已经含有条件表达式其他分支写入的行。函数只能更新 rows 指定的行,并借助 EvalCtx 的可写性协议保留其余结果。

生产实现还可以增加两类常见 fast path:两个 Flat 输入时直接访问连续 raw values;一个 Constant、一个 Flat 时把常量提到循环外。原文尝试展示的输入复用也有价值,但必须证明输入及其 buffer 可安全复用,而不是只凭“输入存在”就覆盖它。选择某条快速路径后,应避免再次落入 generic path 重算。

快速路径示意:Constant + Flat

// 已确认两输入在 rows 上无 NULL,且 result 已准备为可写 FlatVector。
const double constant = constantInput->valueAt(0);
const double* input = flatInput->rawValues();
double* output = flatResult->mutableRawValues();
rows.applyToSelected([&](vector_size_t row) {
  output[row] = constant + input[row];
});

源码:VectorFunction.hEvalCtx.hDecodedArgs.h

使用 Simple Function Interface 后,函数作者通常只需声明逐值计算。对于这里限定的 DOUBLE 加法,核心可以缩成下面这样;注册、类型检查与 Adapter 实例化把它接入表达式框架。它不是在运行时创建一个新的函数对象并简单循环调用,而是围绕函数 holder、类型 Reader / Writer 和执行路径进行模板适配。

教学示例:DOUBLE 的 Simple Function

template <typename T>
struct PlusDoubleFunction {
  FOLLY_ALWAYS_INLINE void call(
      double& result, const double& a, const double& b) {
    result = a + b;
  }
};

C++ 的 double 是这个函数实现消费的 native 值;SQL 层的 DOUBLE 类型仍由函数签名和表达式 TypePtr 描述。字符串会映射成 StringView / 字符串 writer,复杂类型则使用相应 reader / writer 视图,不必把整个 ARRAY / MAP 复制成 STL 容器。

原文的 Adapter 伪代码表达了“外层批量遍历,内层逐值 call”的关系;当前真正的类是 SimpleFunctionAdapter<FUNC>。它还要处理常量、结果复用、null-free、ASCII、字符串 ownership、错误和初始化状态,这些都不是业务函数作者需要每次重写的样板。

不要把这个 DOUBLE 示例直接扩展成所有整数的 a + b:整数溢出行为属于函数方言语义。Presto 的实际 PlusFunction 调用对应 plus(a,b) 实现,而不是让任意有符号溢出落到 C++ 未定义行为。

源码:SimpleFunctionAdapter.hSimpleFunctionAdapter.hArithmetic.hUdf.h

5. 表达式从构建到求值

Velox的expression evaluation引擎可以被用在3种场景:

  1. FilterProject算子中过滤(filter)与投影(project)的表达式求值
  2. TableScan / Connector 的过滤,包括转换成底层 reader 能处理的过滤条件,以及需要表达式求值的 residual filter。并非每个表达式都能完整下推到存储层。

  3. 独立的组件给只需要表达式求值功能的引擎,比如实时计算,ML场景的数据预处理

表达式处理分为构建可执行表达式与逐批次求值两个阶段。这里的 compilation 主要是把 typed expression tree 绑定和组织成 exec::Expr / ExprSet;不能仅凭“编译”一词就认为它必定为每条查询生成新的机器码。

5.1 构建可执行表达式

输入包含已经确定类型的表达式树,构建阶段解析调用的函数和 special form,确定子表达式关系,并建立可复用的执行对象。共同子表达式、常量折叠和 AND / OR 展平都在这个结构上发挥作用。确定性、错误与 null 语义必须一并维护。

  • 共同子表达式消除:例如 strpos(upper(a), 'FOO') > 0 OR strpos(upper(a), 'BAR') > 0,两个调用可共享确定性的 upper(a)。由于两个分支的活动行集合不一定相同,共享不等于“任何情况下整批只调用一次”;实现需要记录已经计算过哪些行,对缺少的行补算。

  • 常量折叠:例如 upper(a) > upper('Foo') 中,确定且不依赖输入的 upper('Foo') 可以变成常量 'FOO'。有副作用或非确定性的表达式不能机械折叠;可能抛出的错误也要保留 TRY、条件分支等语义,不能为了预计算改变异常出现的条件。

  • AND / OR 展平与运行时重排:编译阶段可把 AND(AND(AND(a,b),c),AND(d,e)) 展平成 AND(a,b,c,d,e);真正的自适应排序发生在后续批次求值中。ConjunctExpr 统计各子式的处理成本与剩余活动行,配置允许时优先执行更便宜地缩小活动集合的子式。

    当前 SelectivityInfo::timeToDropValue() 在输入输出行数不同时使用 timeClocks / (numIn - numOut);没有消除行时返回 timeClocks。原文写成带 1 + 的统一分母,与当前代码不符。AND 丢弃已经确定 false 的行,OR 则让已确定 true 的行不再参与后续子式;NULL 和错误仍需要单独合并处理。

源码摘录 · SelectivityInfo.h

float timeToDropValue() const {
    if (numIn_ == numOut_) {
      return timeClocks_;
    }
    return timeClocks_ / static_cast<float>(numIn_ - numOut_);
  }

源码:ConjunctExpr.cppConjunctExpr.cppExprCompiler.cpp

5.2 在选中行与编码上求值

求值接收编译好的表达式、输入向量和活动行集合,返回结果向量。沿表达式树下降时,SelectivityVector 控制当前节点需要处理哪些行;活动行不等于非空行。默认 null-in-null-out 函数可由框架提前排除 null 行,is_null、coalesce、条件表达式等则必须保留自己的 null 规则。

这种接口让计算能在不同粒度上减少工作:只算活动行、只加载需要的 LazyVector 行、在 Dictionary 的 base 行上计算,以及复用已经算过的共享子表达式或字典结果。它们依赖的是可证明的语义与结构条件。

  • Peeling:对于满足条件的确定性表达式,框架可以剥离兼容的 Dictionary / Constant 包装,把外层活动行映射到内层行,在底层值上求值,再重新包装结果。多输入是否共享兼容的封装、函数是否传播 null、是否允许 peeling,都决定这条优化能否使用。

    保留原例:color 的 1,000 行通过索引引用 [red, green, blue]。计算 upper(color) 时,如果这三个 base 值都需要,最多只需把它们变成 [RED, GREEN, BLUE],再用原来的索引映射得到逻辑长度仍为 1,000 的结果。若本批只用到一部分 base rows,也不必计算全部字典。

  • Memoization:多个输入批次可能引用同一个 base Vector,只是外层 indices 不同。Expr 可以在满足条件时缓存 base 行的计算结果,在后续批次重新包装或只补算未缓存行。原例中数百万行重复引用 red / green / blue,正是这种优化希望利用的形态。

    当前实现不是“第一次计算后,无论什么字典都直接复用”。它检查 base 对象身份与生命周期,并逐步建立重复使用的缓存;base 换了或者 weak pointer 失效时要失效重建。只有成功计算的行才能进入缓存,出错行要遵守本次 TRY / 非 TRY 的执行语义。

    入口还要求 dictionary memoization 已启用、表达式只依赖一个 distinct field、包装为 Dictionary,而且 base 未禁用 memo。多输入 peeling 与单字段 dictionary memoization 是相关但不同的优化,不能合并成一个无限制的“字典缓存”。

表达式 peeling:在底层值上计算,再重新包装
图 13 · 表达式 peeling:在底层值上计算,再重新包装 打开原图

源码:Expr.cppExpr.cppExpr.cppPeeledEncoding.h

6. 批量计算、tight loop 与 SIMD

6.1 分支外提、循环展开与自动向量化

Velox 的“向量化”首先指一次处理一批列式数据,从而摊薄函数调用、解码和分派开销。批内的 tight loop 又给编译器留下内联、常量传播、循环展开和自动 SIMD 的机会;必要时,开发者直接编写 SIMD 指令或显式交错多条独立工作。

这些机制要分开看:循环展开增加了指令级并行,不必然产生 SIMD;使用 XMM / YMM 寄存器也不意味着每条指令同时处理多个值。原文三张 Compiler Explorer 图使用 GCC 13.1 与 -O3 -mavx2 -mfma,出现的 vsqrtsd 是标量 double 平方根。把它称为自动向量化的证据会误导读者。

可运行教学示例:null-aware 与全选路径

// Teaching example. This byte flag uses 1 = null, unlike Velox's null bitmap.
template <bool kHasNulls, bool kAllRowsActive>
void squareRootKernel(
    const int32_t* __restrict positions,
    size_t count,
    const double* __restrict input,
    const uint8_t* __restrict isNull,
    double* __restrict output) {
  for (size_t i = 0; i < count; ++i) {
    const size_t row = kAllRowsActive ? i : positions[i];
    if constexpr (kHasNulls) {
      if (isNull[row]) {
        continue;
      }
    }
    output[row] = std::sqrt(input[row]);
  }
}

原例的 positions = [0,2,5,6,7],input 为 [0,1,4,8,0,1,4,8],null 标志为 [1,0,1,0,1,0,1,0]。稀疏且有 null 的版本只计算 row 5 和 row 7;未选中的行及 NULL 行保持原输出槽。这个教学 byte array 用 1 表示 NULL,不是 Velox 的 null bitmap。

kHasNulls=false 表示调用者已经证明活动行不需要 null 检查;kAllRowsActive=true 表示可以直接按连续下标访问。模板布尔参数把这些事实暴露给编译器,让它删去循环内的无效分支。__restrict 还声明输入输出不重叠,调用者必须兑现。

通用循环:行号映射与 null 分支
图 14 · 通用循环:行号映射与 null 分支 打开原图
专门化循环:把条件移到循环外
图 15 · 专门化循环:把条件移到循环外 打开原图

在当前机器使用 GCC 12.2.0、-O3 -mavx2 编译这个动态长度循环,squareRootAll 没有生成 packed vsqrtpd;加入 -fno-math-errno 后出现 packed 平方根。这里改变的是是否必须维护数学函数的 errno 语义,不是输入被换成了另一种 Vector。

两种构建均通过稀疏 null、连续行和 259 行尾部处理的结果检查。验证范围是该教学 kernel,不是整套 Velox 的性能结论;未对负数、NaN、浮点异常和应用依赖 errno 的行为作等价性承诺。工程里是否能用相应编译选项,必须由算子的语义契约决定。

本地验证命令

g++ -std=c++17 -O3 -mavx2 kernel.cpp -o kernel-strict
./kernel-strict
g++ -std=c++17 -O3 -mavx2 -fno-math-errno kernel.cpp -o kernel-no-math-errno
./kernel-no-math-errno
g++ -std=c++17 -O3 -mavx2 -fno-math-errno \
    -S -masm=intel -fopt-info-vec-optimized kernel.cpp -o kernel.s

下载完整 kernel.cpp(含结果检查)

SelectivityVector:一次分支选择两条遍历路径
图 16 · SelectivityVector:一次分支选择两条遍历路径 打开原图

原第三张截图还演示了一个同名的简化 SelectivityVector,内部固定保存五个位置。这是讲解用的模型,不是 Velox 的真实实现。当前源码在全选路径使用连续 for 循环,在稀疏路径枚举 bitmap;源码特别保留 const auto end = end_;,将循环边界提到局部变量,以利于展开与向量化。

行集合已经决定了访问方式,循环体就可以专注于值计算。这种“把分派放在批次边界,把稳定工作留在内层循环”的组织方式,比仅在代码旁写一个 vectorized 标签更关键。

源码:SelectivityVector.h

源码摘录:group probe 的四路交错 · HashTable.cpp

ProbeState state1;
  ProbeState state2;
  ProbeState state3;
  ProbeState state4;
  int32_t probeIndex = 0;
  int32_t numProbes = lookup.rows.size();
  auto rows = lookup.rows.data();
  for (; probeIndex + 4 <= numProbes; probeIndex += 4) {
    int32_t row = rows[probeIndex];
    state1.preProbe(*this, lookup.hashes[row], row);
    row = rows[probeIndex + 1];
    state2.preProbe(*this, lookup.hashes[row], row);
    row = rows[probeIndex + 2];
    state3.preProbe(*this, lookup.hashes[row], row);
    row = rows[probeIndex + 3];
    state4.preProbe(*this, lookup.hashes[row], row);

    state1.firstProbe<ProbeState::Operation::kInsert>(*this, 0);
    state2.firstProbe<ProbeState::Operation::kInsert>(*this, 0);
    state3.firstProbe<ProbeState::Operation::kInsert>(*this, 0);
    state4.firstProbe<ProbeState::Operation::kInsert>(*this, 0);

    fullProbe<false>(lookup, state1, false);
    fullProbe<false>(lookup, state2, true);
    fullProbe<false>(lookup, state3, true);
    fullProbe<false>(lookup, state4, true);
  }

四个 ProbeState 先依次 preProbe、再 firstProbe、最后 fullProbe。这样在等待某个 bucket 或 payload 进入 cache 时,可以推进其他行的独立工作,把访存延迟与计算交错。每个 state 内部再使用 SIMD 比较一组 tags。这里同时存在批量处理、软件流水式的交错与 SIMD 三种机制。

后续 state 使用额外检查,是因为前面的插入可能已经改变 bucket 内容,不能继续盲用 firstProbe 时的旧候选信息。这是批量流水优化必须维护的数据依赖,而非单纯复制四份代码。

源码:HashTable.cpp

6.2 HashTable 的编码、布局与探测

HashTable 用在 Join、Aggregation、RowNumber 等算子中,但调用语义不相同:聚合的 groupProbe 会发现并创建新 group;Join 的 probe 通常查询已经由 build 侧完成的表。建表 / 增长时,VectorHasher 对 key 的范围和 distinct 数作分析,HashTable 再选择索引模式。

  • kArray:把可编码的 key 组合成稠密 value ID,直接索引 row-pointer 数组。在映射成立的范围内不需要普通 hash 冲突处理,也不需要逐行比较原始 key;“完美哈希”可作为直觉,但实现上更准确的是范围 / distinct ID 编码后直接寻址。

  • kNormalizedKey:把多列 value IDs 组合成 64 位的无歧义编码,再使用 bucket 索引定位并比较 normalized key。它不是把任意复合 key 做一次普通 64 位 hash 后就免去冲突验证。

  • kHash:计算 hash 定位 bucket,通过 tag 初筛候选,再比较 RowContainer 中的实际 key。它承接无法使用前两种表示的情况,而不是另一个独立的算子。

  • 对于每列,VectorHasher 可以使用值域偏移或 distinct value ID,再用混合进制的方式把各列 ID 组合成一个整数:例如 id = id0 + range0 * id1,更多列继续乘上前面范围的乘积。范围乘积溢出与新增值超出已选范围都需要处理。

  • 当前 kArrayHashMaxSize = 2L << 20,即 2,097,152 个值空间,不是原文“约 1M”的固定分界;决策还会参考预留比例、range / distinct 的组合、已见 distinct 数、禁用 range-array 的设置等。不能只拿这个常量单独预测模式。

  • 当数组直接寻址不合适,但组合 ID 仍可放入 64 位时,通常可选择 normalized-key 路径。当前实现也有针对单列高基数等形态的例外,会回退到 kHash;实际选择以 decideHashMode() 的条件顺序为准。

  • 范围和 distinct 编码都不能在 64 位内安全组合,或者分析 / 运行条件不满足时,回退到 kHash。执行过程中对表示的切换必须同步处理已有行和索引,不能只改一个模式标志。

非 kArray 模式使用连续的 bucket 数组。当前每个 Bucket 为 128 字节:16 个 1 字节 tag、16 个 6 字节 row pointer,以及 16 字节尾部 padding,正好是两个 64 字节 cache line。tags 连在一起,方便一次 SIMD load / compare;真实 key 与 payload 存在 RowContainer。

“HashTable 不存 key”只能限定为“bucket 不内嵌完整 key”;整个 HashTable 对象管理的 RowContainer 当然保留真实 key,normalized-key 路径还保存对应的规范化表示。6 字节指针依赖当前实现的低 48 位地址表示假设,不能把它概括成所有 64 位 CPU / OS 都无条件适用。

非 kArray 模式的容量是 bucket 数乘以 16,bucket 数按 2 的幂组织。当前装载因子常量为 0.7,原文的 7/8 已不适用于本次源码快照。扩容会重新计算 / 增大索引,容量通常按倍数增长;初始估算、新批次的新增 group 数和具体模式也参与容量选择,不能脱离这些条件断言每次恰好只翻一倍。

HashTable:索引目录与真实行分离
图 17 · HashTable:索引目录与真实行分离 打开原图
Bucket 的 128 字节布局
图 18 · Bucket 的 128 字节布局 打开原图
16 个 tag lane 同时比较,之后才读取候选行
图 19 · 16 个 tag lane 同时比较,之后才读取候选行 打开原图

源码:HashTable.hHashTable.hHashTable.hHashTable.cpp

以Aggregation算子的probe为例子,一个hash probe过程大概如下:

  1. Pre probe,把需要probe的tag计算出来,然后用SIMD intrinsic广播到一个SIMD register,prefetch一个bucket(实际根据cache line size决定)

源码摘录:preProbe 的 tag 广播与预取 · HashTable.cpp

const auto tag = BaseHashTable::hashTag(hash);
    wantedTags_ = BaseHashTable::TagVector::broadcast(tag);
    group_ = nullptr;
    indexInTags_ = kNotSet;
    __builtin_prefetch(
        reinterpret_cast<uint8_t*>(table.table_) + bucketOffset_);
  }
  1. First probe,load一个bucket中所有tag到一个SIMD register,然后和步骤 1中register对比(SIMD intrinsic),如果有命中,prefetch第一个命中的payload row。

源码摘录:firstProbe 的 SIMD 比较 · HashTable.cpp

tagsInTable_ = BaseHashTable::loadTags(
        reinterpret_cast<uint8_t*>(table.table_), bucketOffset_);
    table.incrementTagLoads();
    hits_ = simd::toBitMask(tagsInTable_ == wantedTags_);
    if (hits_) {
      loadNextHit<op>(table, firstKey);
    }
  }
  1. Full probe:先比较已经预取的候选行,再枚举其余匹配 tag 的候选,读取真实 key 验证。候选都不相等后,如果当前 bucket 有真正的 empty tag,可判定查找未命中;插入路径则选择可用位置。如果 bucket 没有 empty,继续探测下一个 bucket。

    删除产生的 tombstone 不能直接当作 empty,否则会截断其他 key 的探测链。当前使用 0x00 作为 empty,0x7f 作为 tombstone。插入还要处理此前记录的 tombstone、并行构建的分区范围和同批其他插入造成的变化;只有 tag 匹配仍然远远不够。

7. Task、Pipeline 与 Driver 调度

7.1 从逐行迭代到批量推进

理解 Pipeline 之前,可以先看逐 tuple 的 iterator model:算子通过反复调用 next 获取一行,再把这行交给下一层逻辑。这种接口简洁、组合方便,但在大量数据上可能把调用、分派和控制流成本放大到每一行。真正要优化的是这些成本以及数据访问方式,而不只是把“pull”标签换成“push”。

  • 调用与分派:逐行 next 的调用频率高,虚函数和复杂控制流也可能阻碍内联。批量接口把部分固定成本摊到一整个 RowVector 上,但批内仍需高效实现。

  • 寄存器活跃值与 ABI:跨函数调用可能需要保存活跃寄存器值,也可能导致寄存器溢出到栈;具体取决于调用约定、内联和编译器的寄存器分配。不是每次调用都必须保存所有寄存器,更不能把任何参数传递都称为 register spilling。把相邻计算留在较小的循环内,有助于减少不必要的跨调用活跃状态。

  • 分支预测:不稳定的间接调用目标、数据相关分支和复杂控制流可能增加预测开销。稳定的虚调用也可能被预测得很好,因此需要看实际代码与 profile,不能断言虚函数天然导致预测失败。

  • 代码与数据局部性:逐行在多套算子逻辑间跳转,可能不利于指令和数据 cache。批量执行可以在同一段逻辑内反复处理连续字段,但过大的中间批次也会占 cache、增加物化成本,所以批量大小是取舍。

Spark 的 Whole-Stage Codegen 把一段相容的操作融合进生成代码,减少运行时接口成本;Photon 等原生向量化执行则更多依赖批处理、tight loop、显式 SIMD 和自适应策略。Snowflake、DuckDB 等也体现了不同形式的原生执行设计。代码生成、push / pull、批量执行属于不同维度,不应理解为所有系统必然沿同一条路线升级。

Velox 可以从数据流角度理解为 push-based vectorized pipeline:Driver 把上游的 RowVector 送入下游的 addInput()。但从调度控制看,Driver 还会调用上游 getOutput(),检查下游 needsInput(),并沿算子序列选择下一步。所以“push-based”不是说算子会递归调用下游,也不是说不存在取输出的动作。

计划片段经 LocalPlanner 拆成若干线性的 Pipeline,每个 Pipeline 可以实例化多个 Driver。每个 Driver 持有独立的 Operator 对象和推进状态,在 executor 上运行;等待 I/O、其他 Driver、下游消费者或内存仲裁时,可以通过 future 协作让出线程。耗时同步函数仍然占用调用线程,框架不会在任意 C++ 指令中途抢占它。

7.2 Driver 如何推进算子

一个本地 Task 管理多个 Pipeline 的 Driver 实例。Pipeline 描述一组具有相同算子结构的执行实例;Driver 则拥有真正的 Operator 对象。Driver 不绑定一条独占线程:某次运行被阻塞后可以离开线程,future 满足后重新排队,下一次可能由另一条线程执行;同一个 Driver 的状态推进仍需满足单线程执行约束。

每个 Driver 上的 Operators 通过该 Driver 相互通信。Driver 从某个 Operator 获取输出,跟踪统计信息,并将输出传递给下一个 Operator。Operators 之间传递的数据由vectors组成。具体来说,Operator 生产/消费的是 RowVector,它是一个vector,其中为关系的每一列都包含一个child vector,这相当于 Arrow 中的 RecordBatch。所有向量都是 velox::BaseVector 的子类。

Pipeline 的第一个 Operator 是 source,最后一个承担 sink 职责。TableScan、Exchange 等 source 通过 splits 或相应接口获得数据范围,但 Values、LocalExchange 等 source 不必都接收文件 splits。sink 也不只有 PartitionedOutput:HashBuild、LocalPartition、CallbackSink 等分别把数据放入共享表、局部交换或回调消费者。

对于跨 Task 的交换,PartitionedOutput 把分区结果交给输出缓冲管理,远端 Exchange 再作为 source 消费。FilterProject、HashProbe、HashAggregation 等可以位于中间。算子的职责位置由计划决定;不能简单把所有 source 等同于读文件,也不能把所有 sink 等同于网络发送。

Operator 暴露 isBlocked()needsInput()getOutput()isFinished()noMoreInput() 等状态接口,Driver 负责解释这些状态并推进执行,算子不通过一层套一层的递归调用驱动邻居。因此阻塞后要保存的是 Driver / Operator 的显式状态,而不是一串必须恢复的嵌套执行栈。

这使框架能够在算子调用边界协作挂起,并在之后继续推进;不是说在任意时刻都能展开并恢复任意 C++ 调用帧。状态接口的约定,尤其是 getOutput 后再次检查阻塞,决定了这种显式状态机能否正确运行。

一个 Task 如何组合执行与资源组件
图 20 · 一个 Task 如何组合执行与资源组件 打开原图
Driver:阻塞 future 如何让出并重新获得线程
图 21 · Driver:阻塞 future 如何让出并重新获得线程 打开原图

在相邻算子的常见推进路径中,Driver 先检查当前 op 是否阻塞,再检查 nextOp 是否阻塞、是否需要输入。只有下游能够消费,才向当前 op 请求输出;输出非空就通过 nextOp->addInput() 交下去,并优先尝试推进收到数据的下游。

getOutput() 返回 nullptr 之后的 isBlocked() 是不可省略的。getOutput 可能刚刚发起异步读,或者正在等待生产者;nullptr 也可能只是当前没有输出。此时 Driver 再取出 BlockingReason 与 ContinueFuture:确实阻塞则保存 BlockingState 并离开线程;未阻塞时再问 isFinished,只有真正结束才通知下游 noMoreInput。未结束则继续按状态机寻找可推进的上游工作。

调用 isBlocked(&future) 本身不会让出 CPU。真正的 off-thread 由 Driver 的阻塞处理完成。异步执行路径退出本次运行后,BlockingState::setResume 给 future 安装 continuation。事件完成时清除 blocking 状态:Task 没有暂停就 enqueue Driver,随后由 executor 再次调用 run;若 Task 暂停,则等待 Task::resume 负责入队。这个过程既不是忙轮询 future,也不要求回到原线程。

future 以异常完成时,continuation 将错误送到 Task::setError;取消、暂停和终止状态也参与恢复判断。下游 needsInput()==false 则产生背压,Driver 不应继续盲目向前读入,尤其要避免持有旧 LazyVector 的同时推进底层 reader。

源码:Driver.cppDriver.cppDriver.cppDriver.cppTask.cpp

计划树切成 Pipeline,再实例化多个 Driver
图 22 · 计划树切成 Pipeline,再实例化多个 Driver 打开原图

上图保留原文的三路扫描、两层 HashJoin 与多个 Driver 的架构信息,并明确 build / probe 的依赖。一个 build pipeline 完成后,通过 JoinBridge 解除相应 probe 的等待;共享 HashTable 不等于跨 pipeline 共享同一套 Operator 实例。

源码:LocalPlanner.cppHashJoinBridge.h

8. 内存管理与 spill 恢复

8.1 Pool、Arbitrator、Allocator 的分工

内存管理要解决的是:同一个长运行进程中,多个快慢不同、状态大小不同的查询如何共享有限容量,并在内存压力下继续推进。秒级、分钟级、小时级查询共存时,仅靠某个算子“尽量少申请”不够,需要查询级上限、细粒度记账、可回收状态和失败处理共同配合。

Velox 的相关设计与 Presto 等多查询执行场景紧密相关,但“公平”不是一个不受配置约束的绝对保证。仲裁策略、查询优先级、可 spill 的状态、超时与容量上限决定实际行为;资源实在不足时,查询仍可能失败或被 abort。

  • 分配效率:Allocator 使用适合不同大小的字节 / page 分配路径,MmapAllocator 的 SizeClass 管理页段;MemoryPool 的 quantization 则减少 reservation 传播与锁竞争。这是不同层面的优化,不能统称为“小块内存都按 MiB 物理分配”。

  • 共享与约束:MemoryPool 跟踪查询和算子的使用与预留,Arbitrator 在参与的 root pools 之间调配 capacity。容量约束降低了无界增长的风险,但无法覆盖所有进程内存:外部分配、allocator 元数据、线程栈、文件映射和缓存等仍需计入进程预算,不能因此保证 OOM killer 永远不会触发。

内存管理器主要由三个部分组成:

  • 内存池(Memory Pool),主要用于精细地跟踪内存使用情况;
  • 仲裁器(Arbitrator),负责协调各个内存单元之间的共享与仲裁;
  • 分配器(Allocator),负责实际的物理内存分配工作。

QueryCtx 持有查询使用的 root MemoryPool。集成方可以传入这个 pool;默认路径也可通过 MemoryManager 创建。Task 在 root 下创建 aggregate pool,再按 plan node 和 Operator 实例建立下层池。leaf 申请更多 reservation、root capacity 不够时,才可能进入增长和仲裁流程;增长不等于一定发生 spill。

默认树的主干是 Query root → Task aggregate → PlanNode aggregate → Operator leaf。同一个 plan node 可以对应不同 Driver 上的多个 Operator leaf;pipeline ID 和 driver ID 体现在池名字里,而不是自动多出两层固定的 PipelinePool、DriverPool。HashJoin 的 split-group 池、Connector 子池和 custom memory resource 还可以扩展默认结构。

默认 MemoryPool 树:Query → Task → PlanNode → Operator
图 23 · 默认 MemoryPool 树:Query → Task → PlanNode → Operator 打开原图

源码:QueryCtx.hQueryCtx.cppTask.cppTask.cppTask.cpp

Leaf Pool 面向实际分配接口,Aggregate Pool 汇总和组织其子池;root 也是 aggregate pool,并承担 capacity 控制。分配通常从 leaf 发起,经 reservation 协议检查 root 的预算,再调用底层 allocator。原文“实际申请都从根池发起”容易把配额检查与物理分配入口混淆。

树状记账允许定位查询、Task、plan node 和具体 Operator 的内存使用;代价是 reservation 调整、同步和回收协调。实现通过复用已有 reservation、批量量化和缓存可回收信息等方式减少频繁沿树传播的开销。读统计时必须确认层级:root 的已保留容量与 leaf 的实际使用量不是同一个字段。

8.2 Reservation 与 capacity 如何配合

MemoryPool::quantizedSize 按累计需求量化 reservation:小于 16 MiB 时按 1 MiB 向上取整;16 MiB 到不足 64 MiB 时按 4 MiB;64 MiB 及以上按 8 MiB。这决定向上预留多少,不代表一个 100-byte 的物理分配请求就一定分配 1 MiB。

在稳定时刻,可以依次区分:实际 used bytes、为 leaf 预留的 reserved bytes、root 当前获得的 capacity、root 最大允许的 maxCapacity。已获得 capacity 中未被 reservation 占用的部分,可供查询继续增长;已预留但尚未使用的部分,也不等同于任意其他查询可直接取走的容量。

只有需要的 reservation 无法在当前预算内得到满足时,才需要尝试增长 root capacity 或回收状态。这个动作可能使用仲裁器的空闲容量、收缩其他 root 的闲置 capacity,或在配置允许时通过 spill 回收已用内存。物理内存分配还要经过 allocator,不能仅靠移动图上的容量边界就获得实际 buffer。

Reservation 的量化,不是 malloc 的分配粒度
图 24 · Reservation 的量化,不是 malloc 的分配粒度 打开原图
四个量:used、reserved、capacity、maxCapacity
图 25 · 四个量:used、reserved、capacity、maxCapacity 打开原图

源码摘录:reservation 量化 · MemoryPool.h

FOLLY_ALWAYS_INLINE static uint64_t quantizedSize(uint64_t size) {
    if (size < 16 * kMB) {
      return bits::roundUp(size, kMB);
    }
    if (size < 64 * kMB) {
      return bits::roundUp(size, 4 * kMB);
    }
    return bits::roundUp(size, 8 * kMB);
  }

8.3 仲裁如何安全地调用算子回收

Arbitrator 管理参与仲裁的 root pools 与可分配 capacity。回收可以只收回未使用的 capacity,也可以通过 Reclaimer 让可 spill 的算子释放已用状态;必要时还会依据策略 abort 查询。查询初始 capacity 可以减少启动阶段频繁仲裁,但它不是保证查询之后永远无需协调的预付物理内存。

当前 SharedArbitrator 先检查请求方自身、容量限制和全局空闲容量,再尝试回收闲置 capacity。全局仲裁未启用时,会在条件满足时尝试请求方自身的回收;启用时可交由全局仲裁流程选择其他参与者。具体顺序由代码分支和配置决定,不能概括成“先 spill 自己,再一定 spill 别人”。

以 Query A 的 Operator 申请内存、最后需要从 Query B 回收为例:A 的 leaf reservation 需求到达 root 后触发 grow,MemoryManager 把请求交给 Arbitrator;后者选择有可回收能力的候选 B。接下来沿 B 的 Query / Task / Node reclaimer 到达相关 Operator。

首先建立安全边界,再回收状态。Task::MemoryReclaimer 请求 Task pause 并等待 Driver 到达允许回收的状态;随后再在安全范围内遍历、选择可回收状态,调用算子的 reclaim / spill。不是从任意线程直接清空另一个线程正在使用的 HashTable,也不是暂停后把所有算子无条件 spill。

回收有超时和异常路径。若 reclaim 抛出异常,Task 先记录错误,避免恢复后继续使用不一致的状态,然后执行 resume 收尾;成功则释放相应 reservation / capacity,允许等待增长的请求继续。算子内部的非回收临界区与可回收量判断,是这个协议成立的另一部分。

一次内存增长:先争取容量,再执行物理分配
图 26 · 一次内存增长:先争取容量,再执行物理分配 打开原图

源码:SharedArbitrator.cppSharedArbitrator.cppTask.cppTask.cppMemoryReclaimer.cppOperator.cpp

8.4 SizeClass、虚拟地址与物理页

Allocator 负责真正提供和回收地址空间 / buffer。Velox 可使用 MallocAllocator 或 MmapAllocator,字节分配、连续页分配和非连续页分配各有路径;MmapAllocator 的 SizeClass 主要用于按页组织的分配,不能代替所有小对象 allocator。

固定大小 buffer 便于管理,但不适合所有请求;完全任意大小的 arena 又增加分割、合并和碎片管理的复杂度。按大小类组织页段,并通过虚拟地址与物理页分离管理,可以在分配速度、内部碎片和灵活性之间折中。原分享借助 Umbra 一类可变长度 buffer 管理的思路解释这个动机;当前行为仍应以 MmapAllocator 的实现和配置为准。

可以设想为每个大小类预留一大片虚拟地址范围,例如全局允许 100 GiB 的页容量时,各 class 的地址空间分别按其容量布局。匿名 mmap 通常不会在建立映射的一刻就触碰全部对应物理页;实际使用时再按需建立页面,所以虚拟地址总量可以大于实际使用量。

但这不是“物理内存和系统开销绝对为零”:页表、管理 bitmap、allocator 元数据和已触碰页面仍有成本;过度承诺与操作系统策略也有影响。Velox 需要全局 allocated / mapped 计数限制实际可保留的页面,不能因为每个 class 都有一大片地址空间就无限用下去。

Page SizeClass 与 150-page 请求
图 27 · Page SizeClass 与 150-page 请求 打开原图
MmapAllocator:free 与归还物理页不是同一步
图 28 · MmapAllocator:free 与归还物理页不是同一步 打开原图

保留原文的 150-page 例子:最小 SizeClass 为 4 pages 时,当前算法可组合成 128×1 + 16×1 + 4×2 = 152 pages,只比需求多 2 pages。原图写的 64×2 + 16×1 + 4×2 也等于 152,而不是图中标出的 150;当前 descending-size 选择并不保证复现原图的 64-page 分解。

调用者 free 后,页面可以变成“空闲但仍 mapped”的状态,供下一次分配快速复用;这和把物理页交还操作系统是两步。需要减少 mapped 页占用时,MmapAllocator 对可释放页执行 madvise(..., MADV_DONTNEED)。原文的 madvice 拼写应为 madvise。

因此应把 allocated、free-but-mapped、advised / not-mapped 三种状态分开;跟踪这些状态的 bitmap 不只是一个“用没用”的布尔值。图中的 4 KiB page 用于解释例子,实际运行仍须符合实现与平台的页大小约束。

源码:MemoryAllocator.cppMmapAllocator.cppMmapAllocator.cppMmapAllocator.cpp

8.5 Spill 保存的是算子恢复所需的状态

Spill 的核心不是“把所有内存原样倒进文件”,而是把以后恢复语义所需的数据序列化。HashTable 相关算子常使用 RowContainer 保存 key、payload 或 aggregate state;但是并非所有 Operator 的状态都必须放进 RowContainer,输入 Vector 和其他结构也可以由相应 Spiller 路径写出。

RowContainer 把列式输入转换成适合按 group / row 地址访问的状态位置,这常称为 pivot。聚合按 group 找到 accumulator,Join 按 key 找到 build payload,RowNumber 按 partition key 找到 counter。外部接口仍是 RowVector,这种内部行式状态并不否定执行引擎整体的向量化设计。

HashAggregation:输入的 group keys(例如 a、b、c)通过 HashTable 定位到 RowContainer 中的 accumulator。avg、max、array_agg 等函数的状态大小与更新方式不同;avg 的 spill 表示需要 sum / count 等可合并中间态,不能只存一个最终平均数。原分享观察到的聚合收益是特定 workload 经验,不能脱离基数、数据宽度和函数实现推广成统一加速倍数。

典型输入阶段 spill 先按 hash bits 分区,再在每个 partition 内按 grouping keys 排序,写出一个或多个有序 run。相同 group 可能在多次 spill 中出现,排序服务于后续归并时合并中间态,不意味着 SQL GROUP BY 输出自带 ORDER BY。若输入阶段已经 spill,noMoreInput 会把剩余内存状态也冲出。

分区写出可以并行,所以当前实现用 HashStringAllocator::freezeAndExecute 约束 spill 期间对相关分配器的修改,避免多个分区并行序列化时触发线程安全问题。这是吞吐优化必须遵守的共享状态约束。

恢复阶段使用 createOrderedReader 逐个处理 partition,对该分区的 runs 做多路 merge,把相同 key 的中间态组合后形成结果,再进入下一分区。这样不必同时把所有分区恢复为一个巨大的 HashTable。实现中还有 distinct 专用路径和输出阶段 spill,不能把这一典型流程当成 GroupingSet 的唯一状态机。

HashAggregation:分区 spill 与有序归并恢复
图 29 · HashAggregation:分区 spill 与有序归并恢复 打开原图

源码:GroupingSet.cppGroupingSet.cppGroupingSet.cppGroupingSet.cppGroupingSet.cppGroupingSet.cpp

RowNumber:原文这里把 RowNumber 算子写成了 RowContainer,需要纠正。存在 partition keys 时,当前 RowNumber 确实使用 HashTable,将分组 key 映射到一个 BIGINT counter;没有 partition keys 时则不必建立这张分组表。它为每行延续分区内的计数,不应与通用的带 ORDER BY 窗口执行混为一谈。

发生 spill 后有两类数据需要区别:已经处理过输入所对应的 keys + counters,以及尚未处理的后续输入 Vector。前者保存的是计算进度,不是要再次输出的一份已完成结果。进入 spill 模式后,后续输入也写入对应的分区文件,等待恢复处理。

恢复当前 partition 时,先以 unordered reader 读取状态文件,重建 HashTable 并恢复 counters;再读取同一 partition 的输入文件,执行 probe、更新 counter 并产生新的 row-number 输出。这个状态机无需借助按 grouping key 排序的 merge,和 HashAggregation 的中间态归并路径不同。

如果当前分区仍然太大,可以在允许的层数内继续递归分区 / spill。输出发给下游消费者,而不是原文所写的“上游”;是否需要额外排序取决于计划语义,不能把 RowNumber 的 unordered restore 当作所有窗口算子的通用方案。

RowNumber:恢复 partition counter,再继续输入
图 30 · RowNumber:恢复 partition counter,再继续输入 打开原图

源码:RowNumber.cppRowNumber.cppRowNumber.cppRowNumber.cppRowNumber.cpp

9. 为什么共享执行层值得做

现代数据工作负载的多样性在不断增加,数据集指数级增长,导致了专门的查询和计算引擎的泛滥,每个引擎针对特定类型的工作负载。数据处理需求从简单的事务处理和分析(批处理和交互式),发展到ETL和大规模数据移动,再到实时流处理,以及用于监控用例的日志和时间序列处理,最近还包括大量的人工智能(AI)和机器学习(ML)用例,包括数据预处理和特征工程。这种演变导致了一个由数十个专门引擎组成的孤立数据生态系统,这些引擎使用不同的框架和库构建,彼此之间几乎没有共享,使用不同的编程语言编写,并由不同的工程团队维护。

这种碎片化也会影响用户和开发团队:不同系统的数据类型、函数参数、null 传播、类型转换和异常规则不一致,任务在系统间迁移时就不只是换一个语法。Velox 早期论文曾用 Meta 内部至少十余种 substr 实现举例,说明重复实现和语义差异并存;这是论文背景中的观察,不是本次对所有系统重新测量的结果。

共享执行库的目标是减少重复建设,同时让差异显式存在。Presto 与 Spark 函数包可以共享 Vector、Expr 和内存机制,但仍需要各自兼容的函数语义,不能为了“统一”强行把 0-based / 1-based 下标、异常处理和 null 规则改成同一份行为。

专用引擎的差异通常集中在语言前端、优化器、分布式运行时、存储接入及对外语义,而本地执行层有大量共同需求:表达 scalar / complex 类型、组织列式内存、计算表达式、实现 Join / Aggregation / Sort,以及序列化、spilling、缓存和内存管理。The Composable Data Management System Manifesto 所倡导的组件化,可以从这些边界理解。

Velox 的价值不是让所有系统最终长成一个完整数据库,而是提供能独立集成的执行机制:统一的 Type 与 Vector 让组件交换数据,函数框架把语义和样板分开,显式 Driver 状态机让等待与恢复可管理,Reclaimer 协议让资源治理能安全触达算子内部状态。

可组合的数据系统:接口连接组件
图 31 · 可组合的数据系统:接口连接组件 打开原图

9.1 从实现看设计上的取舍

保留表示信息,才有机会少做工作。Dictionary、Constant、Lazy 和 SelectivityVector 把重复、未加载和未选中的事实显式传给执行层,peeling、memoization 与稀疏求值据此减少计算。代价是所有消费者都必须遵守编码、null、ownership 和生命周期协议,不能只按 FlatVector 的理想输入编写。

把变化放在外层,把稳定工作留在内层。Type 分派、编码选择和 all-selected 判断适合在批次或循环外完成,内层尽量保留简单、连续、可内联的工作。HashTable 的四路 probe 又说明,简单 tight loop 之外仍有必要显式安排预取和独立访存;这种优化应由真实瓶颈支持。

显式状态比隐式调用栈更适合协作调度。Operator 报告 blocked / needsInput / finished,Driver 集中处理;future 只表达事件完成,是否重新入队还要尊重 Task 状态。收益是可以让出线程、实施背压和处理取消,代价是状态转换、结果 ownership 和异常传播必须写得严密。

资源治理必须知道算子的状态。Allocator 只知道地址和页,无法判断哪些聚合状态能 spill;Pool 只记账,也不能在任意时刻安全清空状态。Arbitrator、Task pause 和 Operator reclaimer 把策略、安全边界和具体回收连接起来。回收可能触发 I/O、延迟和错误,因此不是一项免费能力。

阅读这些实现时,最值得保留的不是某个固定性能口号,而是它们各自成立的条件:值范围是否适合编码、selected rows 是否连续、字典是否相容、结果 buffer 是否可写、Task 是否已安全暂停、恢复是否需要排序。条件讲清楚,设计收益与代价才有可检查的依据。

9.2 源码与延伸阅读

Velox 官方网站Velox: Meta’s Unified Execution Engine(PVLDB 2022)The Composable Data Management System Manifesto。历史背景与应用范围应结合这些资料的发表时间理解。

原分享的相关资料:KDB+CWIPhoton 论文Apache ArrowSimple Function Interface 论文Filter Representation in Vectorized Query Execution

布局与字符串:Umbra 论文StringView(上)StringView(下)Hash table 布局与性能比较嵌套 Array 编码改进

更细的模块分析可继续阅读:Type SystemHashTableTask & DriverMemory Pool and ArbitratorSpiller。这些文章分别注明自己的源码快照,遇到默认值差异以对应版本为准。