Skip to main content

Hudi logo

The native Rust implementation for Apache Hudi, with C++ & Python API bindings.

hudi-rs ci hudi-rs codecov join hudi slack follow hudi x/twitter follow hudi linkedin

The Hudi-rs project aims to standardize the core Apache Hudi APIs, and broaden the Hudi integration in the data ecosystems for a diverse range of users and projects.

Source Downloads Installation Command
PyPi.org pip install hudi
Crates.io cargo add hudi

The hudi crate carries two features: datafusion (off by default, see Apache DataFusion) and spill-rocksdb (on by default), the merge map's on-disk tier, which a merge-on-read merge spills to when a file group's log records exceed hoodie.memory.merge.max.size. RocksDB is built from source with bindgen, so it needs libclang and a C++ toolchain; default-features = false drops it, and a merge that would have spilled then fails instead.

Usage Examples

[!NOTE] These examples expect a Hudi table exists at /tmp/trips_table, created using the quick start guide.

For the full reader API reference (ReadOptions, filter expressions, behavioral guarantees), see docs/reader-spec.md.

Snapshot Query

Snapshot query reads the latest version of the data from the table. The table API also accepts column filters that drive partition + file pruning and row-level filtering.

Python

from hudi import HudiReadOptions, HudiTableBuilder
import pyarrow as pa

