库用户指南

构建逻辑计划

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

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

构建逻辑计划

逻辑计划是数据库查询的一种结构化表示,它描述了从数据库或数据源中检索数据所需的高层操作与转换。它抽象掉了具体的实现细节,专注于查询的逻辑流程,包括过滤、排序以及连接表等操作。

该逻辑计划是生成经过优化的物理执行计划之前的中间步骤。这一点在架构指南的查询规划与执行概述部分有更详细的说明。

手动构建逻辑计划

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 文档。

  • filter
  • limit
  • sort
  • distinct
  • join

下面的示例演示了如何构建与前一个示例相同的简单查询计划,即先进行表扫描,再进行过滤。

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。

评论

登录后参与评论

正在加载评论…