库用户指南

扩展 SQL 语法

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

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

扩展 SQL 语法

DataFusion 提供了一套灵活的扩展机制,让你无需修改核心代码即可自定义 SQL 的解析与规划。当你有以下需求时,它会非常有用:

  • 支持其他 SQL 方言中的自定义操作符(例如 PostgreSQL 用于 JSON 的 ->)
  • 添加原生不支持的自定义数据类型
  • 实现 TABLESAMPLE、PIVOT/UNPIVOT 或 MATCH_RECOGNIZE 之类的 SQL 结构

关于这一主题,你可以阅读博文《在 DataFusion 中扩展 SQL:从 ->> 到 TABLESAMPLE》了解更多内容。

架构概览

DataFusion 在处理一条 SQL 查询时,会经历以下几个阶段:

┌─────────────┐    ┌─────────┐    ┌──────────────────────┐    ┌─────────────┐
│ SQL String  │───▶│ Parser  │───▶│      SqlToRel        │───▶│ LogicalPlan │
└─────────────┘    └─────────┘    │ (SQL to LogicalPlan) │    └─────────────┘
                                  └──────────────────────┘
                                              │
                                              │ uses
                                              ▼
                                  ┌───────────────────────┐
                                  │  Extension Planners   │
                                  │  • ExprPlanner        │
                                  │  • TypePlanner        │
                                  │  • RelationPlanner    │
                                  └───────────────────────┘

扩展规划器会在 SqlToRel 阶段拦截 SQL AST 中的特定部分,让你可以自定义它们转换为 DataFusion 逻辑计划的方式。

扩展点

DataFusion 提供了三个用于扩展 SQL 的规划器 trait:

Trait用途注册方式
ExprPlanner自定义表达式和运算符ctx.register_expr_planner()
TypePlanner自定义 SQL 数据类型SessionStateBuilder::with_type_planner()
RelationPlanner自定义 FROM 子句元素(关系)ctx.register_relation_planner()

规划器优先级:可以注册多个 ExprPlanner 和 RelationPlanner;它们会按注册顺序的逆序调用(最后注册的优先)。返回 Original(...) 可将处理权委托给下一个规划器。同一时间只能有一个 TypePlanner 处于激活状态。

ExprPlanner:自定义表达式和运算符

使用 ExprPlanner 可以自定义 SQL 表达式转换为 DataFusion 逻辑表达式的方式。它适用于以下场景:

  • 自定义二元运算符(例如 ->、->>、@>、?)
  • 自定义字段访问模式
  • 自定义聚合函数或窗口函数的处理

可用方法

类别方法
运算符plan_binary_op、plan_any
字面量plan_array_literal、plan_dictionary_literal、plan_struct_literal
函数plan_extract、plan_substring、plan_overlay、plan_position、plan_make_map
标识符plan_field_access、plan_compound_identifier
聚合/窗口plan_aggregate、plan_window

完整的方法签名请参阅 ExprPlanner API 文档。

示例:自定义箭头运算符

这个示例将 -> 运算符映射为字符串拼接:

use datafusion_expr::planner::{ExprPlanner, PlannerResult, RawBinaryExpr};

#[derive(Debug)]
struct MyCustomPlanner;

impl ExprPlanner for MyCustomPlanner {
    fn plan_binary_op(
        &self,
        expr: RawBinaryExpr,
        _schema: &DFSchema,
    ) -> Result<PlannerResult<RawBinaryExpr>> {
        match &expr.op {
            // Map `->` to string concatenation
            BinaryOperator::Arrow => {
                Ok(PlannerResult::Planned(Expr::BinaryExpr(BinaryExpr {
                    left: Box::new(expr.left.clone()),
                    right: Box::new(expr.right.clone()),
                    op: Operator::StringConcat,
                })))
            }
            _ => Ok(PlannerResult::Original(expr)),
        }
    }
}

