Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 7 additions & 4 deletions CLAUDE.md
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,7 @@ This file provides guidance to Claude Code (claude.ai/code) when working with co
Cargo workspace with two crates:

- `ruhvro/` — the core Rust library (published as `ruhvro` on crates.io). Pure Rust API for serializing/deserializing schemaless Avro to/from Arrow `RecordBatch`es. Has no Python deps.
- `src/lib.rs` (top-level) — the `pyruhvro` PyO3 extension module that wraps `ruhvro` and exposes it to Python via maturin. Re-exports three functions: `deserialize_array`, `deserialize_array_threaded`, `serialize_record_batch`.
- `src/lib.rs` (top-level) — the `pyruhvro` PyO3 extension module that wraps `ruhvro` and exposes it to Python via maturin. Exposes `deserialize_array`, `deserialize_array_threaded`, `serialize_record_batch`, plus `_spawn` variants that use `tokio::spawn` instead of `spawn_blocking`. Also maintains a `String -> Arc<Schema>` cache so repeat calls don't re-parse the schema JSON, and releases the GIL via `Python::detach` around every Rust call so multiple Python threads can run concurrently.

Keep Python-facing concerns in the top-level crate; keep the Avro↔Arrow logic in `ruhvro/`. The PyO3 wrappers should stay thin — convert PyArrow types, call into `ruhvro`, return PyArrow types.

Expand Down Expand Up @@ -42,15 +42,18 @@ The pipeline is **schemaless Avro bytes ⇄ Arrow `RecordBatch`**, driven by a p

Key modules in `ruhvro/src/`:

- `deserialize.rs` — public entry points `parse_schema`, `per_datum_deserialize` (single-threaded), `per_datum_deserialize_threaded` (rayon, splits input into `num_chunks` slices and returns one `RecordBatch` per chunk).
- `serialize.rs` — public entry point `serialize_record_batch`. Converts the `RecordBatch` into a `StructArray`, slices it into `num_chunks`, and serializes each slice in parallel via rayon into a `GenericBinaryArray<i32>` of Avro datums.
- `deserialize.rs` — public entry points `parse_schema`, `per_datum_deserialize` (single-threaded), `per_datum_deserialize_threaded` (tokio `spawn_blocking`, splits input into `num_chunks` slices and returns one `RecordBatch` per chunk), plus `_spawn` variant on the work-stealing async pool. Threaded variants take an `Arc<Schema>` so callers can share one parsed schema across many calls without re-cloning.
- `serialize.rs` — public entry point `serialize_record_batch` (same `Arc<Schema>` convention). Converts the `RecordBatch` into a `StructArray`, slices it into `num_chunks`, and serializes each slice in parallel via tokio `spawn_blocking` into a `GenericBinaryArray<i32>` of Avro datums.
- `fast_decode.rs` / `fast_encode.rs` — schema-walking decoder/encoder that bypasses the `apache_avro::Value` tree and writes straight into Arrow builders. Gated by `is_supported(schema)`; falls back to the `Value`-based path for schemas containing types outside the supported subset.
- `schema_translate.rs` — converts an `apache_avro::Schema` into an `arrow::datatypes::Schema`. This is the source of truth for type mapping (e.g. nullable-union → nullable Arrow field, multi-variant unions → Arrow `Union`, Avro `map` → Arrow `Map`, logical types → Arrow temporal types).
- `complex.rs` — `AvroToArrowBuilder` and its `Struct`/`List`/`Union`/`Map`/`Primitive` variants. This is the *deserialize* side: walks Avro `Value`s into Arrow builders. The `add_val!` macro and `get_val_from_possible_union` helper handle the common "value might be wrapped in a union" case.
- `serialization_containers.rs` — the *serialize* side: `ArrayContainers` walks Arrow arrays column-wise and re-emits `apache_avro::types::Value`s, then `to_avro_datum` encodes each row.

### Threading model

