用户指南
示例用法
登录后可跨设备保存划线和私人笔记登录
示例用法
在这个示例中,我们将对 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 |
+---+--------+评论
登录后参与评论
正在加载评论…
KnowForge