库用户指南

处理 Expr

师成师成· 更新于 2026-09-28· 阅读 20 分钟· 0 次阅读

登录后可跨设备保存划线和私人笔记登录

使用 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 非常有用。我们还介绍了如何测试我们的规则,以确保它能够正常工作。

评论

登录后参与评论

正在加载评论…