Both threaded paths use rayon and require the caller to pass `num_chunks` explicitly — the library does not infer it from `rayon::current_num_threads()`. The Python wrappers release the GIL implicitly via PyO3 while inside `ruhvro` calls, which is the whole point of the project per the README.
A single global tokio multi-thread runtime (`OnceLock` in `ruhvro/src/lib.rs`) services all parallel work — created on first call, alive for the process lifetime. Worker threads are parked when idle, so the runtime costs nothing when unused.

Both threaded paths require the caller to pass `num_chunks` explicitly — the library does not infer it. `num_chunks` is clamped to `[1, max(rows, 1)]` so `0` doesn't panic and overshooting the row count doesn't spawn empty tasks. The Python wrappers explicitly release the GIL with `py.detach(...)` around every Rust call, so multiple Python threads can call into pyruhvro concurrently and benefit from the internal parallelism.

### Avro union handling

Expand Down
11 changes: 10 additions & 1 deletion Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 1 addition & 1 deletion ruhvro/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,7 @@ description = "Fast, multi-threaded deserialization of schema-less avro encoded
repository = "https://github.com/Tyler-Sch/pyruhvro"

[dependencies]
rayon = "1.10"
tokio = { version = "1", features = ["rt-multi-thread"] }
apache-avro = "0.21"
arrow = "58"
anyhow = "1.0"
Expand Down
82 changes: 45 additions & 37 deletions ruhvro/benches/deserialize.rs
Original file line number Diff line number Diff line change
@@ -1,74 +1,82 @@
//! Criterion benchmarks for ruhvro's deserialize path.
//!
//! Each bench:
//! 1. Builds N avro-encoded records once (in setup).
//! 2. Measures `per_datum_deserialize` (single-threaded) and
//! `per_datum_deserialize_threaded` (8 chunks).
//!
//! Run with `cargo bench -p ruhvro --bench deserialize`.
//! HTML reports land in `target/criterion`.
//! Compares single-threaded, spawn_blocking, tokio::spawn, and rayon
//! across 1k / 10k / 100k record counts.

mod common;

use std::sync::Arc;

use apache_avro::Schema as AvroSchema;
use criterion::{black_box, criterion_group, criterion_main, Criterion, Throughput};
use ruhvro::deserialize::{per_datum_deserialize, per_datum_deserialize_threaded};
use ruhvro::deserialize::{
per_datum_deserialize,
per_datum_deserialize_threaded,
per_datum_deserialize_threaded_spawn,
};

const N: usize = 1_000;
const SIZES: &[usize] = &[1_000, 10_000];