hudi_table = HudiTableBuilder.from_base_uri("/tmp/trips_table").build()
batches = hudi_table.read(
    HudiReadOptions(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::ReadOptions;
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 options = ReadOptions::new().with_filters([("city", "=", "san_francisco")])?;
    let batches = hudi_table.read(&options).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(())
}

To run read-optimized (RO) query on Merge-on-Read (MOR) tables, set hoodie.read.use.read_optimized.mode in ReadOptions.

Python

from hudi import HudiReadOptions

batches = hudi_table.read(
    HudiReadOptions(hudi_options={"hoodie.read.use.read_optimized.mode": "true"})
)

Rust

let options = ReadOptions::new()
    .with_hudi_option("hoodie.read.use.read_optimized.mode", "true");
let batches = hudi_table.read(&options).await?;

Time-Travel Query

Time-travel query reads the data at a specific timestamp from the table. The table API also accepts column filters that drive partition + file pruning and row-level filtering.

Python

batches = hudi_table.read(
    HudiReadOptions(filters=[("city", "=", "san_francisco")])
    .with_as_of_timestamp("20241231123456789")
)

Rust

let options = ReadOptions::new()
    .with_as_of_timestamp("20241231123456789")
    .with_filters([("city", "=", "san_francisco")])?;
let batches = hudi_table.read(&options).await?;
Supported timestamp formats

The supported formats for the timestamp argument are:

  • Hudi Timeline format (highest matching precedence): yyyyMMddHHmmssSSS or yyyyMMddHHmmss.
  • Unix epoch time in seconds, milliseconds, microseconds, or nanoseconds.
  • RFC 3339 / ISO 8601 with timezone offset, including:
    • yyyy-MM-dd'T'HH:mm:ss.SSS+00:00
    • yyyy-MM-dd'T'HH:mm:ss.SSSZ
    • yyyy-MM-dd'T'HH:mm:ss+00:00
    • yyyy-MM-dd'T'HH:mm:ssZ

Timestamp strings without a timezone offset (for example yyyy-MM-dd'T'HH:mm:ss) and date-only strings (for example yyyy-MM-dd) are not accepted.

Incremental Query

Incremental query reads the changed data from the table for a given time range.

Python

from hudi import HudiQueryType

# read the records between t1 (exclusive) and t2 (inclusive)
batches = hudi_table.read(
    HudiReadOptions()
    .with_query_type(HudiQueryType.Incremental)
    .with_start_timestamp(t1)
    .with_end_timestamp(t2)
)

# read the records after t1 (end defaults to the latest commit)
batches = hudi_table.read(
    HudiReadOptions()
    .with_query_type(HudiQueryType.Incremental)
    .with_start_timestamp(t1)
)

# with column filters applied to the changed records
batches = hudi_table.read(
    HudiReadOptions(filters=[("city", "=", "san_francisco")])
    .with_query_type(HudiQueryType.Incremental)
    .with_start_timestamp(t1)
    .with_end_timestamp(t2)
)

Rust

use hudi::table::QueryType;

// read the records between t1 (exclusive) and t2 (inclusive)
let options = ReadOptions::new()
    .with_query_type(QueryType::Incremental)
    .with_start_timestamp(t1)
    .with_end_timestamp(t2);
let batches = hudi_table.read(&options).await?;

// read the records after t1 (end defaults to the latest commit)
let options = ReadOptions::new()
    .with_query_type(QueryType::Incremental)
    .with_start_timestamp(t1);
let batches = hudi_table.read(&options).await?;

// with column filters applied to the changed records
let options = ReadOptions::new()
    .with_query_type(QueryType::Incremental)
    .with_start_timestamp(t1)
    .with_end_timestamp(t2)
    .with_filters([("city", "=", "san_francisco")])?;
let batches = hudi_table.read(&options).await?;

Incremental queries support the same timestamp formats as time-travel queries.

Streaming Read

Streaming reads yield RecordBatches one at a time without loading the full result into memory. The same ReadOptions knobs apply, plus batch_size and projection.

Python

options = (
    HudiReadOptions(
        filters=[("city", "=", "san_francisco")],
        projection=["rider", "city", "ts", "fare"],
    )
    .with_batch_size(4096)
)
for batch in hudi_table.read_stream(options):
    print(batch.num_rows)

Rust

use futures::StreamExt;

let options = ReadOptions::new()
    .with_filters([("city", "=", "san_francisco")])?
    .with_projection(["rider", "city", "ts", "fare"])
    .with_batch_size(4096)?;
let mut stream = hudi_table.read_stream(&options).await?;
while let Some(batch) = stream.next().await {
    let batch = batch?;
    println!("{}", batch.num_rows());
}

File Group Reading (Experimental)

File group reading allows you to read data from a specific file slice. This is useful when integrating with query engines, where the plan provides file paths.

Python

from hudi import HudiFileGroupReader

reader = HudiFileGroupReader(
    "/table/base/path", {"hoodie.read.start.timestamp": "0"})

# Returns a PyArrow RecordBatch
record_batch = reader.read_file_slice_from_paths("relative/path.parquet", [])

Rust

use hudi::file_group::reader::FileGroupReader;
use hudi::table::ReadOptions;

// Inside an async context
let reader = FileGroupReader::new_with_options(
    "/table/base/path", [("hoodie.read.start.timestamp", "0")]).await?;

// Returns an Arrow RecordBatch
let record_batch = reader
    .read_file_slice_from_paths(
        "relative/path.parquet",
        Vec::<&str>::new(),
        &ReadOptions::new(),
    )
    .await?;

C++

#include "cxx.h"
#include "src/lib.rs.h"
#include "arrow/c/abi.h"

// Functions may throw rust::Error on failure
auto reader = new_file_group_reader_with_options(
    "/table/base/path", {"hoodie.read.start.timestamp=0"});

// Returns an ArrowArrayStream pointer
std::vector<std::string> log_file_paths{};
ArrowArrayStream* stream_ptr = reader->read_file_slice_from_paths("relative/path.parquet", log_file_paths);

Query Engine Integration

Hudi-rs provides APIs to support integration with query engines. The sections below highlight some commonly used APIs.

Table API

Create a Hudi table instance using its constructor or the TableBuilder API.

All read APIs accept a ReadOptions (Rust) / HudiReadOptions (Python) value. It stores three fields — filters, projection, and hudi_options — and exposes chainable with_* builders for the rest. The available knobs:

  • query_type (with_query_type) — Snapshot (default) or Incremental. Drives dispatch in read, read_stream, and get_file_slices.
  • filters — column filters as (field, op, value) tuples. The field can be any column (partition or data). Used for partition pruning, file-level stats pruning (snapshot only), and row-level filtering.
  • projection — columns to return. Streaming pushes the projection down to the parquet reader; eager reads project after merging.
  • batch_size (with_batch_size) — rows per batch (streaming only; eager reads return one batch per file slice).
  • as_of_timestamp (with_as_of_timestamp) — snapshot/time-travel timestamp (defaults to latest commit).
  • start_timestamp / end_timestamp (with_start_timestamp / with_end_timestamp) — incremental range (defaults to earliest…latest).
  • hudi_options — Hudi configs for this read (e.g. hoodie.read.use.read_optimized.mode). A config that selects which read to perform — hoodie.read.query.type and the as-of/start/end timestamps — is per-read only and is dropped when set on the table. The rest describe how to read: set them on the table and override them here. See Read configs.
Stage API Description
Query planning get_file_slices(options) Get the file slices the read targets, dispatched on options.query_type. To bucket for parallel reads, call hudi::util::collection::split_into_chunks on the result.
compute_table_stats(options) Estimated (num_rows, byte_size) for scan planning, derived from the metadata table. Snapshot only. Returns None for incremental queries, and whenever the estimate cannot be computed (no metadata table, non-Parquet base files, footer sampling failure).
Query execution create_file_group_reader_with_options(read_options, extra_storage_overrides) Create a file group reader with the table's configs. In Python both args are optional; in Rust read_options is an Option and extra_storage_overrides takes a (possibly empty) iterator. Timestamps are resolved automatically (e.g. AsOfTimestamp → EndTimestamp), so callers can pass the same options used for get_file_slices.
read(options) / read_stream(options) Record-read APIs. Dispatch on options.query_type. read_stream errors on Incremental for now. Per-slice streaming lives on FileGroupReader.

Read configs

Read configs reach a read either through ReadOptions / HudiReadOptions or, for those scoped table or read, through the table (TableBuilder, hoodie.properties, hudi-defaults.conf), where a per-read value wins. A per read config set on the table is dropped: baked in there, it would silently redirect every later read.

Config Default Scope Notes
hoodie.read.query.type snapshot per read snapshot or incremental.
hoodie.read.as.of.timestamp latest commit per read Snapshot time-travel point.
hoodie.read.start.timestamp / hoodie.read.end.timestamp earliest / latest per read Incremental window, half-open (start, end].
hoodie.read.file.group.reader.version 2 table or read Which file group reader merges a slice. Version 2 is the default; a read it cannot serve falls back to version 1.
hoodie.read.use.read_optimized.mode false table or read Read base files only, skipping the log files, on MOR tables.
hoodie.read.stream.batch_size 1024 table or read Rows per batch for streaming reads.
hoodie.read.file.slice.read.concurrency 4 table or read File slices read concurrently.
hoodie.read.scan.max.memory.size unset table or read Total bytes a whole scan may use for concurrent slice reads; when set, the concurrency is derived from it. Reaching the limit lowers throughput, it never fails the read.
hoodie.read.input.partitions 0 table or read How many partitions the DataFusion table provider buckets the file slices into. 0 defers to DataFusion's target_partitions.
hoodie.merge.use.record.positions false table or read Match a log record to the base row it updates by position rather than by record key. Honored by reader version 2 only; a log block written without positions is merged by key regardless.

A table whose hoodie.record.merge.mode is CUSTOM needs a merger for its payload class. When neither reader has one, the read fails rather than returning wrong rows; set hoodie.read.file.group.reader.version=1 to read it the way it was read before, unless its base files are HFile, which version 1 cannot read.

Base files are read from Parquet, Lance (hoodie.table.base.file.format=lance), and HFile (metadata tables only).

File Group API

Create a Hudi file group reader instance using its constructor or the Hudi table API create_file_group_reader_with_options().

Stage API Description
Query execution read_file_slice() Read records from a given file slice; based on the configs, read records from only base file, or from base file and log files, and merge records based on the configured strategy.
read_file_slice_from_paths() Read records from an explicit base file path and a list of log file paths. Pass an empty log path list to read just the base file.
read_file_slice_stream() Streaming version of read_file_slice(). A base-file-only or read-optimized slice streams straight from the base file. A MOR slice with log files streams merged chunks under file group reader version 2 (the default), and collapses to a single merged batch under version 1, whose merge has no incremental form.
read_file_slice_from_paths_stream() Streaming version of read_file_slice_from_paths().

Apache DataFusion

Enabling the hudi crate with datafusion feature will provide a DataFusion extension to query Hudi tables.

Add crate hudi with datafusion feature to your application to query a Hudi table.
cargo new my_project --bin && cd my_project
cargo add tokio@1 datafusion@54
cargo add hudi --features datafusion

Update src/main.rs with the code snippet below then cargo run.

Add python hudi with datafusion feature to your application to query a Hudi table.
pip install hudi[datafusion]

Rust

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(())
}

