Metadata-Version: 2.4
Name: hudi
Version: 0.5.0
Classifier: Programming Language :: Python :: 3
Classifier: Programming Language :: Python :: 3.10
Classifier: Programming Language :: Python :: 3.11
Classifier: Programming Language :: Python :: 3.12
Classifier: Programming Language :: Python :: 3.13
Requires-Dist: pyarrow>=11.0.0
Requires-Dist: pyarrow-hotfix
Requires-Dist: datafusion>=54.0.0,<55 ; extra == 'datafusion'
Requires-Dist: pytest==9.1.1 ; extra == 'devel'
Requires-Dist: coverage==7.16.0 ; extra == 'devel'
Requires-Dist: ruff==0.15.18 ; extra == 'devel'
Requires-Dist: mypy==2.3.1 ; extra == 'devel'
Requires-Dist: pre-commit ; extra == 'devel'
Requires-Dist: ruff==0.15.18 ; extra == 'lint'
Requires-Dist: mypy==2.3.1 ; extra == 'lint'
Provides-Extra: datafusion
Provides-Extra: devel
Provides-Extra: lint
Summary: Native Python binding for Apache Hudi, based on hudi-rs.
Keywords: apachehudi,hudi,datalake,arrow
Home-Page: https://github.com/apache/hudi-rs
License-Expression: Apache-2.0
Requires-Python: >=3.10
Description-Content-Type: text/markdown; charset=UTF-8; variant=GFM
Project-URL: repository, https://github.com/apache/hudi-rs/tree/main/python/

<!--
  ~ Licensed to the Apache Software Foundation (ASF) under one
  ~ or more contributor license agreements.  See the NOTICE file
  ~ distributed with this work for additional information
  ~ regarding copyright ownership.  The ASF licenses this file
  ~ to you under the Apache License, Version 2.0 (the
  ~ "License"); you may not use this file except in compliance
  ~ with the License.  You may obtain a copy of the License at
  ~
  ~   http://www.apache.org/licenses/LICENSE-2.0
  ~
  ~ Unless required by applicable law or agreed to in writing,
  ~ software distributed under the License is distributed on an
  ~ "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
  ~ KIND, either express or implied.  See the License for the
  ~ specific language governing permissions and limitations
  ~ under the License.
-->

<p align="center">
  <a href="https://hudi.apache.org/">
    <img src="https://hudi.apache.org/assets/images/hudi_logo_transparent_1400x600.png" alt="Hudi logo" height="120px">
  </a>
</p>
<p align="center">
  The native Rust implementation for Apache Hudi, with C++ & Python API bindings.
  <br>
  <br>
  <a href="https://github.com/apache/hudi-rs/actions/workflows/ci.yml">
    <img alt="hudi-rs ci" src="https://github.com/apache/hudi-rs/actions/workflows/ci.yml/badge.svg">
  </a>
  <a href="https://codecov.io/github/apache/hudi-rs">
    <img alt="hudi-rs codecov" src="https://codecov.io/github/apache/hudi-rs/graph/badge.svg">
  </a>
  <a href="https://hudi.apache.org/slack">
    <img alt="join hudi slack" src="https://img.shields.io/badge/slack-%23hudi-72eff8?logo=slack&color=48c628">
  </a>
  <a href="https://x.com/apachehudi">
    <img alt="follow hudi x/twitter" src="https://img.shields.io/twitter/follow/apachehudi?label=apachehudi">
  </a>
  <a href="https://www.linkedin.com/company/apache-hudi">
    <img alt="follow hudi linkedin" src="https://img.shields.io/badge/apache%E2%80%93hudi-0077B5?logo=linkedin">
  </a>
</p>

The Hudi-rs project aims to standardize the core [Apache Hudi](https://github.com/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**][pypi]    | [![][pypi-badge]][pypi]     | `pip install hudi`   |
| [**Crates.io**][crates] | [![][crates-badge]][crates] | `cargo add hudi`     |

[pypi]: https://pypi.org/project/hudi/
[pypi-badge]: https://img.shields.io/pypi/dm/hudi?style=flat-square&color=51AEF3
[crates]: https://crates.io/crates/hudi
[crates-badge]: https://img.shields.io/crates/d/hudi?style=flat-square&color=163669

The `hudi` crate carries two features: `datafusion` (off by default, see
[Apache DataFusion](#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](https://hudi.apache.org/docs/quick-start-guide).

For the full reader API reference (`ReadOptions`, filter expressions, behavioral guarantees), see [docs/reader-spec.md](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

```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

```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

```python
from hudi import HudiReadOptions

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

#### Rust

```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

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

#### Rust

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

<details>
<summary>Supported timestamp formats</summary>

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.
</details>

### Incremental Query

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

#### Python

```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

```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 `RecordBatch`es one at a time without loading the full result into memory.
The same `ReadOptions` knobs apply, plus `batch_size` and `projection`.

#### Python

```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

```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

```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

```rust,ignore
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++

```cpp
#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](#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](https://datafusion.apache.org/)
extension to query Hudi tables.

<details>
<summary>Add crate hudi with datafusion feature to your application to query a Hudi table.</summary>

```shell
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`.

</details>

<details>
<summary>Add python hudi with datafusion feature to your application to query a Hudi table.</summary>

```shell
pip install hudi[datafusion]
```
</details>

#### Rust
```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
```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

- [Daft](https://docs.daft.ai/en/stable/connectors/hudi/)
- [Ray](https://docs.ray.io/en/latest/data/api/doc/ray.data.read_hudi.html#ray.data.read_hudi)

### 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

```python
from hudi import HudiTableBuilder

hudi_table = (
    HudiTableBuilder
    .from_base_uri("s3://bucket/trips_table")
    .with_option("aws_region", "us-west-2")
    .build()
)
```

#### Rust

```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](./CONTRIBUTING.md) for all the details about making contributions to the project.

