☰
Apache DataFusion `Expr` 表达式全指南:从构造、求值到重写与内联 UDF
2026/9/25 3:41:06 网站建设 项目流程
  • 大数据
  • 数据分析
  • 后端

【免费下载链接】datafusion

Apache DataFusion SQL Query Engine

项目地址:https://gitcode.com/gh_mirrors/datafu/datafusion
点击查看免费下载

Expr(expression,表达式)是 Apache DataFusion 中最核心的抽象之一,用于表示一切计算逻辑,遵循主流编译器与数据库通用的“表达式树”(expression tree)模型。本篇技术指南将带你以库开发者的视角完整掌握 DataFusionExpr:理解其与 ArrowSchema/DFSchema的关系、编程式构造与求值方式、通过transform重写表达式,以及将自定义标量 UDF 内联为普通二元表达式的完整OptimizerRule实现与测试方案。读完本文,你将具备在 DataFusion 中自定义表达式变换、编写优化规则并对齐源码级实现原理的实战能力。

理解Expr:DataFusion 的表达式树抽象

Expr是 “expression” 的缩写,是 DataFusion 中表示一次计算的核心抽象。SQL 表达式a + b会被表示为一个BinaryExpr变体的Expr,其中包含左、右两个子Expr以及一个操作符(operator)。

作为第二个示例,SQL 表达式a + b * c同样被表示为BinaryExpr变体的Expr:其左侧子表达式为a,右侧子表达式又是一个BinaryExpr(b * c),构成经典的表达式树:

┌────────────────────┐ │ BinaryExpr │ │ op: + │ └────────────────────┘ ▲ ▲ ┌───────┘ └────────────────┐ │ │ ┌────────────────────┐ ┌────────────────────┐ │ Expr::Col │ │ BinaryExpr │ │ col: a │ │ op: * │ └────────────────────┘ └────────────────────┘ ▲ ▲ ┌────────┘ └─────────┐ │ │ ┌────────────────────┐ ┌────────────────────┐ │ Expr::Col │ │ Expr::Col │ │ col: b │ │ col: c │ └────────────────────┘ └────────────────────┘

从源码结构看,Expr枚举定义 包含约 30 个变体,覆盖了常见运算的全部形态:Column(列引用)、Literal(常量)、BinaryExpr(二元运算,如age > 21)、Cast/TryCast(类型转换)、ScalarFunction(标量函数调用)、AggregateFunction(聚合函数)、WindowFunction(窗口函数)、Case、Between、InList、Not/IsNull/IsNotNull等逻辑表达式,以及ScalarSubquery、InSubquery、Exists等子查询表达式。作为库开发者,你可以用Expr表示任何想要执行的计算,并借助 DataFusion 提供的整套 API 对其进行构造、求值、简化和分析。

Arrow Schema 与 DataFusion DFSchema

在深入Expr之前,需要先理解 DataFusion 中承载字段信息的两种 Schema 结构。

Schema 与 DFSchema 的区别

  • Schema(Arrow Schema):Apache Arrow 的底层组件,定义数据集的结构,指明列名及其数据类型。详细的 API 文档可参考 arrow-schema crate 的struct.Schema。
  • DFSchema(DataFusion Schema):在Schema基础上扩展,额外携带列限定符(column qualifiers)与函数依赖(functional dependencies)等信息。列限定符是到表的多段路径,例如table.schema.catalog;函数依赖描述表内各属性(特征)之间的关联关系。这在跨表查询管理中尤其有价值。

从 DFSchema 源码定义 可以确认,DFSchema内部持有三部分数据:底层的 ArrowSchemaRef(inner)、与字段一一对应的可选TableReference限定符列表(field_qualifiers)、以及存储函数依赖的FunctionalDependencies。这也是它与纯 ArrowSchema的本质差异所在。

Schema 与 DFSchema 的相互转换

