用户指南

示例用法

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

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

示例用法

在这个示例中,我们将对 example.csv 文件执行一些简单的处理。

项目还附带了更多代码示例。

添加已发布的 DataFusion 依赖

请在 DataFusion 的 crates.io 页面上查找最新可用的 DataFusion 版本。将该依赖添加到你的 Cargo.toml 文件中:

datafusion = "55.1.0"
tokio = { version = "1.0", features = ["rt-multi-thread"] }

对存储在 CSV 中的数据运行 SQL 查询

use datafusion::prelude::*;

#[tokio::main]
async fn main() -> datafusion::error::Result<()> {
  // register the table
  let ctx = SessionContext::new();
  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 and print results
  df.show().await?;
  Ok(())
}

有关 SQL API 的更多信息,请参阅库指南的 SQL API 章节。

使用 DataFrame API 处理存储在 CSV 中的数据

use datafusion::prelude::*;
use datafusion::functions_aggregate::expr_fn::min;

#[tokio::main]
async fn main() -> datafusion::error::Result<()> {
  // create the dataframe
  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))?;

  // execute and print results
  df.show().await?;
  Ok(())
}

两个示例的输出

+---+--------+
| a | MIN(b) |
+---+--------+
| 1 | 2      |
+---+--------+

Arrow 版本

DataFusion 的许多公开 API 使用了 arrow 和 parquet crate 中的类型,因此如果你的项目中使用了 arrow,其版本必须与 DataFusion 所使用的版本一致。你可以在 DataFusion 的 crates.io 页面上查看所需的版本。

确保版本一致的最简单方法是使用 DataFusion 导出的 arrow,例如:

use datafusion::arrow::datatypes::Schema;

例如,DataFusion 26.0.0 依赖要求 arrow 为 40.0.0。如果你在自己的项目中改用 arrow 41.0.0,可能会遇到如下错误:

mismatched types [E0308] expected `Schema`, found `arrow_schema::Schema` Note: `arrow_schema::Schema` and `Schema` have similar names, but are actually distinct types Note: `arrow_schema::Schema` is defined in crate `arrow_schema` Note: `Schema` is defined in crate `arrow_schema` Note: perhaps two different versions of crate `arrow_schema` are being used? Note: associated function defined here

或者对 ArrayRef 调用 downcast_ref 可能会意外返回 None。

标识符与大小写

请注意,在 SQL 中所有标识符实际上都会被转换为小写,因此如果你的 CSV 文件包含大写字母(例如 Name),则必须将列名放在双引号中,否则示例将无法运行。

通过将 datafusion.sql_parser.enable_ident_normalization 配置设置为 false,可以禁用此行为。详情请参阅配置设置。

为了说明此行为,请考虑 capitalized_example.csv 文件:

  • 使用 SQL API:
use datafusion::prelude::*;

#[tokio::main]
async fn main() -> datafusion::error::Result<()> {
  // register the table
  let ctx = SessionContext::new();
  ctx.register_csv("example", "tests/data/capitalized_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\" <= c GROUP BY \"A\" LIMIT 100").await?;

  // execute and print results
  df.show().await?;
  Ok(())
}
  • 使用 DataFrame API:
use datafusion::prelude::*;
use datafusion::functions_aggregate::expr_fn::min;

#[tokio::main]
async fn main() -> datafusion::error::Result<()> {
  // create the dataframe
  let ctx = SessionContext::new();
  let df = ctx.read_csv("tests/data/capitalized_example.csv", CsvReadOptions::new()).await?;

  let df = df
      // col will parse the input string, hence requiring double quotes to maintain the capitalization
      .filter(col("\"A\"").lt_eq(col("c")))?
      // alternatively use ident to pass in an unqualified column name directly without parsing
      .aggregate(vec![ident("A")], vec![min(col("b"))])?
      .limit(0, Some(100))?;

  // execute and print results
  df.show().await?;
  Ok(())
}

两个示例的输出:

+---+--------+
| A | MIN(b) |
+---+--------+
| 2 | 1      |
| 1 | 2      |
+---+--------+

评论

登录后参与评论

正在加载评论…