ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

Apache DataFusion `Expr` 表达式全指南:从构造、求值到重写与内联 UDF

Apache DataFusion `Expr` 表达式全指南:从构造、求值到重写与内联 UDF 大数据数据分析后端【免费下载链接】datafusionApache DataFusion SQL Query Engine项目地址https://gitcode.com/gh_mirrors/datafu/datafusion点击查看免费下载Exprexpression表达式是 Apache DataFusion 中最核心的抽象之一用于表示一切计算逻辑遵循主流编译器与数据库通用的“表达式树”expression tree模型。本篇技术指南将带你以库开发者的视角完整掌握 DataFusionExpr理解其与 ArrowSchema/DFSchema的关系、编程式构造与求值方式、通过transform重写表达式以及将自定义标量 UDF 内联为普通二元表达式的完整OptimizerRule实现与测试方案。读完本文你将具备在 DataFusion 中自定义表达式变换、编写优化规则并对齐源码级实现原理的实战能力。理解ExprDataFusion 的表达式树抽象Expr是 “expression” 的缩写是 DataFusion 中表示一次计算的核心抽象。SQL 表达式a b会被表示为一个BinaryExpr变体的Expr其中包含左、右两个子Expr以及一个操作符operator。作为第二个示例SQL 表达式a b * c同样被表示为BinaryExpr变体的Expr其左侧子表达式为a右侧子表达式又是一个BinaryExprb * 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 的区别SchemaArrow SchemaApache Arrow 的底层组件定义数据集的结构指明列名及其数据类型。详细的 API 文档可参考 arrow-schema crate 的struct.Schema。DFSchemaDataFusion Schema在Schema基础上扩展额外携带列限定符column qualifiers与函数依赖functional dependencies等信息。列限定符是到表的多段路径例如table.schema.catalog函数依赖描述表内各属性特征之间的关联关系。这在跨表查询管理中尤其有价值。从 DFSchema 源码定义 可以确认DFSchema内部持有三部分数据底层的 ArrowSchemaRefinner、与字段一一对应的可选TableReference限定符列表field_qualifiers、以及存储函数依赖的FunctionalDependencies。这也是它与纯 ArrowSchema的本质差异所在。Schema 与 DFSchema 的相互转换从 Schema 到 DFSchema使用DFSchema::try_from_qualified_schema传入表名与原始 schema即可得到带限定符的 DFSchema可参考 docs.rs 上DFSchema文档中的 creating-qualified-schemas 示例。从 DFSchema 到 SchemaDFSchema实现了IntoSchematrait可通过as_arrow()直接拿到内部 Arrow Schema 引用见 dfschema.rs 源码因此转换非常直接参考 converting-back-to-arrow-schema 示例。创建与求值ExprDataFusion 提供了datafusion-examples下的注释详尽示例代码 expr_api.rs完整演示了Expr的创建、求值、简化与分析构造fluent APIcol(a) lit(5)一行即可生成表达式其等价的手写形式是Expr::BinaryExpr(BinaryExpr::new(Box::new(col(a)), Operator::Plus, Box::new(Expr::Literal(ScalarValue::Int32(Some(5)), None))))两种方式断言相等。expr_fn 聚合 APIfirst_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:0000: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创建Expruse 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]) - ResultColumnarValue { // 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_typesVecDataType函数接受的参数类型列表此处为单个Int64return_type函数返回类型此处为Int64volatility波动性Volatility决定优化器能否在某些场景下对该函数做优化。Immutable表示相同输入永远返回相同结果如本函数随机数生成器应标记为Volatile相同输入也可能返回不同值依赖外部环境的可标记Stablefun函数实现本体。从源码看create_udf内部只是将参数包装进SimpleScalarUDFexpr_fn.rs它把参数列表包装为Signature::exact(input_types, volatility)并在invoke_with_args中直接调用传入的函数闭包。若需要更灵活的功能如多签名、类型强制、文档注解则应直接实现ScalarUDFImpltrait见 advanced_udf.rs 对应文档 中的 “Adding byimpl ScalarUDFImpl” 小节。获得ScalarUDF后还需通过ctx.register_udf(add_one_udf)注册到SessionContextSQL 中即可按名调用。若想先系统了解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、返回ResultExpr的函数若表达式不需要重写用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) - ResultTransformedExpr { 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(...)来自TreeNodetraitdatafusion::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返回ResultTransformedLogicalPlan。规则若能优化计划返回包装了优化后计划的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) - ResultTransformedExpr { 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, ) - ResultTransformedLogicalPlan { // Map over the expressions and rewrite them let new_expressions: VecExpr plan .expressions() .into_iter() .map(|expr| rewrite_add_one(expr)) .collect::ResultVec_()? // returns VecTransformedExpr .into_iter() .map(|transformed| transformed.data) .collect(); let inputs plan.inputs().into_iter().cloned().collect::Vec_(); let plan: ResultLogicalPlan 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]) - ResultColumnarValue { // 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: rows1#; 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: rows1——注意测试中规则处于注释状态得到的是常量折叠后的结果取消ctx.add_optimizer_rule(Arc::new(AddOneInliner {}))的注释后即可观察到内联中间形态。完整可运行示例含注册批数据、对 Filter 谓词重写与结果断言的模式可参考 optimizer_rule.rs。获取表达式的数据类型表达式的arrow::datatypes::DataType可以通过调用get_type获得前提是传入一个实现了Expr::Schemable的对象例如DFSchemause 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依类型提升规则得到Float32col(c)在Utf8与Int32两种 schema 下分别得到Utf8与Int32更多用例见 expr_api.rs 的 expression_type_demo。值得一提的相关能力是类型强制type coercionExpr直接转物理表达式求值时不会自动做类型提升如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 三个可运行示例。赞分享大数据数据分析后端【免费下载链接】datafusionApache DataFusion SQL Query Engine项目地址https://gitcode.com/gh_mirrors/datafu/datafusion点击查看免费下载相关推荐PointNet架构深度解析揭秘多尺度分组与层次化特征提取PointNet架构深度解析揭秘多尺度分组与层次化特征提取 PointNet是斯坦福大学提出的深度学习模型专门用于3D点云数据的分类和分割任务。作为Apache DataFusion Expr API 实战指南在 Rust 中构建、求值与优化逻辑表达式Apache DataFusion Expr API 实战指南在 Rust 中构建、求值与优化逻辑表达式 DataFusion 的 DataFrame 方法大数据数据分析后端Apache Arrow DataFusion表达式(Expr)操作完全指南Apache Arrow DataFusion表达式 Expr 操作完全指南 还在为复杂的数据处理表达式而头疼吗Apache Arrow DataFusion上一篇探索Rnote重新定义数字手写笔记的无限可能下一篇QtScrcpy快捷键导出功能自定义配置备份与分享创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表