查询优化器
查询优化器
DataFusion 是一个用 Rust 编写的可扩展查询执行框架,使用 Apache Arrow 作为其内存格式。
DataFusion 采用模块化设计,各个 crate 可以在其他项目中重复使用。
该 crate 是 DataFusion 的一个子模块,为逻辑计划提供查询优化器,并包含大量 OptimizerRule 和 PhysicalOptimizerRule,它们可以重写计划和/或其表达式,使计划在计算结果不变的前提下执行得更快。
关于内置分析器规则、逻辑优化器规则和物理优化器规则的参考列表,请参阅优化器规则参考。
如需更深入地了解优化器架构、规则类型和谓词,请参阅《在 DataFusion 中优化 SQL(和 DataFrame),第 1 部分》、《第 2 部分》、《在 Apache DataFusion 中利用排序信息获得更优计划》以及《动态过滤器:在执行过程中于算子之间传递信息,让查询快 25 倍》。
运行优化器
以下代码演示了基本流程:使用默认优化规则集创建优化器,并将其应用于逻辑计划,从而得到优化后的逻辑计划。
use std::sync::Arc;
use datafusion::logical_expr::{col, lit, LogicalPlan, LogicalPlanBuilder};
use datafusion::optimizer::{OptimizerRule, OptimizerContext, Optimizer};
// We need a logical plan as the starting point. There are many ways to build a logical plan:
//
// The `datafusion-expr` crate provides a LogicalPlanBuilder
// The `datafusion-sql` crate provides a SQL query planner that can create a LogicalPlan from SQL
// The `datafusion` crate provides a DataFrame API that can create a LogicalPlan
let initial_logical_plan = LogicalPlanBuilder::empty(false).build().unwrap();
// use builtin rules or customized rules
let rules: Vec<Arc<dyn OptimizerRule + Send + Sync>> = vec![];
let optimizer = Optimizer::with_rules(rules);
let config = OptimizerContext::new().with_max_passes(16);
let optimized_plan = optimizer.optimize(initial_logical_plan.clone(), &config, observer);
fn observer(plan: &LogicalPlan, rule: &dyn OptimizerRule) {
println!(
"After applying rule '{}':\n{}",
rule.name(),
plan.display_indent()
)
}编写优化规则
请参考 optimizer_rule.rs 示例,了解编写优化器规则的一般方法,然后再研究现有的规则。
OptimizerRule 将一个 LogicalPlan 转换为另一个能够计算出相同结果、但可能效率更高的逻辑计划。如果没有适合输入计划的转换,优化器可以直接原样返回该计划。
所有规则都必须实现 OptimizerRule trait。
#[derive(Default, Debug)]
struct MyOptimizerRule {}
impl OptimizerRule for MyOptimizerRule {
fn name(&self) -> &str {
"my_optimizer_rule"
}
fn rewrite(
&self,
plan: LogicalPlan,
_config: &dyn OptimizerConfig,
) -> Result<Transformed<LogicalPlan>> {
unimplemented!()
}
}提供自定义规则
优化器可以使用自定义的规则集来创建。
let optimizer = Optimizer::with_rules(vec![
Arc::new(MyOptimizerRule {})
]);通用准则
规则通常会遍历逻辑计划,并遍历运算符内部的表达式树,然后有选择地修改单个运算符或表达式。
有时会先进行一次初始遍历访问计划并构建状态,随后在第二次遍历中执行实际的优化。这种方法用于投影下推和过滤下推。
表达式命名
DataFusion 中的每个表达式都有一个名称,该名称用作列名。例如,在下面这个示例中,输出包含一个名为 "COUNT(aggregate_test_100.c9)" 的单列:
> select count(c9) from aggregate_test_100;
+------------------------------+
| COUNT(aggregate_test_100.c9) |
+------------------------------+
| 100 |
+------------------------------+这些名称既用于指代两个子查询中的列,也用于 LogicalPlan 从一个阶段到另一个阶段的内部引用。例如:
> select "COUNT(aggregate_test_100.c9)" + 1 from (select count(c9) from aggregate_test_100) as sq;
+--------------------------------------------+
| sq.COUNT(aggregate_test_100.c9) + Int64(1) |
+--------------------------------------------+
| 101 |
+--------------------------------------------+推论
由于 DataFusion 使用字符串名称来标识列,因此优化器在重写表达式时绝不能改变表达式的名称,这一点至关重要。通常的做法是为重写后的表达式添加别名来重命名。
下面是一个简单的重写示例。表达式 1 + 2 在内部可以被简化为 3,但对外仍必须显示为与 1 + 2 相同的形式:
> select 1 + 2;
+---------------------+
| Int64(1) + Int64(2) |
+---------------------+
| 3 |
+---------------------+查看 EXPLAIN 的输出可以看到,优化器实际上已将 1 + 2 重写为等价于 3 as "1 + 2" 的形式:
> explain format indent select 1 + 2;
+---------------+-------------------------------------------------+
| plan_type | plan |
+---------------+-------------------------------------------------+
| logical_plan | Projection: Int64(3) AS Int64(1) + Int64(2) |
| | EmptyRelation |
| physical_plan | ProjectionExec: expr=[3 as Int64(1) + Int64(2)] |
| | PlaceholderRowExec |
| | |
+---------------+-------------------------------------------------+如果表达式名称未被保留,就会出现 devlive-community/knowforge#3704 和 devlive-community/knowforge#3555 这类缺陷,导致找不到预期的列。
构建表达式名称
目前有两种方式可以为逻辑计划中的表达式创建名称。
impl Expr {
/// Returns the name of this expression as it should appear in a schema. This name
/// will not include any CAST expressions.
pub fn display_name(&self) -> Result<String> {
Ok("display_name".to_string())
}
/// Returns a full and complete string representation of this expression.
pub fn canonical_name(&self) -> String {
"canonical_name".to_string()
}
}在比较表达式以确定它们是否等价时,应使用 canonical_name;而在创建用于模式(schema)中的名称时,应使用 display_name。
工具方法
提供了许多工具方法,用于处理一些常见任务。
递归遍历表达式树
TreeNode API 提供了一种便捷的方式,可以递归遍历表达式树或计划树。
例如,要查找逻辑计划中的所有子查询引用,可以使用以下代码:
// Return all subquery references in an expression
fn extract_subquery_filters(expression: &Expr) -> Result<Vec<&Expr>> {
let mut extracted = vec![];
expression.apply(|expr| {
if let Expr::InSubquery(_) = expr {
extracted.push(expr);
}
Ok(TreeNodeRecursion::Continue)
})?;
Ok(extracted)
}同样,你可以使用 TreeNode API 来重写 LogicalPlan 或 ExecutionPlan。
// Return all joins in a logical plan
fn find_joins(overall_plan: &LogicalPlan) -> Result<Vec<&Join>> {
let mut extracted = vec![];
overall_plan.apply(|plan| {
if let LogicalPlan::Join(join) = plan {
extracted.push(join);
}
Ok(TreeNodeRecursion::Continue)
})?;
Ok(extracted)
}重写表达式
TreeNode API 同样提供了一种便捷的方式来重写表达式和计划。例如,要重写所有如下形式的表达式
col BETWEEN x AND y进入
col >= x AND col <= y你可以使用以下代码:
// Recursively rewrite all BETWEEN expressions
// returns Transformed::yes if any changes were made
fn rewrite_between(expr: Expr) -> Result<Transformed<Expr>> {
// transform_up does a bottom up rewrite
expr.transform_up(|expr| {
// only handle BETWEEN expressions
let Expr::Between(Between {
negated,
expr,
low,
high,
}) = expr else {
return Ok(Transformed::no(expr))
};
let rewritten_expr = if negated {
// don't rewrite NOT BETWEEN
Expr::Between(Between::new(expr, negated, low, high))
} else {
// rewrite to (expr >= low) AND (expr <= high)
expr.clone().gt_eq(*low).and(expr.lt_eq(*high))
};
Ok(Transformed::yes(rewritten_expr))
})
}编写测试
新规则所在的同一文件中应包含单元测试,用于单独测试该规则应用于计划后的效果(不应用其他任何规则)。
integration-tests.rs 中还应有一个测试,用于在整体优化流程的背景下测试该规则。
调试
可以使用 EXPLAIN VERBOSE 命令查看每条优化规则对查询产生的效果。
在下面的示例中,type_coercion 和 simplify_expressions 阶段对计划进行了简化,使其直接返回常量 "3.2",而不是在执行时进行计算。
> explain verbose select cast(1 + 2.2 as string) as foo;
+------------------------------------------------------------+---------------------------------------------------------------------------+
| plan_type | plan |
+------------------------------------------------------------+---------------------------------------------------------------------------+
| initial_logical_plan | Projection: CAST(Int64(1) + Float64(2.2) AS Utf8) AS foo |
| | EmptyRelation |
| logical_plan after type_coercion | Projection: CAST(CAST(Int64(1) AS Float64) + Float64(2.2) AS Utf8) AS foo |
| | EmptyRelation |
| logical_plan after simplify_expressions | Projection: Utf8("3.2") AS foo |
| | EmptyRelation |
| logical_plan after unwrap_cast_in_comparison | SAME TEXT AS ABOVE |
| logical_plan after decorrelate_where_exists | SAME TEXT AS ABOVE |
| logical_plan after decorrelate_where_in | SAME TEXT AS ABOVE |
| logical_plan after scalar_subquery_to_join | SAME TEXT AS ABOVE |
| logical_plan after subquery_filter_to_join | SAME TEXT AS ABOVE |
| logical_plan after simplify_expressions | SAME TEXT AS ABOVE |
| logical_plan after eliminate_filter | SAME TEXT AS ABOVE |
| logical_plan after reduce_cross_join | SAME TEXT AS ABOVE |
| logical_plan after common_sub_expression_eliminate | SAME TEXT AS ABOVE |
| logical_plan after eliminate_limit | SAME TEXT AS ABOVE |
| logical_plan after projection_push_down | SAME TEXT AS ABOVE |
| logical_plan after rewrite_disjunctive_predicate | SAME TEXT AS ABOVE |
| logical_plan after reduce_outer_join | SAME TEXT AS ABOVE |
| logical_plan after filter_push_down | SAME TEXT AS ABOVE |
| logical_plan after limit_push_down | SAME TEXT AS ABOVE |
| logical_plan after single_distinct_aggregation_to_group_by | SAME TEXT AS ABOVE |
| logical_plan | Projection: Utf8("3.2") AS foo |
| | EmptyRelation |
| initial_physical_plan | ProjectionExec: expr=[3.2 as foo] |
| | PlaceholderRowExec |
| | |
| physical_plan after aggregate_statistics | SAME TEXT AS ABOVE |
| physical_plan after join_selection | SAME TEXT AS ABOVE |
| physical_plan after coalesce_batches | SAME TEXT AS ABOVE |
| physical_plan after repartition | SAME TEXT AS ABOVE |
| physical_plan after add_merge_exec | SAME TEXT AS ABOVE |
| physical_plan | ProjectionExec: expr=[3.2 as foo] |
| | PlaceholderRowExec |
| | |
+------------------------------------------------------------+---------------------------------------------------------------------------+关于查询优化的思考
DataFusion 中的查询优化采用基于代价的模型。基于代价的模型依赖表级和列级的统计信息来估计选择率;选择率估计是过滤与投影代价分析中的重要一环,因为它们使我们能够估算连接和过滤操作的代价。
构建这类估计的一个关键环节是边界分析,它利用区间算术,将诸如 a > 2500 AND a <= 5000 这样的表达式转化为准确的选择率估计,进而用于寻找更高效的执行计划。
AnalysisContext API
AnalysisContext 在表达式求值和边界分析过程中充当共享的知识库。可以把它看作一个动态仓库,用于维护以下信息:
- 列与表达式当前已知的边界
- 已收集或已推断出的统计信息
- 一个可随分析推进而更新的可变状态
AnalysisContext 尤为强大的地方在于它能够沿表达式树传播信息。表达式树中的每个节点在被分析时,都可以从这个共享上下文中读取信息,也可以向其中写入信息,从而支持复杂的边界分析与推断。
用 ColumnStatistics 进行基数估计
列统计信息是优化决策的基础。DataFusion 的 ColumnStatistics 不仅跟踪简单的度量,还提供了一整套丰富的信息,包括:
- 空值数量
- 最大值与最小值
- 数值之和(针对数值列)
- 不同值的数量
其中每一项统计信息都封装在 Precision 类型中,用于表明该值是精确值还是估计值,从而让优化器能够对基数估计的可靠性做出明智的判断。
边界分析流程
边界分析过程分为若干阶段,每个阶段都建立在前一阶段所收集信息的基础之上。随着分析在表达式树中推进,AnalysisContext 会被持续更新。
表达式边界分析
在分析表达式时,DataFusion 使用区间算术来执行边界分析。考虑这样一个简单的谓词:age > 18 AND age <= 25。其分析流程如下:
上下文初始化
- 从已知的列统计信息开始
- 根据列约束设定初始边界
- 初始化共享的分析上下文
遍历表达式树
- 分析表达式树中的每个节点
- 向上传播边界信息
- 允许子节点影响父节点的边界
边界更新
- 每个表达式都可以更新共享的上下文
- 这些变更会流经整个表达式树
- 最终的边界信息为优化决策提供依据
使用分析 API
下面的示例展示了如何对物理表达式执行分析遍历,以推断该表达式的选取率(selectivity)及其可能取值的范围。
fn analyze_filter_example() -> Result<()> {
// Create a schema with an 'age' column
let age = Field::new("age", DataType::Int64, false);
let schema = Arc::new(Schema::new(vec![age]));
// Define column statistics
let column_stats = ColumnStatistics::default()
.with_min_value(Precision::Exact(ScalarValue::Int64(Some(14))))
.with_max_value(Precision::Exact(ScalarValue::Int64(Some(79))))
.with_null_count(Precision::Exact(0));
// Create expression: age > 18 AND age <= 25
let expr = col("age")
.gt(lit(18i64))
.and(col("age").lt_eq(lit(25i64)));
// Initialize analysis context
let initial_boundaries = vec![ExprBoundaries::try_from_column(
&schema, &column_stats, 0)?];
let context = AnalysisContext::new(initial_boundaries);
// Analyze expression
let df_schema = DFSchema::try_from(schema)?;
let physical_expr = SessionContext::new().create_physical_expr(expr, &df_schema)?;
let analysis = analyze(&physical_expr, context, df_schema.as_ref())?;
Ok(())
}评论
登录后参与评论
KnowForge