#[tokio::main]
async fn main() -> Result<()> {
    // Use postgres dialect to enable `->` operator parsing
    let config = SessionConfig::new()
        .set_str("datafusion.sql_parser.dialect", "postgres");
    let mut ctx = SessionContext::new_with_config(config);

    // Register the custom planner
    ctx.register_expr_planner(Arc::new(MyCustomPlanner))?;

    // Now `->` works as string concatenation
    let results = ctx.sql("SELECT 'hello'->'world'").await?.collect().await?;
    // Returns: "helloworld"
    Ok(())
}

更多详情,请参阅 ExprPlanner API 文档以及 expr_planner 测试示例。

TypePlanner:自定义数据类型

使用 TypePlanner 可以将 SQL 数据类型映射为 Arrow/DataFusion 类型。当你需要支持未被原生识别的 SQL 类型时,这一功能非常有用。

示例:自定义 DATETIME 类型

use datafusion_expr::planner::TypePlanner;

#[derive(Debug)]
struct MyTypePlanner;

impl TypePlanner for MyTypePlanner {
    fn plan_type_field(&self, sql_type: &ast::DataType) -> Result<Option<FieldRef>> {
        match sql_type {
            // Map DATETIME(precision) to Arrow Timestamp
            ast::DataType::Datetime(precision) => {
                let time_unit = match precision {
                    Some(0) => TimeUnit::Second,
                    Some(3) => TimeUnit::Millisecond,
                    Some(6) => TimeUnit::Microsecond,
                    None | Some(9) => TimeUnit::Nanosecond,
                    _ => return Ok(None), // Let default handling take over
                };
                Ok(Some(
                    DataType::Timestamp(time_unit, None).into_nullable_field_ref()
                ))
            }
            _ => Ok(None), // Return None for types we don't handle
        }
    }
}

#[tokio::main]
async fn main() -> Result<()> {
    let state = SessionStateBuilder::new()
        .with_default_features()
        .with_type_planner(Arc::new(MyTypePlanner))
        .build();

    let ctx = SessionContext::new_with_state(state);

    // Now DATETIME type is recognized
    ctx.sql("CREATE TABLE events (ts DATETIME(3))").await?;
    Ok(())
}

示例:支持 UUID 类型

use datafusion_expr::planner::TypePlanner;

#[derive(Debug)]
struct MyTypePlanner;

impl TypePlanner for MyTypePlanner {
    fn plan_type_field(&self, sql_type: &ast::DataType) -> Result<Option<FieldRef>> {
        match sql_type {
            sqlparser::ast::DataType::Uuid => Ok(Some(Arc::new(
                Field::new("", DataType::FixedSizeBinary(16), true).with_metadata(
                    [("ARROW:extension:name".to_string(), "arrow.uuid".to_string())]
                        .into(),
                ),
            ))),
            _ => Ok(None),
        }
    }
}

#[tokio::main]
async fn main() -> Result<()> {
    let state = SessionStateBuilder::new()
        .with_default_features()
        .with_type_planner(Arc::new(MyTypePlanner))
        .build();

    let ctx = SessionContext::new_with_state(state);

    // Now UUID type is recognized
    ctx.sql("CREATE TABLE idx (uuid UUID)").await?;
    Ok(())
}

有关更多细节,请参阅 TypePlanner API 文档。

RelationPlanner:自定义 FROM 子句元素

使用 RelationPlanner 来处理 FROM 子句中的自定义关系(relation)。这使你能够实现诸如以下的 SQL 结构:

  • TABLESAMPLE 用于数据采样
  • PIVOT / UNPIVOT 用于数据重塑
  • MATCH_RECOGNIZE 用于模式匹配
  • sqlparser 解析的任何自定义关系语法

RelationPlannerContext

在实现 RelationPlanner 时,你会收到一个 RelationPlannerContext,它提供了用于规划的工具方法:

方法用途
plan(relation)递归规划嵌套关系
get_cte(name)查找当前查询作用域中可见的 CTE
sql_to_expr(expr, schema)将 SQL 表达式转换为 DataFusion Expr
context_provider()访问会话配置、表、函数

