处理 Expr
使用 Expr
Expr 是 "expression"(表达式)的缩写。它是 DataFusion 中用于表示计算的核心抽象,遵循大多数编译器和数据库中常见的"表达式树"抽象。
例如,SQL 表达式 a + b 会表示为一个 Expr,其变体为 BinaryExpr。BinaryExpr 包含左右两个 Expr 以及一个运算符。
再比如,SQL 表达式 a + b * c 会表示为一个 Expr,其变体为 BinaryExpr。左侧的 Expr 是 a,右侧的 Expr 则是另一个 BinaryExpr,它的左侧 Expr 是 b,右侧 Expr 是 c。作为一个经典的表达式树,它看起来是这样的:
┌────────────────────┐
│ BinaryExpr │
│ op: + │
└────────────────────┘
▲ ▲
┌───────┘ └────────────────┐
│ │
┌────────────────────┐ ┌────────────────────┐
│ Expr::Col │ │ BinaryExpr │
│ col: a │ │ op: * │
└────────────────────┘ └────────────────────┘
▲ ▲
┌────────┘ └─────────┐
│ │
┌────────────────────┐ ┌────────────────────┐
│ Expr::Col │ │ Expr::Col │
│ col: b │ │ col: c │
└────────────────────┘ └────────────────────┘作为库的作者,你可以使用 Expr 来表示你想要执行的计算。本指南将引导你完成如何以 Expr 的形式创建自己的标量 UDF,以及如何重写 Expr 以将简单 UDF 内联展开。
Arrow Schema 与 DataFusion DFSchema
Apache Arrow 的 Schema 提供了一种轻量级结构来定义数据,而 Apache DataFusion 的 DFSchema 在其基础上扩展了额外信息,例如列限定符和函数依赖。列限定符是指向表的多段路径,例如表、模式(schema)、目录(catalog)。函数依赖是指表中各属性(特征)之间的相互关系。
Schema 与 DFSchema 的区别
Schema:作为 Apache Arrow 的基础组件,
Schema定义了数据集的结构,指定列名及其数据类型。有关 Arrow Schema 的详细文档,请参阅 Struct Schema。
DFSchema:
DFSchema在Schema的基础上扩展,纳入了表名等限定符,使其在需要时能够携带额外的上下文信息。这对于管理跨多表的查询尤为重要。有关 DFSchema 的详细文档,请参阅 Struct DFSchema。
如何在 Schema 与 DFSchema 之间转换
从 Schema 转换为 DFSchema:使用 DFSchema::try_from_qualified_schema,传入表名和原始 schema,详细的代码示例请参阅 creating-qualified-schemas。
从 DFSchema 转换为 Schema:由于已为 DFSchema 实现了 Into trait,可将其转换为 Arrow Schema,详细的代码示例请参阅 converting-back-to-arrow-schema。
创建与求值 Expr
有关创建、求值、简化和分析 Expr 的带详细注释的代码,请参阅 expr_api.rs。
标量 UDF 示例
我们将以 ScalarUDF 表达式作为示例。这需要实现一个真正的 UDF,为方便起见,我们将沿用添加 UDF 指南中的同一个示例。
那么,假设你已经写好了这个函数,就可以用它来创建一个 Expr:
use datafusion::logical_expr::{Volatility, create_udf};
use datafusion::arrow::datatypes::DataType;
use datafusion::logical_expr::{col, lit};
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")]);如果你想在深入了解创建和重写 Expr 的细节之前先了解更多关于 Expr 的内容,可以阅读表达式用户指南。
重写 Expr
以下是几个重写和使用 Expr 的示例:
重写表达式是指将一个 Expr 转换为另一个 Expr 的过程。这样做有很多好处,包括:
- 简化
Expr,使其更易于求值 - 优化
Expr,使其求值速度更快 - 将
Expr转换为其他形式,例如将BinaryExpr转换为CastExpr
在我们的示例中,我们将使用重写来更新 add_one UDF,将其重写为一个包含字面量 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),
})
})
}创建 OptimizerRule
在 DataFusion 中,OptimizerRule 是一个 trait,用于改写出现在 LogicalPlan 各个部分中的 Expr。它遵循 DataFusion 一贯的理念:通过 trait 实现来驱动行为。
我们将这个规则命名为 AddOneInliner,并实现 OptimizerRule trait。OptimizerRule trait 包含两个方法:
name——返回该规则的名称rewrite——接收一个LogicalPlan和一个&dyn OptimizerConfig,并返回Result<Transformed<LogicalPlan>>。如果该规则能够优化计划,则返回带有优化后计划的Transformed::yes;如果无法优化计划,则返回Transformed::no。
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};
#[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))
}
}注意这里使用了 rewrite_add_one,它被映射到 plan.expressions() 上以重写这些表达式,然后使用 plan.with_new_exprs 以重写后的表达式创建一个新的 LogicalPlan。
我们快要完成了。接下来只需测试我们的规则是否正常工作。
测试规则
测试规则相当简单:我们可以创建一个包含该规则的 SessionState,然后创建一个 DataFrame 并运行查询。逻辑计划将会被我们的规则所优化。
use datafusion::execution::context::SessionContext;
#[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";
// 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(())
}该计划优化后为:
Projection: add_one(Int64(5)) AS added_one
-> Projection: Int64(5) + Int64(1) AS added_one
-> Projection: Int64(6) AS added_one也就是说 add_one UDF 已被内联到投影中。
获取表达式的数据类型
表达式的 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()));结论
在本指南中,我们介绍了如何以编程方式创建 Expr 以及如何重写它们。这对于简化和优化 Expr 非常有用。我们还介绍了如何测试我们的规则,以确保它能够正常工作。
评论
登录后参与评论
KnowForge