库用户指南

使用 DataFrame API

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

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

使用 DataFrame API

用户指南介绍了 DataFrame API,本节将更深入地描述该 API。

什么是 DataFrame?

如用户指南所述,DataFusion 的 DataFrame 借鉴了 Pandas DataFrame 的接口设计,其实现是对 LogicalPlan 的一层薄封装,用于构建和执行这些计划。

最简单的 dataframe 是扫描某个表的 dataframe,而该表可以位于文件中或内存中。

如何生成 DataFrame

与其他 DataFrame API 类似,你可以通过 API 以编程方式构造 DataFrame。例如,你可以将内存中的 RecordBatch 读取为 DataFrame:

use std::sync::Arc;
use datafusion::prelude::*;
use datafusion::arrow::array::{ArrayRef, Int32Array};
use datafusion::arrow::record_batch::RecordBatch;
use datafusion::error::Result;

#[tokio::main]
async fn main() -> Result<()> {
    let ctx = SessionContext::new();
    // Register an in-memory table containing the following data
    // id | bank_account
    // ---|-------------
    // 1  | 9000
    // 2  | 8000
    // 3  | 7000
    let data = RecordBatch::try_from_iter(vec![
        ("id", Arc::new(Int32Array::from(vec![1, 2, 3])) as ArrayRef),
        ("bank_account", Arc::new(Int32Array::from(vec![9000, 8000, 7000]))),
    ])?;
    // Create a DataFrame that scans the user table, and finds
    // all users with a bank account at least 8000
    // and sorts the results by bank account in descending order
    let dataframe = ctx
        .read_batch(data)?
        .filter(col("bank_account").gt_eq(lit(8000)))? // bank_account >= 8000
        .sort(vec![col("bank_account").sort(false, true)])?; // ORDER BY bank_account DESC

    Ok(())
}

你同样可以从 SQL 查询生成 DataFrame,并使用 DataFrame 的 API 来处理该查询的输出结果。

use std::sync::Arc;
use datafusion::prelude::*;
use datafusion::assert_batches_eq;
use datafusion::arrow::array::{ArrayRef, Int32Array};
use datafusion::arrow::record_batch::RecordBatch;
use datafusion::error::Result;

#[tokio::main]
async fn main() -> Result<()> {
    let ctx = SessionContext::new();
    // Register the same in-memory table as the previous example
    let data = RecordBatch::try_from_iter(vec![
        ("id", Arc::new(Int32Array::from(vec![1, 2, 3])) as ArrayRef),
        ("bank_account", Arc::new(Int32Array::from(vec![9000, 8000, 7000]))),
    ])?;
    ctx.register_batch("users", data)?;
    // Create a DataFrame using SQL
    let dataframe = ctx.sql("SELECT * FROM users;")
        .await?
        // Note we can filter the output of the query using the DataFrame API
        .filter(col("bank_account").gt_eq(lit(8000)))?; // bank_account >= 8000

    let results = &dataframe.collect().await?;

    // use the `assert_batches_eq` macro to show the output
    assert_batches_eq!(
        vec![
            "+----+--------------+",
            "| id | bank_account |",
            "+----+--------------+",
            "| 1  | 9000         |",
            "| 2  | 8000         |",
            "+----+--------------+",
        ],
        &results
    );
    Ok(())
}

Collect / Streaming Exec

DataFusion 的 DataFrame 是"惰性"的,也就是说在被执行之前它们不会做任何处理,这使得额外的优化成为可能。

你可以通过以下三种方式之一来运行 DataFrame:

  1. collect:执行查询,并将所有输出缓冲到一个 Vec<RecordBatch> 中
  2. execute_stream:开始执行并返回一个 SendableRecordBatchStream,它在每次调用 next() 时增量地计算输出
  3. cache:执行查询,并把输出缓冲到一个新的内存中 DataFrame

要将所有输出收集到内存缓冲区中,可以使用 collect 方法:

use datafusion::prelude::*;
use datafusion::error::Result;

#[tokio::main]
async fn main() -> Result<()> {
    let ctx = SessionContext::new();
    // read the contents of a CSV file into a DataFrame
    let df = ctx.read_csv("tests/data/example.csv", CsvReadOptions::new()).await?;
    // execute the query and collect the results as a Vec<RecordBatch>
    let batches = df.collect().await?;
    for record_batch in batches {
        println!("{record_batch:?}");
    }
    Ok(())
}

