构建逻辑计划
构建逻辑计划
逻辑计划是数据库查询的一种结构化表示,它描述了从数据库或数据源中检索数据所需的高层操作与转换。它抽象掉了具体的实现细节,专注于查询的逻辑流程,包括过滤、排序以及连接表等操作。
该逻辑计划是生成经过优化的物理执行计划之前的中间步骤。这一点在架构指南的查询规划与执行概述部分有更详细的说明。
手动构建逻辑计划
DataFusion 的 LogicalPlan 是一个枚举,其中的变体代表了所有受支持的运算符,同时还包含一个 Extension 变体,允许基于 DataFusion 构建的项目添加自定义逻辑运算符。
可以像下面所示那样,通过直接创建 LogicalPlan 枚举的实例来创建逻辑计划,但使用 LogicalPlanBuilder 要方便得多,下一节将对它进行介绍。
下面是一个直接构建逻辑计划的示例:
use datafusion::common::DataFusionError;
use datafusion::arrow::datatypes::{DataType, Field, Schema, SchemaRef};
use datafusion::logical_expr::{Filter, LogicalPlan, TableScan, LogicalTableSource};
use datafusion::prelude::*;
use std::sync::Arc;
fn main() -> Result<(), DataFusionError> {
// create a logical table source
let schema = Schema::new(vec![
Field::new("id", DataType::Int32, true),
Field::new("name", DataType::Utf8, true),
]);
let table_source = LogicalTableSource::new(SchemaRef::new(schema));
// create a TableScan plan
let projection = None; // optional projection
let filters = vec![]; // optional filters to push down
let fetch = None; // optional LIMIT
let table_scan = LogicalPlan::TableScan(TableScan::try_new(
"person",
Arc::new(table_source),
projection,
filters,
fetch,
)?
);
// create a Filter plan that evaluates `id > 500` that wraps the TableScan
let filter_expr = col("id").gt(lit(500));
let plan = LogicalPlan::Filter(Filter::try_new(filter_expr, Arc::new(table_scan)) ? );
// print the plan
println!("{}", plan.display_indent_schema());
Ok(())
}本示例生成以下计划:
Filter: person.id > Int32(500) [id:Int32;N, name:Utf8;N]
TableScan: person [id:Int32;N, name:Utf8;N]使用 LogicalPlanBuilder 构建逻辑计划
DataFusion 的逻辑计划可以通过 LogicalPlanBuilder 结构体来创建。此外还有一个 DataFrame API,它是一个更高级的 API,其底层委托给 LogicalPlanBuilder。
有多个函数可用于创建新的 builder,例如:
empty—— 创建一个不含任何字段的空计划values—— 根据一组字面量值创建计划scan—— 创建一个表示表扫描的计划scan_with_filters—— 创建一个带过滤条件的表扫描计划
创建 builder 之后,可以调用转换方法来声明需要对该计划执行的进一步操作。请注意,在这个阶段我们所做的只是构建逻辑计划的结构,不会执行任何查询。
下面列出了一些转换方法的示例,完整的列表请参阅 LogicalPlanBuilder 的 API 文档。
filterlimitsortdistinctjoin
下面的示例演示了如何构建与前一个示例相同的简单查询计划,即先进行表扫描,再进行过滤。
use datafusion::common::DataFusionError;
use datafusion::arrow::datatypes::{DataType, Field, Schema, SchemaRef};
use datafusion::logical_expr::{LogicalPlanBuilder, LogicalTableSource};
use datafusion::prelude::*;
use std::sync::Arc;
fn main() -> Result<(), DataFusionError> {
// create a logical table source
let schema = Schema::new(vec![
Field::new("id", DataType::Int32, true),
Field::new("name", DataType::Utf8, true),
]);
let table_source = LogicalTableSource::new(SchemaRef::new(schema));
// optional projection
let projection = None;
// create a LogicalPlanBuilder for a table scan
let builder = LogicalPlanBuilder::scan("person", Arc::new(table_source), projection)?;
// perform a filter operation and build the plan
let plan = builder
.filter(col("id").gt(lit(500)))? // WHERE id > 500
.build()?;
// print the plan
println!("{}", plan.display_indent_schema());
Ok(())
}本示例生成以下计划:
Filter: person.id > Int32(500) [id:Int32;N, name:Utf8;N]
TableScan: person [id:Int32;N, name:Utf8;N]将逻辑计划转换为物理计划
逻辑计划无法直接执行。它们必须被"编译"成 ExecutionPlan,这通常被称为"物理计划"。
与 LogicalPlan 相比,ExecutionPlan 包含更多细节,例如具体的算法和详细的优化。给定一个 LogicalPlan,创建 ExecutionPlan 最简单的方式是使用 SessionState::create_physical_plan,如下所示
use datafusion::datasource::{provider_as_source, MemTable};
use datafusion::common::DataFusionError;
use datafusion::physical_plan::display::DisplayableExecutionPlan;
use datafusion::arrow::datatypes::{DataType, Field, Schema, SchemaRef};
use datafusion::logical_expr::{LogicalPlanBuilder, LogicalTableSource};
use datafusion::prelude::*;
use std::sync::Arc;
// Creating physical plans may access remote catalogs and data sources
// thus it must be run with an async runtime.
#[tokio::main]
async fn main() -> Result<(), DataFusionError> {
// create a default table source
let schema = Schema::new(vec![
Field::new("id", DataType::Int32, true),
Field::new("name", DataType::Utf8, true),
]);
// To create an ExecutionPlan we must provide an actual
// TableProvider. For this example, we don't provide any data
// but in production code, this would have `RecordBatch`es with
// in memory data
let table_provider = Arc::new(MemTable::try_new(Arc::new(schema), vec![vec![]])?);
// Use the provider_as_source function to convert the TableProvider to a table source
let table_source = provider_as_source(table_provider);
// create a LogicalPlanBuilder for a table scan without projection or filters
let logical_plan = LogicalPlanBuilder::scan("person", table_source, None)?.build()?;
// Now create the physical plan by calling `create_physical_plan`
let ctx = SessionContext::new();
let physical_plan = ctx.state().create_physical_plan(&logical_plan).await?;
// print the plan
println!("{}", DisplayableExecutionPlan::new(physical_plan.as_ref()).indent(true));
Ok(())
}本示例会生成以下物理计划:
DataSourceExec: partitions=0, partition_sizes=[]表来源
前面的示例使用了 LogicalTableSource,它被 DataFusion 用于测试和文档;如果你使用 DataFusion 构建逻辑计划但不使用 DataFusion 的物理规划器,它也是合适的选择。
不过,更常见的情况是使用 TableProvider。要从 TableProvider 获取 TableSource,请使用 provider_as_source 或 DefaultTableSource。
评论
登录后参与评论
KnowForge