库用户指南

使用 SQL API

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

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

使用 SQL API

DataFusion 提供了完整的 SQL API,让你可以通过 SQL 查询字符串与 DataFusion 交互。使用 SQL API 最简单的方式是采用 SessionContext 结构体,它提供了执行 SQL 查询的最高层 API。

要使用 SQL,首先需要将数据注册为一张表,然后通过 SessionContext::sql 方法运行查询。若需要更底层的控制(例如禁止 DDL),可以使用 SessionContext::sql_with_options 或 SessionState API。

使用 SessionContext::register* 注册数据源

SessionContext::register* 方法用于告诉 DataFusion 数据源的名称以及如何读取数据。注册完成后,你便可以通过 SessionContext::sql 方法执行 SQL 查询,并将数据源作为表来引用。

SessionContext::sql 方法为了便于使用,会返回一个 DataFrame。有关如何使用 DataFrame 的更多信息,请参阅“使用 DataFrame API”一节。

读取 CSV 文件

use datafusion::error::Result;
use datafusion::prelude::*;
use arrow::record_batch::RecordBatch;

#[tokio::main]
async fn main() -> Result<()> {
    let ctx = SessionContext::new();
    // register the "example" table
    ctx.register_csv("example", "tests/data/example.csv", CsvReadOptions::new()).await?;
    // create a plan to run a SQL query
    let df = ctx.sql("SELECT a, min(b) FROM example WHERE a <= b GROUP BY a LIMIT 100").await?;
    // execute the plan and collect the results as Vec<RecordBatch>
    let results: Vec<RecordBatch> = df.collect().await?;
    // Use the assert_batches_eq macro to compare the results with expected output
    datafusion::assert_batches_eq!(vec![
        "+---+----------------+",
        "| a | min(example.b) |",
        "+---+----------------+",
        "| 1 | 2              |",
        "+---+----------------+",
        ],
        &results
    );
  Ok(())
}

读取 Apache Parquet 文件

与 CSV 类似,你可以使用 register_parquet 方法将 Parquet 文件注册为表。

use datafusion::error::Result;
use datafusion::prelude::*;
#[tokio::main]
async fn main() -> Result<()> {
    // create local session context
    let ctx = SessionContext::new();
    let testdata = datafusion::test_util::parquet_test_data();

    // register parquet file with the execution context
    ctx.register_parquet(
        "alltypes_plain",
        &format!("{testdata}/alltypes_plain.parquet"),
        ParquetReadOptions::default(),
    )
    .await?;

    // execute the query
    let df = ctx.sql(
        "SELECT int_col, double_col, CAST(date_string_col as VARCHAR) \
         FROM alltypes_plain \
         WHERE id > 1 AND tinyint_col < double_col",
    ).await?;

    // execute the plan, and compare to the expected results
    let results = df.collect().await?;
    datafusion::assert_batches_eq!(vec![
        "+---------+------------+--------------------------------+",
        "| int_col | double_col | alltypes_plain.date_string_col |",
        "+---------+------------+--------------------------------+",
        "| 1       | 10.1       | 03/01/09                       |",
        "| 1       | 10.1       | 04/01/09                       |",
        "| 1       | 10.1       | 02/01/09                       |",
        "+---------+------------+--------------------------------+",
        ],
        &results
    );
    Ok(())
}

读取 Apache Avro 文件

DataFusion 还可以通过 register_avro 方法读取 Avro 文件。

{
use datafusion::arrow::util::pretty;
use datafusion::error::Result;
use datafusion::prelude::*;

#[tokio::main]
async fn main() -> Result<()> {
    let ctx = SessionContext::new();
    // find the path to the avro test files
    let testdata = datafusion::test_util::arrow_test_data();
    // register avro file with the execution context
    let avro_file = &format!("{testdata}/avro/alltypes_plain.avro");
    ctx.register_avro("alltypes_plain", avro_file, AvroReadOptions::default()).await?;

    // execute the query
    let df = ctx.sql(
        "SELECT int_col, double_col, CAST(date_string_col as VARCHAR) \
         FROM alltypes_plain \
         WHERE id > 1 AND tinyint_col < double_col"
      ).await?;

    // execute the plan, and compare to the expected results
    let results = df.collect().await?;
    datafusion::assert_batches_eq!(vec![
        "+---------+------------+--------------------------------+",
        "| int_col | double_col | alltypes_plain.date_string_col |",
        "+---------+------------+--------------------------------+",
        "| 1       | 10.1       | 03/01/09                       |",
        "| 1       | 10.1       | 04/01/09                       |",
        "| 1       | 10.1       | 02/01/09                       |",
        "+---------+------------+--------------------------------+",
        ],
        &results
    );
    Ok(())
}
}

将多个文件读取为一张表

也可以将多个文件作为单张表来读取。这可以通过 ListingTableProvider 实现,它接收一组文件路径,并将这些文件作为单张表读取,在必要时对各自的结构进行匹配。

即将推出

通过 SQL 使用 CREATE EXTERNAL TABLE 注册数据源

你也可以使用 CREATE EXTERNAL TABLE 语句通过 SQL 注册文件。

use datafusion::error::Result;
use datafusion::prelude::*;
#[tokio::main]
async fn main() -> Result<()> {
    // create local session context
    let ctx = SessionContext::new();
    let testdata = datafusion::test_util::parquet_test_data();

    // register parquet file using SQL
    let ddl = format!(
        "CREATE EXTERNAL TABLE alltypes_plain \
        STORED AS PARQUET LOCATION '{testdata}/alltypes_plain.parquet'"
    );
    ctx.sql(&ddl).await?;

    // execute the query referring to the alltypes_plain table we just registered
    let df = ctx.sql(
        "SELECT int_col, double_col, CAST(date_string_col as VARCHAR) \
         FROM alltypes_plain \
         WHERE id > 1 AND tinyint_col < double_col",
    ).await?;

    // execute the plan, and compare to the expected results
    let results = df.collect().await?;
    datafusion::assert_batches_eq!(vec![
        "+---------+------------+--------------------------------+",
        "| int_col | double_col | alltypes_plain.date_string_col |",
        "+---------+------------+--------------------------------+",
        "| 1       | 10.1       | 03/01/09                       |",
        "| 1       | 10.1       | 04/01/09                       |",
        "| 1       | 10.1       | 02/01/09                       |",
        "+---------+------------+--------------------------------+",
        ],
        &results
    );
    Ok(())
}

评论

登录后参与评论

正在加载评论…