使用 execute_stream 可以一次生成一个 RecordBatch,逐步产出输出:

use datafusion::prelude::*;
use datafusion::error::Result;
use futures::stream::StreamExt;

#[tokio::main]
async fn main() -> Result<()> {
    let ctx = SessionContext::new();
    // read example.csv file into a DataFrame
    let df = ctx.read_csv("tests/data/example.csv", CsvReadOptions::new()).await?;
    // begin execution (returns quickly, does not compute results)
    let mut stream = df.execute_stream().await?;
    // results are returned incrementally as they are computed
    while let Some(record_batch) = stream.next().await {
        println!("{record_batch:?}");
    }
    Ok(())
}

将 DataFrame 写入文件

你也可以将 DataFrame 的内容写入文件。写入文件时,DataFusion 会执行 DataFrame 并将结果流式传输到输出。DataFusion 支持写入 csv、json、arrow、avro 和 parquet 文件,也支持通过 API 写入自定义文件格式(示例见 custom_file_format.rs)。

例如,要读取一个 CSV 文件并将其写入 parquet 文件,可以使用 DataFrame::write_parquet 方法。

use datafusion::prelude::*;
use datafusion::error::Result;
use datafusion::dataframe::DataFrameWriteOptions;

#[tokio::main]
async fn main() -> Result<()> {
    let ctx = SessionContext::new();
    // read example.csv file into a DataFrame
    let df = ctx.read_csv("tests/data/example.csv", CsvReadOptions::new()).await?;
    // stream the contents of the DataFrame to the `example.parquet` file
    let target_path = tempfile::tempdir()?.path().join("example.parquet");
    df.write_parquet(
        target_path.to_str().unwrap(),
        DataFrameWriteOptions::new(),
        None, // writer_options
    ).await;
    Ok(())
}

输出文件的内容如下(输出示例):

> select * from '../datafusion/core/example.parquet';
+---+---+---+
| a | b | c |
+---+---+---+
| 1 | 2 | 3 |
+---+---+---+

LogicalPlan 与 DataFrame 之间的关系

DataFrame 结构体的定义如下:

use datafusion::execution::session_state::SessionState;
use datafusion::logical_expr::LogicalPlan;
pub struct DataFrame {
    // state required to execute a LogicalPlan
    session_state: Box<SessionState>,
    // LogicalPlan that describes the computation to perform
    plan: LogicalPlan,
}

如上所示,DataFrame 是 LogicalPlan 的轻量级封装,因此你可以在两者之间轻松地来回转换。

use datafusion::prelude::*;
use datafusion::error::Result;
use datafusion::logical_expr::LogicalPlanBuilder;

#[tokio::main]
async fn main() -> Result<()>{
    let ctx = SessionContext::new();
    // read example.csv file into a DataFrame
    let df = ctx.read_csv("tests/data/example.csv", CsvReadOptions::new()).await?;
    // You can easily get the LogicalPlan from the DataFrame
    let (_state, plan) = df.into_parts();
    // Just combine LogicalPlan with SessionContext and you get a DataFrame
    // get LogicalPlan in dataframe
    let new_df = DataFrame::new(ctx.state(), plan);
    Ok(())
}

实际上,通过 DataFrame 的方法,你可以构建出与使用 LogicalPlanBuilder 时完全相同的 LogicalPlan:

use datafusion::prelude::*;
use datafusion::error::Result;
use datafusion::logical_expr::LogicalPlanBuilder;

#[tokio::main]
async fn main() -> Result<()>{
    let ctx = SessionContext::new();
    // read example.csv file into a DataFrame
    let df = ctx.read_csv("tests/data/example.csv", CsvReadOptions::new()).await?;
    // Create a new DataFrame sorted by  `id`, `bank_account`
    let new_df = df.select(vec![col("a"), col("b")])?
        .sort_by(vec![col("a")])?;
    // Build the same plan using the LogicalPlanBuilder
    // Similar to `SELECT a, b FROM example.csv ORDER BY a`
    let df = ctx.read_csv("tests/data/example.csv", CsvReadOptions::new()).await?;
    let (_state, plan) = df.into_parts(); // get the DataFrame's LogicalPlan
    let plan = LogicalPlanBuilder::from(plan)
        .project(vec![col("a"), col("b")])?
        .sort_by(vec![col("a")])?
        .build()?;
    // prove they are the same
    assert_eq!(new_df.logical_plan(), &plan);
    Ok(())
}

评论

登录后参与评论

正在加载评论…