从 Schema 到 DFSchema:使用DFSchema::try_from_qualified_schema,传入表名与原始 schema,即可得到带限定符的 DFSchema(可参考 docs.rs 上DFSchema文档中的 "creating-qualified-schemas" 示例)。

从 DFSchema 到 Schema:DFSchema实现了Into<Schema>trait(可通过as_arrow()直接拿到内部 Arrow Schema 引用,见 dfschema.rs 源码),因此转换非常直接(参考 "converting-back-to-arrow-schema" 示例)。

创建与求值Expr

DataFusion 提供了datafusion-examples下的注释详尽示例代码 expr_api.rs,完整演示了Expr的创建、求值、简化与分析:

  • 构造(fluent API):col("a") + lit(5)一行即可生成表达式;其等价的手写形式是Expr::BinaryExpr(BinaryExpr::new(Box::new(col("a")), Operator::Plus, Box::new(Expr::Literal(ScalarValue::Int32(Some(5)), None)))),两种方式断言相等。
  • expr_fn 聚合 API:first_value_udaf().call(vec![col("price")])生成first_value(price);配合ExprFunctionExttrait 还可构造带FILTER与ORDER BY的复杂聚合,如first_value(price) FILTER (WHERE quantity > 100) ORDER BY [ts DESC NULLS LAST]。
  • 求值:先把逻辑Expr通过SessionContext::new().create_physical_expr(expr, &df_schema)转换为物理表达式,再对RecordBatch调用evaluate,得到ColumnarValue。
  • 简化:通过SimplifyContext与ExprSimplifier可完成常量折叠(如ts = to_timestamp("2020-09-08T12:00:00+00:00")简化为与时间戳常量1599566400000000000的比较)、算术简化(i + (1 + 2)→i + 3)、逻辑简化(((i > 5) AND FALSE) OR (i < 10)→i < 10)以及字符串转日期简化(cast('2020-09-01' as date)→Date32(18506))。
  • 范围分析:对谓词表达式调用analyze可推导出列取值区间(如date > '2020-09-01' AND date < '2020-10-01'推导出区间['2020-09-01', '2020-10-01']),配合列统计信息(ColumnStatistics)还能估算选择率(selectivity),支撑剪枝与代价优化。

以标量 UDF 为载体的Expr实战

本指南以一个ScalarUDF表达式为例展开。实现 UDF 的完整教程见 adding-udfs.md(即仓库中的 “Adding User Defined Functions” 指南),本文直接沿用其中的add_one函数(对每个 i64 参数加 1)。

在已实现add_one函数的基础上,可以用create_udf创建Expr:

