用户指南

DataFrame API

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

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

DataFrame API

DataFrame 概述

DataFrame 表示一组具有相同命名列的逻辑行集合,类似于 Pandas DataFrame 或 Spark DataFrame。

DataFrame 通常通过调用 SessionContext 上的方法来创建,例如 read_csv,随后可以通过调用转换方法,例如 filter、select、aggregate 和 limit,逐步构建查询定义。

查询可以通过调用 collect 方法来执行。

DataFusion 的 DataFrame 采用惰性求值,这意味着每次转换都会创建一个新计划,但不会立即执行任何实际操作。这种方式使得整体计划能够在执行之前得到优化。计划会在调用动作方法时(例如 collect)被求值(执行)。更多详情请参阅库用户指南。

DataFrame API 在 docs.rs 上的 API 参考中有完善的文档。有关如何构建与 DataFrame API 搭配使用的逻辑表达式(Expr)的更多信息,请参阅表达式参考。

示例

DataFrame 结构体是 DataFusion prelude 的一部分,可以通过以下语句导入。

use datafusion::prelude::*;

下面是一个使用 DataFrame API 执行查询的最小示例。

使用宏 API 从内存中的行创建 DataFrame

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

#[tokio::main]
async fn main() -> Result<()> {
    // Create a new dataframe with in-memory data using macro
    let df = dataframe!(
        "a" => [1, 2, 3],
        "b" => [true, true, false],
        "c" => [Some("foo"), Some("bar"), None]
    )?;
    df.show().await?;
    Ok(())
}

使用标准 API 从文件或内存行创建 DataFrame

use datafusion::arrow::array::{Int32Array, RecordBatch, StringArray};
use datafusion::arrow::datatypes::{DataType, Field, Schema};
use datafusion::error::Result;
use datafusion::functions_aggregate::expr_fn::min;
use datafusion::prelude::*;
use std::sync::Arc;

#[tokio::main]
async fn main() -> Result<()> {
    // Read the data from a csv file
    let ctx = SessionContext::new();
    let df = ctx.read_csv("tests/data/example.csv", CsvReadOptions::new()).await?;
    let df = df.filter(col("a").lt_eq(col("b")))?
        .aggregate(vec![col("a")], vec![min(col("b"))])?
        .limit(0, Some(100))?;
    // Print results
    df.show().await?;

    // Create a new dataframe with in-memory data
    let schema = Schema::new(vec![
      Field::new("id", DataType::Int32, true),
      Field::new("name", DataType::Utf8, true),
    ]);
    let batch = RecordBatch::try_new(
      Arc::new(schema),
      vec![
          Arc::new(Int32Array::from(vec![1, 2, 3])),
          Arc::new(StringArray::from(vec!["foo", "bar", "baz"])),
      ],
    )?;
    let df = ctx.read_batch(batch)?;
    df.show().await?;

    Ok(())
}

评论

登录后参与评论

正在加载评论…