向量化 - 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 与这些系统的联系帮助建立技术脉络。这里保留这层背景,但不把人员经历、商业竞争猜测或示意箭头当成代码继承关系;也不使用“第一个向量化数据库”这样的绝对断言。
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.h;PlanFragment.h;LocalPlanner.cpp;Task.h。
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。
源码:Type.h;Type.h;Type.h;Type.h;VectorEncoding.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.h;LazyVector.h;DecodedVector.h;DecodedVector.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 槽。
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,是三个不同问题。
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.h;ConstantVector.h;ConstantVector.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,就是按这种结构传递一批关系行。
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;
};FlatVector<StringView> 的 values buffer 由 16 字节描述符组成。长串保存 4 字节长度、4 字节前缀和 8 字节指针,指针指向完整字符串;长度不超过 12 字节的字符串直接占用 prefix 与 union 的空间。这里计算的是字节数,不是 Unicode 字符数。
向量的 stringBuffers 维持外部字符串内存的所有权,而 StringView 本身不是 owning string。复制一个长串 StringView 只复制了描述符;如果没有同时维持字符缓冲区生命周期,指针会悬空。下图把原例中的 orange peel 长度修正为 11 字节,其余短串与长串例子仍保留。
- 长度不超过 12 个字节的短字符串完全内联在StringView的value buffer中,这样可以增加内存局部性,cache locality更好。
比较长度和前缀可以快速排除不同字符串,减少访问长串 buffer 的机会;长串长度和前缀相同后仍需比较剩余内容。相等比较与字典序比较的短路条件不同,实际比较次序以 StringView 的操作符实现为准。
trim / substr 等操作在满足实现条件时,可以让输出 StringView 引用输入字符的一段,避免复制长串字符。Adapter 会根据函数声明获取相应输入 stringBuffers 的共享所有权;短串可能直接内联。零拷贝是特定路径的优化,不是这些函数对所有输入的无条件保证。
- 由于 StringView 具有固定大小(16 个字节),因此可以乱序写入 StringView Vector(例如,先在位置 2 写入 StringView,然后是 0 和 1)。
源码:StringView.h;StringView.h;SimpleFunctionAdapter.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.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 的复杂输入无条件使用它。
以两个 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];
});使用 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.h;SimpleFunctionAdapter.h;Arithmetic.h;Udf.h。
5. 表达式从构建到求值
Velox的expression evaluation引擎可以被用在3种场景:
- FilterProject算子中过滤(filter)与投影(project)的表达式求值
TableScan / Connector 的过滤,包括转换成底层 reader 能处理的过滤条件,以及需要表达式求值的 residual filter。并非每个表达式都能完整下推到存储层。
- 独立的组件给只需要表达式求值功能的引擎,比如实时计算,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.cpp;ConjunctExpr.cpp;ExprCompiler.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 是相关但不同的优化,不能合并成一个无限制的“字典缓存”。
源码:Expr.cpp;Expr.cpp;Expr.cpp;PeeledEncoding.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 还声明输入输出不重叠,调用者必须兑现。
在当前机器使用 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原第三张截图还演示了一个同名的简化 SelectivityVector,内部固定保存五个位置。这是讲解用的模型,不是 Velox 的真实实现。当前源码在全选路径使用连续 for 循环,在稀疏路径枚举 bitmap;源码特别保留 const auto end = end_;,将循环边界提到局部变量,以利于展开与向量化。
行集合已经决定了访问方式,循环体就可以专注于值计算。这种“把分派放在批次边界,把稳定工作留在内层循环”的组织方式,比仅在代码旁写一个 vectorized 标签更关键。
源码摘录: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.h;HashTable.h;HashTable.h;HashTable.cpp。
以Aggregation算子的probe为例子,一个hash probe过程大概如下:
- 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_);
}- 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);
}
}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 后再次检查阻塞,决定了这种显式状态机能否正确运行。
在相邻算子的常见推进路径中,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.cpp;Driver.cpp;Driver.cpp;Driver.cpp;Task.cpp。
上图保留原文的三路扫描、两层 HashJoin 与多个 Driver 的架构信息,并明确 build / probe 的依赖。一个 build pipeline 完成后,通过 JoinBridge 解除相应 probe 的等待;共享 HashTable 不等于跨 pipeline 共享同一套 Operator 实例。
源码:LocalPlanner.cpp;HashJoinBridge.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 还可以扩展默认结构。
源码:QueryCtx.h;QueryCtx.cpp;Task.cpp;Task.cpp;Task.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 量化 · 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,允许等待增长的请求继续。算子内部的非回收临界区与可回收量判断,是这个协议成立的另一部分。
源码:SharedArbitrator.cpp;SharedArbitrator.cpp;Task.cpp;Task.cpp;MemoryReclaimer.cpp;Operator.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 都有一大片地址空间就无限用下去。
保留原文的 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.cpp;MmapAllocator.cpp;MmapAllocator.cpp;MmapAllocator.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 的唯一状态机。
源码:GroupingSet.cpp;GroupingSet.cpp;GroupingSet.cpp;GroupingSet.cpp;GroupingSet.cpp;GroupingSet.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.cpp;RowNumber.cpp;RowNumber.cpp;RowNumber.cpp;RowNumber.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 协议让资源治理能安全触达算子内部状态。
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+、CWI、Photon 论文、Apache Arrow、Simple Function Interface 论文、Filter Representation in Vectorized Query Execution。
布局与字符串:Umbra 论文、StringView(上)、StringView(下)、Hash table 布局与性能比较、嵌套 Array 编码改进。
更细的模块分析可继续阅读:Type System、HashTable、Task & Driver、Memory Pool and Arbitrator、Spiller。这些文章分别注明自己的源码快照,遇到默认值差异以对应版本为准。