库用户指南
使用 SQL API
登录后可跨设备保存划线和私人笔记登录
使用 SQL API
DataFusion 提供了完整的 SQL API,让你可以通过 SQL 查询字符串与 DataFusion 交互。使用 SQL API 最简单的方式是采用 SessionContext 结构体,它提供了执行 SQL 查询的最高层 API。
要使用 SQL,首先需要将数据注册为一张表,然后通过 SessionContext::sql 方法运行查询。若需要更底层的控制(例如禁止 DDL),可以使用 SessionContext::sql_with_options 或 SessionState API。
使用 SessionContext::register* 注册数据源
SessionContext::register* 方法用于告诉 DataFusion 数据源的名称以及如何读取数据。注册完成后,你便可以通过 SessionContext::sql 方法执行 SQL 查询,并将数据源作为表来引用。
SessionContext::sql 方法为了便于使用,会返回一个 DataFrame。有关如何使用 DataFrame 的更多信息,请参阅“使用 DataFrame API”一节。
读取 CSV 文件
use datafusion::error::Result;
use datafusion::prelude::*;
use arrow::record_batch::RecordBatch;
#[tokio::main]
async fn main() -> Result<()> {
let ctx = SessionContext::new();
// register the "example" table
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 the plan and collect the results as Vec<RecordBatch>
let results: Vec<RecordBatch> = df.collect().await?;
// Use the assert_batches_eq macro to compare the results with expected output
datafusion::assert_batches_eq!(vec![
"+---+----------------+",
"| a | min(example.b) |",
"+---+----------------+",
"| 1 | 2 |",
"+---+----------------+",
],
&results
);
Ok(())
}读取 Apache Parquet 文件
与 CSV 类似,你可以使用 register_parquet 方法将 Parquet 文件注册为表。
use datafusion::error::Result;
use datafusion::prelude::*;
#[tokio::main]
async fn main() -> Result<()> {
// create local session context
let ctx = SessionContext::new();
let testdata = datafusion::test_util::parquet_test_data();
// register parquet file with the execution context
ctx.register_parquet(
"alltypes_plain",
&format!("{testdata}/alltypes_plain.parquet"),
ParquetReadOptions::default(),
)
.await?;
// execute the query
let df = ctx.sql(
"SELECT int_col, double_col, CAST(date_string_col as VARCHAR) \
FROM alltypes_plain \
WHERE id > 1 AND tinyint_col < double_col",
).await?;
// execute the plan, and compare to the expected results
let results = df.collect().await?;
datafusion::assert_batches_eq!(vec![
"+---------+------------+--------------------------------+",
"| int_col | double_col | alltypes_plain.date_string_col |",
"+---------+------------+--------------------------------+",
"| 1 | 10.1 | 03/01/09 |",
"| 1 | 10.1 | 04/01/09 |",
"| 1 | 10.1 | 02/01/09 |",
"+---------+------------+--------------------------------+",
],
&results
);
Ok(())
}读取 Apache Avro 文件
DataFusion 还可以通过 register_avro 方法读取 Avro 文件。
{
use datafusion::arrow::util::pretty;
use datafusion::error::Result;
use datafusion::prelude::*;
#[tokio::main]
async fn main() -> Result<()> {
let ctx = SessionContext::new();
// find the path to the avro test files
let testdata = datafusion::test_util::arrow_test_data();
// register avro file with the execution context
let avro_file = &format!("{testdata}/avro/alltypes_plain.avro");
ctx.register_avro("alltypes_plain", avro_file, AvroReadOptions::default()).await?;
// execute the query
let df = ctx.sql(
"SELECT int_col, double_col, CAST(date_string_col as VARCHAR) \
FROM alltypes_plain \
WHERE id > 1 AND tinyint_col < double_col"
).await?;
// execute the plan, and compare to the expected results
let results = df.collect().await?;
datafusion::assert_batches_eq!(vec![
"+---------+------------+--------------------------------+",
"| int_col | double_col | alltypes_plain.date_string_col |",
"+---------+------------+--------------------------------+",
"| 1 | 10.1 | 03/01/09 |",
"| 1 | 10.1 | 04/01/09 |",
"| 1 | 10.1 | 02/01/09 |",
"+---------+------------+--------------------------------+",
],
&results
);
Ok(())
}
}将多个文件读取为一张表
也可以将多个文件作为单张表来读取。这可以通过 ListingTableProvider 实现,它接收一组文件路径,并将这些文件作为单张表读取,在必要时对各自的结构进行匹配。
即将推出
通过 SQL 使用 CREATE EXTERNAL TABLE 注册数据源
你也可以使用 CREATE EXTERNAL TABLE 语句通过 SQL 注册文件。
use datafusion::error::Result;
use datafusion::prelude::*;
#[tokio::main]
async fn main() -> Result<()> {
// create local session context
let ctx = SessionContext::new();
let testdata = datafusion::test_util::parquet_test_data();
// register parquet file using SQL
let ddl = format!(
"CREATE EXTERNAL TABLE alltypes_plain \
STORED AS PARQUET LOCATION '{testdata}/alltypes_plain.parquet'"
);
ctx.sql(&ddl).await?;
// execute the query referring to the alltypes_plain table we just registered
let df = ctx.sql(
"SELECT int_col, double_col, CAST(date_string_col as VARCHAR) \
FROM alltypes_plain \
WHERE id > 1 AND tinyint_col < double_col",
).await?;
// execute the plan, and compare to the expected results
let results = df.collect().await?;
datafusion::assert_batches_eq!(vec![
"+---------+------------+--------------------------------+",
"| int_col | double_col | alltypes_plain.date_string_col |",
"+---------+------------+--------------------------------+",
"| 1 | 10.1 | 03/01/09 |",
"| 1 | 10.1 | 04/01/09 |",
"| 1 | 10.1 | 02/01/09 |",
"+---------+------------+--------------------------------+",
],
&results
);
Ok(())
}评论
登录后参与评论
正在加载评论…
KnowForge