Python

from datafusion import SessionContext
from hudi import HudiDataFusionDataSource

table = HudiDataFusionDataSource(
    "/tmp/trips_table", [("hoodie.read.input.partitions", "5")]
)
ctx = SessionContext()
ctx.register_table("trips", table)
ctx.sql("SELECT max(fare), city from trips group by city order by 1 desc").show()

Other Integrations

Hudi is also integrated with

Work with cloud storage

Ensure cloud storage credentials are set properly as environment variables, e.g., AWS_*, AZURE_*, or GOOGLE_*. Relevant storage environment variables will then be picked up. The target table's base uri with schemes such as s3://, az://, or gs:// will be processed accordingly.

Alternatively, you can pass the storage configuration as options via Table APIs.

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::error::Result;
use hudi::table::builder::TableBuilder as HudiTableBuilder;

#[tokio::main]
async fn main() -> Result<()> {
    let hudi_table = HudiTableBuilder::from_base_uri("s3://bucket/trips_table")
        .with_option("aws_region", "us-west-2")
        .build().await?;
    Ok(())
}

Contributing

Check out the contributing guide for all the details about making contributions to the project.

Metadata

Release files for hudi 0.5.0

For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.

Source distribution (sdist)

Source distribution for hudi 0.5.0
File Size Uploaded
hudi-0.5.0.tar.gz 9.7 MB Details

Built distributions (wheels)

