用户指南
DataFrame API
登录后可跨设备保存划线和私人笔记登录
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(())
}评论
登录后参与评论
正在加载评论…
KnowForge