fn run_group(c: &mut Criterion, name: &str, parsed: AvroSchema, encoded: Vec<Vec<u8>>) {
fn run_group(
c: &mut Criterion,
schema_name: &str,
parsed: AvroSchema,
encoded: Vec<Vec<u8>>,
n: usize,
) {
let refs: Vec<&[u8]> = encoded.iter().map(Vec::as_slice).collect();
let schema_arc = Arc::new(parsed.clone());
let group_name = format!("{schema_name}/{n}");
let mut group = c.benchmark_group(&group_name);
group.throughput(Throughput::Elements(n as u64));

let mut group = c.benchmark_group(name);
group.throughput(Throughput::Elements(N as u64));
group.bench_function("single_threaded", |b| {
b.iter(|| black_box(per_datum_deserialize(black_box(&refs), black_box(&parsed)).unwrap()))
});

group.bench_function("per_datum_deserialize", |b| {
group.bench_function("spawn_blocking", |b| {
b.iter(|| {
let rb = per_datum_deserialize(black_box(&refs), black_box(&parsed)).unwrap();
black_box(rb);
black_box(per_datum_deserialize_threaded(
black_box(refs.clone()), black_box(Arc::clone(&schema_arc)), common::NUM_CHUNKS,
).unwrap())
})
});

group.bench_function("per_datum_deserialize_threaded", |b| {
group.bench_function("tokio_spawn", |b| {
b.iter(|| {
let rbs = per_datum_deserialize_threaded(
black_box(refs.clone()),
black_box(&parsed),
common::NUM_CHUNKS,
)
.unwrap();
black_box(rbs);
black_box(per_datum_deserialize_threaded_spawn(
black_box(refs.clone()), black_box(Arc::clone(&schema_arc)), common::NUM_CHUNKS,
).unwrap())
})
});

group.finish();
}

fn bench_flat_primitives(c: &mut Criterion) {
let (parsed, encoded) = common::flat_primitives(N);
run_group(c, "flat_primitives", parsed, encoded);
}

fn bench_nullable_primitives(c: &mut Criterion) {
let (parsed, encoded) = common::nullable_primitives(N);
run_group(c, "nullable_primitives", parsed, encoded);
for &n in SIZES {
let (parsed, encoded) = common::flat_primitives(n);
run_group(c, "flat_primitives", parsed, encoded, n);
}
}

fn bench_nested_struct(c: &mut Criterion) {
let (parsed, encoded) = common::nested_struct(N);
run_group(c, "nested_struct", parsed, encoded);
for &n in SIZES {
let (parsed, encoded) = common::nested_struct(n);
run_group(c, "nested_struct", parsed, encoded, n);
}
}

fn bench_array_and_map(c: &mut Criterion) {
let (parsed, encoded) = common::array_and_map(N);
run_group(c, "array_and_map", parsed, encoded);
for &n in SIZES {
let (parsed, encoded) = common::array_and_map(n);
run_group(c, "array_and_map", parsed, encoded, n);
}
}

criterion_group!(
benches,
bench_flat_primitives,
bench_nullable_primitives,
bench_nested_struct,
bench_array_and_map
bench_array_and_map,
);
criterion_main!(benches);
85 changes: 49 additions & 36 deletions ruhvro/benches/serialize.rs
Original file line number Diff line number Diff line change
@@ -1,78 +1,91 @@
//! Criterion benchmarks for ruhvro's serialize path.
//!
//! Each bench:
//! 1. Builds N avro-encoded records, deserializes them once into a
//! `RecordBatch` (in setup, outside the timing loop).
//! 2. Measures `serialize_record_batch` with 1 chunk (single-threaded) and
//! 8 chunks (rayon-parallel).
//! Compares single-threaded, spawn_blocking (tokio), tokio::spawn, and rayon.
//!
//! Run with `cargo bench -p ruhvro --bench serialize`.

mod common;

use std::sync::Arc;

use apache_avro::Schema as AvroSchema;
use arrow::array::RecordBatch;
use criterion::{black_box, criterion_group, criterion_main, Criterion, Throughput};
use ruhvro::deserialize::per_datum_deserialize;
use ruhvro::serialize::serialize_record_batch;
use ruhvro::serialize::{
serialize_record_batch,
serialize_record_batch_spawn,
};

const N: usize = 1_000;
const SIZES: &[usize] = &[1_000, 10_000];

fn prepare_batch(parsed: &AvroSchema, encoded: &[Vec<u8>]) -> RecordBatch {
let refs: Vec<&[u8]> = encoded.iter().map(Vec::as_slice).collect();
per_datum_deserialize(&refs, parsed).unwrap()
}

fn run_group(c: &mut Criterion, name: &str, parsed: AvroSchema, batch: RecordBatch) {
let mut group = c.benchmark_group(name);
group.throughput(Throughput::Elements(N as u64));
fn run_group(c: &mut Criterion, schema_name: &str, parsed: AvroSchema, batch: RecordBatch, n: usize) {
let schema_arc = Arc::new(parsed);
let group_name = format!("{schema_name}/{n}");
let mut group = c.benchmark_group(&group_name);
group.throughput(Throughput::Elements(n as u64));

group.bench_function("single_threaded", |b| {
b.iter(|| {
black_box(serialize_record_batch(
black_box(batch.clone()), black_box(Arc::clone(&schema_arc)), 1,
).unwrap())
})
});

group.bench_function("serialize_record_batch_1chunk", |b| {
group.bench_function("spawn_blocking", |b| {
b.iter(|| {
// RecordBatch holds Arc'd columns, so .clone() is cheap (refcount bumps).
let bytes = serialize_record_batch(black_box(batch.clone()), black_box(&parsed), 1)
.unwrap();
black_box(bytes);
black_box(serialize_record_batch(
black_box(batch.clone()), black_box(Arc::clone(&schema_arc)), common::NUM_CHUNKS,
).unwrap())
})
});

group.bench_function("serialize_record_batch_8chunks", |b| {
group.bench_function("tokio_spawn", |b| {
b.iter(|| {
let bytes = serialize_record_batch(
black_box(batch.clone()),
black_box(&parsed),
common::NUM_CHUNKS,
)
.unwrap();
black_box(bytes);
black_box(serialize_record_batch_spawn(
black_box(batch.clone()), black_box(Arc::clone(&schema_arc)), common::NUM_CHUNKS,
).unwrap())
})
});

group.finish();
}

fn bench_flat_primitives(c: &mut Criterion) {
let (parsed, encoded) = common::flat_primitives(N);
let batch = prepare_batch(&parsed, &encoded);
run_group(c, "flat_primitives", parsed, batch);
for &n in SIZES {
let (parsed, encoded) = common::flat_primitives(n);
let batch = prepare_batch(&parsed, &encoded);
run_group(c, "flat_primitives", parsed, batch, n);
}
}

fn bench_nullable_primitives(c: &mut Criterion) {
let (parsed, encoded) = common::nullable_primitives(N);
let batch = prepare_batch(&parsed, &encoded);
run_group(c, "nullable_primitives", parsed, batch);
for &n in SIZES {
let (parsed, encoded) = common::nullable_primitives(n);
let batch = prepare_batch(&parsed, &encoded);
run_group(c, "nullable_primitives", parsed, batch, n);
}
}

fn bench_nested_struct(c: &mut Criterion) {
let (parsed, encoded) = common::nested_struct(N);
let batch = prepare_batch(&parsed, &encoded);
run_group(c, "nested_struct", parsed, batch);
for &n in SIZES {
let (parsed, encoded) = common::nested_struct(n);
let batch = prepare_batch(&parsed, &encoded);
run_group(c, "nested_struct", parsed, batch, n);
}
}

fn bench_array_and_map(c: &mut Criterion) {
let (parsed, encoded) = common::array_and_map(N);
let batch = prepare_batch(&parsed, &encoded);
run_group(c, "array_and_map", parsed, batch);
for &n in SIZES {
let (parsed, encoded) = common::array_and_map(n);
let batch = prepare_batch(&parsed, &encoded);
run_group(c, "array_and_map", parsed, batch, n);
}
}

criterion_group!(
Expand Down
3 changes: 2 additions & 1 deletion ruhvro/examples/prof_decode.rs
Original file line number Diff line number Diff line change
Expand Up @@ -40,6 +40,7 @@ fn build_record(i: usize, sch: &AvroSchema) -> Value {

fn main() {
let parsed = parse_schema(SCHEMA).unwrap();
let schema_arc = std::sync::Arc::new(parsed.clone());
let n = 1_000;
let encoded: Vec<Vec<u8>> = (0..n)
.map(|i| to_avro_datum(&parsed, build_record(i, &parsed)).unwrap())
Expand All @@ -49,7 +50,7 @@ fn main() {
let mut total_rows: usize = 0;
for iter in 0..10_000 {
let refs: Vec<&[u8]> = encoded.iter().map(Vec::as_slice).collect();
let rbs = per_datum_deserialize_threaded(refs, &parsed, 8).unwrap();
let rbs = per_datum_deserialize_threaded(refs, std::sync::Arc::clone(&schema_arc), 8).unwrap();
total_rows += rbs.iter().map(|rb| rb.num_rows()).sum::<usize>();
if iter % 1000 == 0 {
eprintln!(" iter {iter}");
Expand Down
Loading
Loading