Apache Arrow Gandiva 表达式、Projector 与 Filter 使用指南(C++)
【免费下载链接】arrowApache Arrow is the universal columnar format and multi-language toolbox for fast data interchange and in-memory analytics项目地址: https://gitcode.com/GitHub_Trending/arrow3/arrow
导读
本文围绕 Apache Arrow C++ 子模块 Gandiva 的表达式(Expression)、投影器(Projector)与过滤器(Filter)展开,介绍如何用TreeExprBuilder构建表达式树、将其包装为Expression或Condition,再通过Projector::Make/Filter::Make编译出可执行内核并批量求值。读完本文,你将掌握 Gandiva 表达式构建的完整 API 组合、SelectionVector 位宽选型原则,以及“先过滤再投影”这一高性能组合用法,并能在 cpp/examples/arrow/gandiva_example.cc 中找到可直接运行验证的完整示例。
Gandiva 是 Apache Arrow 中基于 LLVM 的表达式求值引擎:表达式在构建后被编译为机器码执行,从而以接近原生代码的速度在 Arrow 列式数据上完成投影与过滤。本文严格依据 docs/source/cpp/gandiva/expr_projector_filter.rst 编写,并结合
cpp/src/gandiva目录下的头文件与示例源码进行补充。
构建表达式:TreeExprBuilder 与表达式树
Gandiva 提供了一套通用的表达式表示:表达式由一棵节点树(tree of nodes)构成,所有树的构建都通过TreeExprBuilder完成。表达式树的叶子节点通常是两类:
- 字段引用,由
TreeExprBuilder::MakeField创建; - 字面量,由
TreeExprBuilder::MakeLiteral(以及MakeStringLiteral、MakeBinaryLiteral、MakeDecimalLiteral、MakeNull)创建。
叶子节点可以被组合成更复杂的表达式树,官方文档给出了四类核心组合器:
| 组合器 | 作用 |
|---|---|
TreeExprBuilder::MakeFunction | 创建函数节点;可调用GetRegisteredFunctionSignatures获取合法的函数签名列表 |
TreeExprBuilder::MakeIf | 创建 if-else 分支逻辑 |
TreeExprBuilder::MakeAnd/MakeOr | 创建布尔表达式;取“非”请使用MakeFunction中的not(bool)函数 |
TreeExprBuilder::MakeInExpressionInt32等 | 创建集合成员测试(in 表达式) |
在 tree_expr_builder.h 中可以看到这些 API 的完整声明。其中MakeFunction的原型为:
static NodePtr MakeFunction(const std::string& name, const NodeVector& params, DataTypePtr return_type);即需要显式给出函数名(如"add"、"less_than")、参数节点列表以及返回类型。MakeIf则接收condition、then_node、else_node和result_type四个参数,构成典型的三目分支。
值得留意的是,MakeInExpression*在源码中被细分为多种重载,覆盖不同数据类型:MakeInExpressionInt32/Int64、Decimal、String、Binary、Float、Double、Date32/Date64、Time32/Time64 与 TimeStamp,每种都接收一个节点和一个std::unordered_set常量集合。这意味着集合成员测试同样具备类型安全的重载选择。
每个组合器都会创建新的复合节点,复合节点以叶子节点或其他复合节点作为子节点;通过不断组合,可以构建任意复杂度的表达式树。树构建完成后,根据用途不同,会被包装成两种对象:
Expression:用于投影(projection);Condition:用于过滤(filter)。
下面是从 gandiva_example.cc 中截取的官方示例——创建表示x + 3的Expression与表示x < 3的Condition:
std::shared_ptr<arrow::Field> field_x_raw = arrow::field("x", arrow::int32()); std::shared_ptr<Node> field_x = TreeExprBuilder::MakeField(field_x_raw); std::shared_ptr<Node> literal_3 = TreeExprBuilder::MakeLiteral(3); std::shared_ptr<arrow::Field> field_result = arrow::field("result", arrow::int32()); std::shared_ptr<Node> add_node = TreeExprBuilder::MakeFunction("add", {field_x, literal_3}, arrow::int32()); std::shared_ptr<Expression> expression = TreeExprBuilder::MakeExpression(add_node, field_result); std::shared_ptr<Node> less_than_node = TreeExprBuilder::MakeFunction("less_than", {field_x, literal_3}, arrow::boolean()); std::shared_ptr<Condition> condition = TreeExprBuilder::MakeCondition(less_than_node);注意这里的几个细节:MakeLiteral(3)通过 C++ 重载推断为int32_t字面量;MakeFunction("add", ...)返回类型显式指定为arrow::int32();条件树的根节点则通过MakeCondition包装为Condition。MakeExpression(root_node, result_field)需要同时给出结果字段,用于声明投影输出的列名与类型。
函数签名查询:GetRegisteredFunctionSignatures
文档特别指出:创建函数节点前,可以调用GetRegisteredFunctionSignatures获取合法函数签名列表,从而确认某个函数名、参数类型与返回类型组合是否受支持。在 Gandiva 源码中,函数注册表由 function_registry.cc 及其按领域划分的注册文件(如 function_registry_arithmetic.cc、function_registry_string.cc、function_registry_datetime.cc 等)构建,涵盖算术、字符串、日期时间、哈希、数学函数、时间戳运算等类别。FunctionSignature类型及注册表查询接口定义在 function_signature.h 与 function_registry.h 中。
对于开发者而言,最稳妥的实践是:在写死某个函数名之前,先枚举签名表核对名称与类型组合,避免在Projector::Make/Filter::Make阶段才暴露签名不匹配的错误。
Projector 与 Filter:两个执行内核
Gandiva 提供两个执行内核:
Projector:消费一个 record batch,将其投影(project)为新的 record batch;Filter:消费一个 record batch,产出**SelectionVector**——一个包含所有满足条件行索引的向量。
两者在创建实例时都会完成表达式 IR 的优化,并且针对静态 schema 编译,因此此时必须已知 record batch 的 schema。这一约束在源码中同样体现:Projector::Make与Filter::Make的第一个参数都是SchemaPtr schema(见 projector.h 与 filter.h)。
延续上一节的expression与condition,官方示例创建 Projector 与 Filter 的代码如下(gandiva_example.cc):
std::shared_ptr<arrow::Schema> input_schema = arrow::schema({field_x_raw}); std::shared_ptr<arrow::Schema> output_schema = arrow::schema({field_result}); std::shared_ptr<Projector> projector; Status status; std::vector<std::shared_ptr<Expression>> expressions = {expression}; status = Projector::Make(input_schema, expressions, &projector); ARROW_RETURN_NOT_OK(status); std::shared_ptr<Filter> filter; status = Filter::Make(input_schema, condition, &filter); ARROW_RETURN_NOT_OK(status);Projector::Make接收表达式向量(可一次编译多个投影表达式),Filter::Make接收单个Condition。此外Projector还提供两个可选Make重载,分别用于注入运行时Configuration以及指定SelectionVector::Mode(详见后文“过滤 + 投影组合使用”一节);Filter::Make同样支持自定义Configuration。
一旦 Projector 或 Filter 创建完成,即可对 Arrow record batch 反复求值。这两个执行内核自身是单线程的,但被设计为可复用、并行处理不同的 record batch——即把同一个 Projector/Filter 实例分发给多个线程,各自处理不同的批次,从而获得并行吞吐。
求值投影:Projector::Evaluate
执行通过Projector::Evaluate完成,它输出一个数组向量(arrow::ArrayVector),将该向量连同输出 schema 一起传给arrow::RecordBatch::Make()即可组装出结果批次。源码中Evaluate有多个重载(见 projector.h):
- 无选择向量版本:
Evaluate(batch, pool, &output),输出数组从内存池pool分配; - 无选择向量、调用者自备输出数组版本:
Evaluate(batch, &output); - 带选择向量版本:
Evaluate(batch, selection_vector, pool, &output)与Evaluate(batch, selection_vector, &output),后者要求调用者预先分配具有足够容量的数组。
官方示例(gandiva_example.cc)展示了完整的投影求值流程:
auto pool = arrow::default_memory_pool(); int num_records = 4; arrow::Int32Builder builder; int32_t values[4] = {1, 2, 3, 4}; ARROW_RETURN_NOT_OK(builder.AppendValues(values, 4)); ARROW_ASSIGN_OR_RAISE(std::shared_ptr<arrow::Array> array, builder.Finish()); auto in_batch = arrow::RecordBatch::Make(input_schema, num_records, {array}); arrow::ArrayVector outputs; status = projector->Evaluate(*in_batch, pool, &outputs); ARROW_RETURN_NOT_OK(status); std::shared_ptr<arrow::RecordBatch> result = arrow::RecordBatch::Make(output_schema, outputs[0]->length(), outputs);这里输入批次in_batch的 schema 必须与Projector::Make时传入的 schema 一致;输出数组由Evaluate从pool中分配,因此调用方无需预分配内存。
求值过滤:Filter::Evaluate 与 SelectionVector
Filter::Evaluate产出SelectionVector——一个匹配过滤条件的行索引向量。从源码看,SelectionVector 本质上是Arrow 整数数组的包装器,按位宽参数化(selection_vector.h),其Mode枚举包括MODE_NONE、MODE_UINT16、MODE_UINT32、MODE_UINT64,并提供MakeInt16/MakeInt32等工厂方法(MakeInt16(max_slots, pool, &selection_vector)为最常见用法)。
SelectionVector 必须在传给Evaluate()之前初始化,初始化时需要确定两个关键参数:
- 位宽(bitwidth):决定它能容纳的最大索引值;
- 最大槽位数(max slots):决定它能包含多少条索引。
文档给出的选型原则非常明确:
一般情况下,max slots 应设置为批次大小(batch size),位宽应选择能表示所有小于批次大小的整数的最小整数位宽。例如批次大小为 100k 时,max slots 设为 100k,位宽选 32 位(因为 2^16 = 64k,不足以表示 100k 以内的索引)。
这一点在示例代码中也有印证(gandiva_example.cc):4 行数据的批次使用MakeInt16,并将max_slots设为in_batch->num_rows()。
Evaluate()运行完毕后,SelectionVector 被填充,此时用SelectionVector::ToArray()取出底层数组,再交给arrow::compute::Take()物化出最终输出批次:
std::shared_ptr<gandiva::SelectionVector> result_indices; // Use 16-bit integers for indices. Result can be no longer than input size, // so use batch num_rows as max_slots. status = gandiva::SelectionVector::MakeInt16(/*max_slots=*/in_batch->num_rows(), pool, &result_indices); ARROW_RETURN_NOT_OK(status); status = filter->Evaluate(*in_batch, result_indices); ARROW_RETURN_NOT_OK(status); std::shared_ptr<arrow::Array> take_indices = result_indices->ToArray(); Datum maybe_batch; ARROW_ASSIGN_OR_RAISE(maybe_batch, arrow::compute::Take(Datum(in_batch), Datum(take_indices), TakeOptions::NoBoundsCheck())); result = maybe_batch.record_batch();由于Take使用TakeOptions::NoBoundsCheck(),此处可以避免多余的越界检查开销;依赖Take的接口位于 arrow/compute/api_vector.h(示例文件第 19 行#include)。
过滤 + 投影组合使用
最后,Gandiva 还支持在投影的同时应用选择向量,即带过滤的投影。关键前提有两个:
- 初始化 Projector 时,传入
SelectionVector::GetMode(),使投影器按正确的位宽编译; - 求值时把 SelectionVector 传入
Projector::Evaluate()的带选择向量重载。
Projector::Make的对应重载为Make(schema, exprs, selection_vector_mode, configuration, &projector)(见 projector.h)。官方示例(gandiva_example.cc):
// Make sure the projector is compiled for the appropriate selection vector mode status = Projector::Make(input_schema, expressions, result_indices->GetMode(), ConfigurationBuilder::DefaultConfiguration(), &projector); ARROW_RETURN_NOT_OK(status); arrow::ArrayVector outputs_filtered; status = projector->Evaluate(*in_batch, result_indices.get(), pool, &outputs_filtered); ARROW_RETURN_NOT_OK(status); result = arrow::RecordBatch::Make(output_schema, outputs[0]->length(), outputs_filtered);注意此处显式传入ConfigurationBuilder::DefaultConfiguration()(示例第 33 行using gandiva::ConfigurationBuilder;),其完整定义位于 configuration.cc。该组合的完整运行结果依次打印三段输出:Project result:、Filter result:与Project + filter result:,读者可自行编译 gandiva_example.cc 验证:
- 输入
x ∈ {1, 2, 3, 4}; - 投影
x + 3得到{4, 5, 6, 7}; - 过滤
x < 3选出索引{0, 1},经Take物化后得到{1, 2}; - 过滤 + 投影组合则直接对选中行计算
x + 3,得到{4, 5}。
实践要点小结
- schema 先行:Projector/Filter 均在创建时针对静态 schema 编译(projector.h、filter.h),因此输入批次的 schema 必须与
Make时一致; - 实例复用:执行内核单线程,但同一实例可在多个线程上并行处理不同批次;
- SelectionVector 提前初始化:位宽与 max slots 需在
Evaluate前确定,位宽不足会限制可表示的最大索引; - 组合求值位宽匹配:带选择向量的投影必须用
SelectionVector::GetMode()编译 Projector,否则位宽不匹配; - 完整示例:所有可运行代码集中在 gandiva_example.cc,配套构建入口见 cpp/examples/arrow/CMakeLists.txt;Gandiva 的 C API 绑定与测试分别位于 c_glib/gandiva-glib 与 c_glib/test/gandiva,可作为跨语言调用与行为验证的参考。
【免费下载链接】arrowApache Arrow is the universal columnar format and multi-language toolbox for fast data interchange and in-memory analytics项目地址: https://gitcode.com/GitHub_Trending/arrow3/arrow
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考