Table of built distributions (wheels) for hudi 0.5.0
File
hudi-0.5.0-cp310-abi3-win_amd64.whl CPython 3.10 abi3 Windows x86-64 Details
hudi-0.5.0-cp310-abi3-manylinux_2_17_x86_64.manylinux2014_x86_64.whl CPython 3.10 abi3 Linux glibc 2.17+ x86-64 Details
hudi-0.5.0-cp310-abi3-manylinux_2_17_aarch64.manylinux2014_aarch64.whl CPython 3.10 abi3 Linux glibc 2.17+ ARM64 Details
hudi-0.5.0-cp310-abi3-macosx_11_0_arm64.whl CPython 3.10 abi3 macOS 11.0+ ARM64 Details
hudi-0.5.0-cp310-abi3-macosx_10_13_x86_64.whl CPython 3.10 abi3 macOS 10.13+ x86-64 Details

Total release size: 217.5 MB

Release files / hudi-0.5.0.tar.gz

Download URL hudi-0.5.0.tar.gz
Size 9.7 MB
Tags Source
SHA-256 checksum
How to use checksums
dc8e1f3405d262491b97e8f009bcae2cbd66279dd2a8300b738ecdc3d9cba51d
BLAKE2b-256 checksum
How to use checksums
c8c199b8e45b3f20f1d136f3036430f3706e933174e6c8476d932cbc3d6cbc02
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via maturin/1.15.0

Release files / hudi-0.5.0-cp310-abi3-win_amd64.whl

Download URL hudi-0.5.0-cp310-abi3-win_amd64.whl
Size 37.4 MB
Tags CPython 3.10 Windows x86-64 abi3
SHA-256 checksum
How to use checksums
a17831c43bfb99460f251e53e8900b4b52bb0907dcb41235df48809c2748ad43
BLAKE2b-256 checksum
How to use checksums
7acfc2f25d1f828005ea6a2d2edf4e166a78cebc92404475decd867c4292a505
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via maturin/1.15.0

Release files / hudi-0.5.0-cp310-abi3-manylinux_2_17_x86_64.manylinux2014_x86_64.whl

Download URL hudi-0.5.0-cp310-abi3-manylinux_2_17_x86_64.manylinux2014_x86_64.whl
Size 45.5 MB
Tags CPython 3.10 Linux glibc 2.17+ x86-64 abi3
SHA-256 checksum
How to use checksums
241356603960ad33962fa9fc48507c414552169f95082a4d4ba42bf62c66536f
BLAKE2b-256 checksum
How to use checksums
6bb8fb73e886d6c20203dceff0ce96e16ca6fd24060e05357ba9943e65c90f09
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via maturin/1.15.0

Release files / hudi-0.5.0-cp310-abi3-manylinux_2_17_aarch64.manylinux2014_aarch64.whl

Download URL hudi-0.5.0-cp310-abi3-manylinux_2_17_aarch64.manylinux2014_aarch64.whl
Size 44.8 MB
Tags CPython 3.10 Linux glibc 2.17+ ARM64 abi3
SHA-256 checksum
How to use checksums
6020363c1642b14fb9fe84170923208ab48fde0505aac4c18c0c8df450b5a17c
BLAKE2b-256 checksum
How to use checksums
8d4f462abe2717f366da5569d2f0e1710e9bdec46a050f4bb03f71830c40785f
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via maturin/1.15.0

Release files / hudi-0.5.0-cp310-abi3-macosx_11_0_arm64.whl

Download URL hudi-0.5.0-cp310-abi3-macosx_11_0_arm64.whl
Size 39.0 MB
Tags CPython 3.10 abi3 macOS 11.0+ ARM64
SHA-256 checksum
How to use checksums
89619571d0848922f679e7cade3a24dd32894755cd5739701908e35a3e48d322
BLAKE2b-256 checksum
How to use checksums
fa5c9c63c967604c0f771cd031f9be1bf2593890a1f85476b1d42be632dd3d6e
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via maturin/1.15.0

Release files / hudi-0.5.0-cp310-abi3-macosx_10_13_x86_64.whl

Download URL hudi-0.5.0-cp310-abi3-macosx_10_13_x86_64.whl
Size 41.3 MB
Tags CPython 3.10 abi3 macOS 10.13+ x86-64
SHA-256 checksum
How to use checksums
1e74026233673856b263f926ea41ebf8d1c74edb3c5fc5994b6c79864e3c8543
BLAKE2b-256 checksum
How to use checksums
7f54d84bbba15df24191d3ed14830b6e5d6dfd566ab44608f3ef4cb50d9f9e8f
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via maturin/1.15.0
Anthropic, PBC Visionary sponsor Bloomberg Visionary sponsor Hudson River Trading Visionary sponsor Meta Visionary sponsor NVIDIA Visionary sponsor Microsoft Sustainability sponsor Depot Continuous Integration AWS Cloud computing and Security Sponsor Datadog Monitoring Fastly CDN Google Download Analytics Sentry Error logging StatusPage Status page