Python/Rust 快速入门
本指南将帮助你快速上手 Hudi-rs——Apache Hudi 的原生 Rust 实现,并提供 Python 绑定。你将学习如何安装、设置,以及如何分别使用 Python 和 Rust 接口执行基本操作。
安装
# Python
pip install hudi
# Rust
cargo add hudi使用示例
[!NOTE] 这些示例假设在
/tmp/trips_table处已存在一个 Hudi 表,该表可按照快速入门指南创建。
快照查询
快照查询从表中读取数据的最新版本。表 API 同样支持分区过滤条件。
Python
from hudi import HudiTableBuilder
import pyarrow as pa
hudi_table = HudiTableBuilder.from_base_uri("/tmp/trips_table").build()
batches = hudi_table.read_snapshot(filters=[("city", "=", "san_francisco")])
# convert to PyArrow table
arrow_table = pa.Table.from_batches(batches)
result = arrow_table.select(["rider", "city", "ts", "fare"])
print(result)Rust
use hudi::error::Result;
use hudi::table::builder::TableBuilder as HudiTableBuilder;
use arrow::compute::concat_batches;
#[tokio::main]
async fn main() -> Result<()> {
let hudi_table = HudiTableBuilder::from_base_uri("/tmp/trips_table").build().await?;
let batches = hudi_table.read_snapshot(&[("city", "=", "san_francisco")]).await?;
let batch = concat_batches(&batches[0].schema(), &batches)?;
let columns = vec!["rider", "city", "ts", "fare"];
for col_name in columns {
let idx = batch.schema().index_of(col_name).unwrap();
println!("{}: {}", col_name, batch.column(idx));
}
Ok(())
}要在读时合并(MOR)表上运行读优化(RO)查询,请在创建表时设置 hoodie.read.use.read_optimized.mode。
Python
hudi_table = (
HudiTableBuilder
.from_base_uri("/tmp/trips_table")
.with_option("hoodie.read.use.read_optimized.mode", "true")
.build()
)Rust
let hudi_table =
HudiTableBuilder::from_base_uri("/tmp/trips_table")
.with_option("hoodie.read.use.read_optimized.mode", "true")
.build().await?;[!NOTE] 目前读取 MOR 表仅限于数据块为 Parquet 格式的表。
时间旅行查询
时间旅行查询用于读取表中某个特定时间戳处的数据。表 API 同样接受分区过滤条件。
Python
batches = (
hudi_table
.read_snapshot_as_of("20241231123456789", filters=[("city", "=", "san_francisco")])
)Rust
let batches =
hudi_table
.read_snapshot_as_of("20241231123456789", &[("city", "=", "san_francisco")]).await?;增量查询
增量查询用于读取表中给定时间范围内的变更数据。
Python
# read the records between t1 (exclusive) and t2 (inclusive)
batches = hudi_table.read_incremental_records(t1, t2)
# read the records after t1
batches = hudi_table.read_incremental_records(t1)Rust
// read the records between t1 (exclusive) and t2 (inclusive)
let batches = hudi_table.read_incremental_records(t1, Some(t2)).await?;
// read the records after t1
let batches = hudi_table.read_incremental_records(t1, None).await?;[!NOTE] 目前时间戳参数仅支持 Hudi Timeline 格式:
yyyyMMddHHmmssSSS或yyyyMMddHHmmss。
查询引擎集成
Hudi-rs 提供了用于与查询引擎集成的 API。下面几节介绍一些常用的 API。
Table API
使用构造函数或 TableBuilder API 创建 Hudi 表实例。
| 阶段 | API | 描述 |
|---|---|---|
| 查询规划 | get_file_slices() | 对于快照查询,获取文件切片列表。 |
get_file_slices_splits() | 对于快照查询,获取按 split 划分的文件切片列表。 | |
get_file_slices_as_of() | 对于时间旅行查询,获取指定时间点的文件切片列表。 | |
get_file_slices_splits_as_of() | 对于时间旅行查询,获取指定时间点按 split 划分的文件切片列表。 | |
get_file_slices_between() | 对于增量查询,获取时间范围内发生变更的文件切片列表。 | |
| 查询执行 | create_file_group_reader_with_options() | 使用表实例的配置创建文件组读取器实例。 |
File Group API
使用构造函数或 Hudi 表 API create_file_group_reader_with_options() 创建 Hudi 文件组读取器实例。
| 阶段 | API | 描述 |
|---|---|---|
| 查询执行 | read_file_slice() | 从给定的文件切片中读取记录;根据配置,仅从基础文件读取记录,或从基础文件和日志文件读取记录,并基于配置的策略合并记录。 |
Apache DataFusion
启用 hudi crate 的 datafusion 特性后,将提供一个 DataFusion 扩展,用于查询 Hudi 表。
在你的应用中添加带 datafusion 特性的 hudi crate,即可查询 Hudi 表。
cargo new my_project --bin && cd my_project
cargo add tokio@1 datafusion@43
cargo add hudi --features datafusion使用下面的代码片段更新 src/main.rs,然后执行 cargo run。
use std::sync::Arc;
use datafusion::error::Result;
use datafusion::prelude::{DataFrame, SessionContext};
use hudi::HudiDataSource;
#[tokio::main]
async fn main() -> Result<()> {
let ctx = SessionContext::new();
let hudi = HudiDataSource::new_with_options(
"/tmp/trips_table",
[("hoodie.read.input.partitions", "5")]).await?;
ctx.register_table("trips_table", Arc::new(hudi))?;
let df: DataFrame = ctx.sql("SELECT * from trips_table where city = 'san_francisco'").await?;
df.show().await?;
Ok(())
}其他集成
Hudi 还与以下组件集成了:
使用云存储
请确保云存储凭证已作为环境变量正确设置,例如 AWS_*、AZURE_* 或 GOOGLE_*。随后,相关的存储环境变量将被自动读取,目标表的 base URI 中带有 s3://、az:// 或 gs:// 癉式的路径也会得到相应的处理。
也可以通过表 API 以 options 的形式传入存储配置。
Python
from hudi import HudiTableBuilder
hudi_table = (
HudiTableBuilder
.from_base_uri("s3://bucket/trips_table")
.with_option("aws_region", "us-west-2")
.build()
)Rust
use hudi::table::builder::TableBuilder as HudiTableBuilder;
async fn main() -> Result<()> {
let hudi_table =
HudiTableBuilder::from_base_uri("s3://bucket/trips_table")
.with_option("aws_region", "us-west-2")
.build().await?;
}参与贡献
关于为该项目做出贡献的所有详细信息,请查看贡献指南。
评论
登录后参与评论
KnowForge