use std::sync::Arc; use datafusion::arrow::datatypes::DataType; use datafusion::logical_expr::{Volatility, create_udf}; use datafusion::logical_expr::{col, lit}; use datafusion::logical_expr::ColumnarValue; use datafusion::common::Result; pub fn add_one(args: &[ColumnarValue]) -> Result<ColumnarValue> { // Error handling omitted for brevity let args = ColumnarValue::values_to_arrays(args)?; let i64s = as_int64_array(&args[0])?; let new_array = i64s .iter() .map(|array_elem| array_elem.map(|value| value + 1)) .collect::<Int64Array>(); Ok(ColumnarValue::from(Arc::new(new_array) as ArrayRef)) } let add_one_udf = create_udf( "add_one", vec![DataType::Int64], DataType::Int64, Volatility::Immutable, Arc::new(add_one), ); // make the expr `add_one(5)` let expr = add_one_udf.call(vec![lit(5)]); // make the expr `add_one(my_column)` let expr = add_one_udf.call(vec![col("my_column")]);

关于create_udf的五个参数,结合 expr_fn.rs 源码 与 adding-udfs.md 的说明:

  • name:函数名,即 SQL 查询中使用的名称;
  • input_types:Vec<DataType>,函数接受的参数类型列表,此处为单个Int64;
  • return_type:函数返回类型,此处为Int64;
  • volatility:波动性(Volatility),决定优化器能否在某些场景下对该函数做优化。Immutable表示相同输入永远返回相同结果(如本函数);随机数生成器应标记为Volatile(相同输入也可能返回不同值);依赖外部环境的可标记Stable;
  • fun:函数实现本体。

从源码看,create_udf内部只是将参数包装进SimpleScalarUDF(expr_fn.rs):它把参数列表包装为Signature::exact(input_types, volatility),并在invoke_with_args中直接调用传入的函数闭包。若需要更灵活的功能(如多签名、类型强制、文档注解),则应直接实现ScalarUDFImpltrait(见 advanced_udf.rs 对应文档 中的 “Adding byimpl ScalarUDFImpl” 小节)。

获得ScalarUDF后,还需通过ctx.register_udf(add_one_udf)注册到SessionContext,SQL 中即可按名调用。

若想先系统了解Expr的各类构造函数(列引用、字面量、布尔/位运算、比较、算术、字符串、聚合、窗口等),可阅读 表达式用户指南。

重写Expr

表达式重写(Rewriting Expressions)是指将一个Expr变换为另一个Expr的过程。其典型动机包括:

  • 简化Expr,使其更易求值;
  • 优化Expr,使其求值更快;
  • 转换Expr形态,例如把BinaryExpr转为CastExpr。

除本文的示例外,仓库还提供以下重写与Expr操作示例:expr_api.rs、analyzer_rule.rs、optimizer_rule.rs。

在本文示例中,我们通过重写将add_oneUDF 更新为带字面量1的BinaryExpr,即把 UDF“内联”为普通算术表达式。

使用transform重写

要实现内联,需要编写一个接收Expr、返回Result<Expr>的函数:

  • 若表达式不需要重写,用Transformed::no包装原Expr;
  • 若表达式需要重写,用Transformed::yes包装新Expr。
use datafusion::common::Result; use datafusion::common::tree_node::{Transformed, TreeNode}; use datafusion::logical_expr::{col, lit, Expr}; use datafusion::logical_expr::{ScalarUDF}; fn rewrite_add_one(expr: Expr) -> Result<Transformed<Expr>> { expr.transform(&|expr| { Ok(match expr { Expr::ScalarFunction(scalar_func) if scalar_func.func.inner().name() == "add_one" => { let input_arg = scalar_func.args[0].clone(); let new_expression = input_arg + lit(1i64); Transformed::yes(new_expression) } _ => Transformed::no(expr), }) }) }

这里的expr.transform(...)来自TreeNodetrait(datafusion::common::tree_node::{Transformed, TreeNode}),它负责在整棵表达式树上递归地应用闭包,并把Transformed标志沿调用链向上传播(底层的TreeNodeAPI 设计可参考datafusion-common的tree_node模块)。模式匹配Expr::ScalarFunction(scalar_func)加守卫条件func.inner().name() == "add_one"用于识别 UDF 调用节点,取第一个参数scalar_func.args[0]后拼上+ lit(1i64)完成替换。

创建OptimizerRule

在 DataFusion 中,OptimizerRule是一个 trait,用于支持重写LogicalPlan各部分中出现的Expr,是 DataFusion “用 trait 实现驱动行为”设计哲学的又一体现。trait 的完整定义见 optimizer.rs 源码,包含两个核心方法:

  • name:返回规则名称;
  • rewrite:接收LogicalPlan与&dyn OptimizerConfig,返回Result<Transformed<LogicalPlan>>。规则若能优化计划,返回包装了优化后计划的Transformed::yes;否则返回Transformed::no。

下面实现名为AddOneInliner的规则:

use std::sync::Arc; use datafusion::common::Result; use datafusion::common::tree_node::{Transformed, TreeNode}; use datafusion::logical_expr::{col, lit, Expr, LogicalPlan, LogicalPlanBuilder}; use datafusion::optimizer::{OptimizerRule, OptimizerConfig, OptimizerContext, Optimizer}; fn rewrite_add_one(expr: Expr) -> Result<Transformed<Expr>> { expr.transform(&|expr| { Ok(match expr { Expr::ScalarFunction(scalar_func) if scalar_func.func.inner().name() == "add_one" => { let input_arg = scalar_func.args[0].clone(); let new_expression = input_arg + lit(1i64); Transformed::yes(new_expression) } _ => Transformed::no(expr), }) }) } #[derive(Default, Debug)] struct AddOneInliner {} impl OptimizerRule for AddOneInliner { fn name(&self) -> &str { "add_one_inliner" } fn rewrite( &self, plan: LogicalPlan, _config: &dyn OptimizerConfig, ) -> Result<Transformed<LogicalPlan>> { // Map over the expressions and rewrite them let new_expressions: Vec<Expr> = plan .expressions() .into_iter() .map(|expr| rewrite_add_one(expr)) .collect::<Result<Vec<_>>>()? // returns Vec<Transformed<Expr>> .into_iter() .map(|transformed| transformed.data) .collect(); let inputs = plan.inputs().into_iter().cloned().collect::<Vec<_>>(); let plan: Result<LogicalPlan> = plan.with_new_exprs(new_expressions, inputs); plan.map(|p| Transformed::yes(p)) } }

注意这里先通过plan.expressions()取出计划内所有表达式,将rewrite_add_one映射到每个表达式上,再用plan.with_new_exprs(new_expressions, inputs)构造携带重写后表达式的新LogicalPlan。with_new_exprs是LogicalPlan提供的标准重建接口:给定新表达式列表与子计划列表,返回等价但表达式已被替换的计划。

值得一提的进阶点:从 OptimizerRule 源码 看,现代规则还可以覆盖apply_order()来声明应用顺序(如ApplyOrder::BottomUp/TopDown),由优化器负责计划树的递归遍历;而本文示例是规则自身处理递归的经典写法。另外,源码注释明确建议规则尽量避免基于函数名的特判(如func.name() == "sum"),因为函数可能被覆盖导致语义不同(例如datafusion-sparkcrate 注册的sum),应优先使用ScalarUDFImpl/AggregateUDFImpl提供的方法。更完整的现代OptimizerRule示例(含supports_rewrite、apply_order、map_expressions用法)见 optimizer_rule.rs。

测试规则

测试规则相当直接:创建一个带有该规则的SessionState(或SessionContext),创建DataFrame并运行查询,逻辑计划会被规则优化。

use std::sync::Arc; use datafusion::common::Result; use datafusion::common::tree_node::{Transformed, TreeNode}; use datafusion::logical_expr::{col, lit, Expr, LogicalPlan, LogicalPlanBuilder}; use datafusion::optimizer::{OptimizerRule, OptimizerConfig, OptimizerContext, Optimizer}; use datafusion::arrow::array::{ArrayRef, Int64Array}; use datafusion::common::cast::as_int64_array; use datafusion::logical_expr::ColumnarValue; use datafusion::logical_expr::{Volatility, create_udf}; use datafusion::arrow::datatypes::DataType; use datafusion::execution::context::SessionContext; // ... rewrite_add_one 与 AddOneInliner 定义同上 ... pub fn add_one(args: &[ColumnarValue]) -> Result<ColumnarValue> { // Error handling omitted for brevity let args = ColumnarValue::values_to_arrays(args)?; let i64s = as_int64_array(&args[0])?; let new_array = i64s .iter() .map(|array_elem| array_elem.map(|value| value + 1)) .collect::<Int64Array>(); Ok(ColumnarValue::from(Arc::new(new_array) as ArrayRef)) } #[tokio::main] async fn main() -> Result<()> { let ctx = SessionContext::new(); // 取消注释下一行即可启用规则: // ctx.add_optimizer_rule(Arc::new(AddOneInliner {})); let add_one_udf = create_udf( "add_one", vec![DataType::Int64], DataType::Int64, Volatility::Immutable, Arc::new(add_one), ); ctx.register_udf(add_one_udf); let sql = "SELECT add_one(5) AS added_one"; // 若想对比未优化计划,可改用 into_unoptimized_plan(): // let plan = ctx.sql(sql).await?.into_unoptimized_plan().clone(); let plan = ctx.sql(sql).await?.into_optimized_plan()?.clone(); let expected = r#"Projection: Int64(6) AS added_one EmptyRelation: rows=1"#; assert_eq!(plan.to_string(), expected); Ok(()) }

启用规则后,计划被优化为如下形态,可见add_oneUDF 已被内联进投影:

Projection: add_one(Int64(5)) AS added_one -> Projection: Int64(5) + Int64(1) AS added_one -> Projection: Int64(6) AS added_one

即add_oneUDF 已被内联为Int64(5) + Int64(1),最终被常量折叠为Int64(6)。上例断言验证了ctx.sql("SELECT add_one(5) AS added_one").await?.into_optimized_plan()的输出恰好是Projection: Int64(6) AS added_one加EmptyRelation: rows=1——注意测试中规则处于注释状态,得到的是常量折叠后的结果;取消ctx.add_optimizer_rule(Arc::new(AddOneInliner {}))的注释后,即可观察到内联中间形态。完整可运行示例(含注册批数据、对 Filter 谓词重写与结果断言的模式)可参考 optimizer_rule.rs。

获取表达式的数据类型

表达式的arrow::datatypes::DataType可以通过调用get_type获得,前提是传入一个实现了Expr::Schemable的对象(例如DFSchema):

use arrow::datatypes::{DataType, Field}; use datafusion::common::DFSchema; use datafusion::logical_expr::{col, ExprSchemable}; use std::collections::HashMap; // Get the type of an expression that adds 2 columns. Adding an Int32 // and Float32 results in Float32 type let expr = col("c1") + col("c2"); let schema = DFSchema::from_unqualified_fields( vec![ Field::new("c1", DataType::Int32, true), Field::new("c2", DataType::Float32, true), ] .into(), HashMap::new(), ).unwrap(); assert_eq!("Float32", format!("{}", expr.get_type(&schema).unwrap()));

ExprSchemable::get_type需要 schema 是因为表达式的类型取决于输入表达式的类型:例如Int32 + Float32依类型提升规则得到Float32,col("c")在Utf8与Int32两种 schema 下分别得到Utf8与Int32(更多用例见 expr_api.rs 的 expression_type_demo)。

值得一提的相关能力是类型强制(type coercion):Expr直接转物理表达式求值时不会自动做类型提升(如Int8 > Int32会报 “Invalid comparison operation”),需要通过SessionContext::create_physical_expr、ExprSimplifier::coerce、TypeCoercionRewriter或手写transform显式处理(完整对比见 expr_api.rs 的 type_coercion_demo)。

结语

本指南系统展示了如何编程式创建Expr、如何重写它们(这对简化和优化Expr非常有用),以及如何编写测试验证自定义规则是否正常工作。掌握Expr及其背后的TreeNode遍历与OptimizerRule扩展机制,是深入 DataFusion 库开发——无论是自定义 UDF、表达式简化、还是计划优化——的必经之路。进一步的延伸阅读包括:表达式用户指南(Expr构造函数全集)、添加 UDF 指南(标量/窗口/聚合/表函数注册)、查询优化器指南(优化规则体系),以及datafusion-examples/examples/query_planning/目录下的 expr_api.rs、analyzer_rule.rs 与 optimizer_rule.rs 三个可运行示例。

  • 大数据
  • 数据分析
  • 后端

【免费下载链接】datafusion

Apache DataFusion SQL Query Engine

项目地址:https://gitcode.com/gh_mirrors/datafu/datafusion
点击查看免费下载
上一篇:探索Rnote:重新定义数字手写笔记的无限可能
下一篇:QtScrcpy快捷键导出功能:自定义配置备份与分享

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

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

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

立即咨询