扩展 SQL 语法
扩展 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 时,主要有两种方式:
- 重写为标准 SQL:将自定义语法转换为 DataFusion 已经知道如何执行的等价标准操作(例如,PIVOT → 使用 CASE 表达式的 GROUP BY)。只要可行,这是最简单的方式。
- 自定义逻辑与物理节点:创建一个
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'))回顾
- 使用
ExprPlanner处理自定义运算符和表达式 - 使用
TypePlanner处理自定义 SQL 数据类型 - 使用
RelationPlanner处理自定义 FROM 子句语法(TABLESAMPLE、PIVOT 等) - 通过
SessionContext或SessionStateBuilder注册规划器
另请参阅
- API 文档:
ExprPlanner、TypePlanner、RelationPlanner - relation_planner 示例 - 完整的 TABLESAMPLE、PIVOT/UNPIVOT 实现
- expr_planner 测试示例 - 自定义运算符示例
- UDF 指南中的自定义表达式规划
评论
登录后参与评论
KnowForge