有关 normalize_ident() 和 object_name_to_table_reference() 等其他方法,请参阅 RelationPlanner API 文档。

关系规划器会在 DataFusion 内置的 CTE 和目录查找之前被调用。如果你的规划器是按名称识别普通关系的,那么在确定某个关系是普通名称(且该名称原本会被你的规划器处理)之后,请检查 get_cte()。被有意保留的名称和自定义语法可以维持其原有的优先级。当某个 CTE 可见时,返回 RelationPlanning::Original,以便下一个规划器来处理该关系。当链中的每个基于名称识别普通关系的规划器都执行此检查后,规划流程便会到达 DataFusion 内置的 CTE 解析。请先对解析出的名称进行规范化,使查找遵循 DataFusion 现有的名称匹配规则。此操作仅适用于不带函数参数的表:名为 numbers 的 CTE 会遮蔽 numbers,但不会遮蔽 numbers(10) 这样的调用。

// Run this after recognizing `relation` as an ordinary name handled here.
if let TableFactor::Table { name, args: None, .. } = &relation {
    let table_ref = ctx.object_name_to_table_reference(name.clone())?;
    if ctx.get_cte(&table_ref).is_some() {
        return Ok(RelationPlanning::Original(Box::new(relation)));
    }
}

实现策略

在实现 RelationPlanner 时,主要有两种方式:

  1. 重写为标准 SQL:将自定义语法转换为 DataFusion 已经知道如何执行的等价标准操作(例如,PIVOT → 使用 CASE 表达式的 GROUP BY)。只要可行,这是最简单的方式。
  2. 自定义逻辑与物理节点:创建一个 UserDefinedLogicalNode 来在逻辑计划中表示该操作,并配合一个自定义的 ExecutionPlan 来执行它。要实现端到端执行,两者缺一不可。

示例:基本的 RelationPlanner 结构

use datafusion_expr::planner::{
    PlannedRelation, RelationPlanner, RelationPlannerContext, RelationPlanning,
};
use datafusion_sql::sqlparser::ast::TableFactor;

#[derive(Debug)]
struct MyRelationPlanner;

impl RelationPlanner for MyRelationPlanner {
    fn plan_relation(
        &self,
        relation: TableFactor,
        ctx: &mut dyn RelationPlannerContext,
    ) -> Result<RelationPlanning> {
        match relation {
            // Handle your custom relation
            TableFactor::Pivot { table, alias, .. } => {
                // Plan the input table
                let input = ctx.plan(*table)?;

                // Transform or wrap the plan as needed
                // ...

                Ok(RelationPlanning::Planned(PlannedRelation::new(input, alias)))
            }

            // Return Original for relations you don't handle
            other => Ok(RelationPlanning::Original(other)),
        }
    }
}

#[tokio::main]
async fn main() -> Result<()> {
    let ctx = SessionContext::new();

    // Register the custom planner
    ctx.register_relation_planner(Arc::new(MyRelationPlanner))?;

    Ok(())
}

完整示例

DataFusion 仓库中包含了演示每种方式的完整示例:

TABLESAMPLE(自定义逻辑节点与物理节点)

table_sample.rs 示例展示了如何支持如下查询的端到端完整实现:

SELECT * FROM table TABLESAMPLE BERNOULLI(10 PERCENT) REPEATABLE(42)

PIVOT/UNPIVOT(重写策略)

pivot_unpivot.rs 示例演示了如何将自定义语法重写为标准 SQL,以处理如下查询:

SELECT * FROM sales
  PIVOT (SUM(amount) FOR quarter IN ('Q1', 'Q2', 'Q3', 'Q4'))

回顾

  1. 使用 ExprPlanner 处理自定义运算符和表达式
  2. 使用 TypePlanner 处理自定义 SQL 数据类型
  3. 使用 RelationPlanner 处理自定义 FROM 子句语法(TABLESAMPLE、PIVOT 等)
  4. 通过 SessionContext 或 SessionStateBuilder 注册规划器

另请参阅

评论

登录后参与评论

正在加载评论…