快速入门

Python/Rust 快速入门

师成师成· 更新于 2026-09-29· 阅读 13 分钟· 0 次阅读

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

本指南将帮助你快速上手 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?;
}

参与贡献

关于为该项目做出贡献的所有详细信息,请查看贡献指南。

评论

登录后参与评论

正在加载评论…