Apache Arrow Gandiva 表达式、Projector 与 Filter 使用指南(C++)
2026/9/14 17:29:15 网站建设 项目流程

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构建表达式树、将其包装为ExpressionCondition,再通过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(以及MakeStringLiteralMakeBinaryLiteralMakeDecimalLiteralMakeNull)创建。

叶子节点可以被组合成更复杂的表达式树,官方文档给出了四类核心组合器:

组合器作用
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则接收conditionthen_nodeelse_noderesult_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 + 3Expression与表示x < 3Condition

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包装为ConditionMakeExpression(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::MakeFilter::Make的第一个参数都是SchemaPtr schema(见 projector.h 与 filter.h)。

延续上一节的expressioncondition,官方示例创建 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 一致;输出数组由Evaluatepool中分配,因此调用方无需预分配内存。

求值过滤:Filter::Evaluate 与 SelectionVector

Filter::Evaluate产出SelectionVector——一个匹配过滤条件的行索引向量。从源码看,SelectionVector 本质上是Arrow 整数数组的包装器,按位宽参数化(selection_vector.h),其Mode枚举包括MODE_NONEMODE_UINT16MODE_UINT32MODE_UINT64,并提供MakeInt16/MakeInt32等工厂方法(MakeInt16(max_slots, pool, &selection_vector)为最常见用法)。

SelectionVector 必须在传给Evaluate()之前初始化,初始化时需要确定两个关键参数:

  1. 位宽(bitwidth):决定它能容纳的最大索引值;
  2. 最大槽位数(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 还支持在投影的同时应用选择向量,即带过滤的投影。关键前提有两个:

  1. 初始化 Projector 时,传入SelectionVector::GetMode(),使投影器按正确的位宽编译
  2. 求值时把 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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询