apache / apache/arrow-rs

[arrow-avro] Add AsyncWriter

Open
#9,212 3 comments 1 reaction 0 assignees View on GitHub
enhancement
Dominant language
Rust
Stars
3.6k
Forks
1.3k
Avg merge
2d 18h
Merged PRs (30d)
169

Description

**Is your feature request related to a problem or challenge? Please describe what you are trying to do.**

`arrow-avro` currently provides synchronous writing APIs (`Writer` and the aliases `AvroWriter` for OCF and `AvroStreamWriter` for SOE). In async-first applications (Tokio-based services, DataFusion execution, writing to `object_store`, HTTP streaming responses, etc.), these sync APIs force one of the following suboptimal approaches:
* Use blocking I/O (`std::fs::File`) inside an async runtime (risking thread starvation / reduced throughput).
* Use `spawn_blocking` to isolate blocking writes (adds complexity and makes backpressure/error handling harder).
* Buffer the entire output into memory (e.g., write to `Vec`) before uploading to object storage (bad for large outputs).

In parallel, PR #8930 is adding an `arrow-avro` **AsyncReader** for Avro OCF (with `object_store` integration). To enable end-to-end async pipelines, `arrow-avro` should also have a corresponding **AsyncWriter**.

Concrete use cases:

* Write Avro OCF results directly to `tokio::fs::File` without blocking.
* Stream an Avro OCF response body over the network (chunked transfer, back-pressure aware).
* Write Avro OCF to S3/GCS/Azure via `object_store` multipart upload (symmetry with the new AsyncReader).

**Describe the solution you'd like**

Add a new `arrow-avro` async writer API that mirrors the patterns used by Parquet’s async writer (`parquet/src/arrow/async_writer`) and interoperates with the `arrow-avro` AsyncReader being introduced.

Proposed high-level design:

1. **New module and feature gating**

* Add `arrow-avro/src/writer/async_writer.rs` (or `writer/async_writer/mod.rs`) behind `cfg(feature = "async")`
* Align feature flags with the AsyncReader PR (`async`, and optionally `object_store`)

2. **Async sink abstraction (match Parquet’s pattern)**

* Introduce an `AsyncFileWriter` trait analogous to Parquet’s:
* Accepts `bytes::Bytes` payloads
* Returns `futures::future::BoxFuture<'_, Result<(), ArrowError>>`
* Has `write(Bytes)` and `complete()` methods
* Provide:

* `impl AsyncFileWriter for Box`
* A blanket impl for Tokio sinks (like Parquet): `impl AsyncFileWriter for T`

3. **AsyncWriter types and API parity**

* Implement a generic `AsyncWriter` where `W: AsyncFileWriter` and `F: AvroFormat`
* Provide aliases consistent with the existing sync types:
* `type AsyncAvroWriter = AsyncWriter`
* `type AsyncAvroStreamWriter = AsyncWriter`
* Public API should be intentionally close to the existing sync `Writer`:
* `WriterBuilder::build_async::(writer)` (preferred: reuses existing builder options)
* `async fn write(&mut self, batch: &RecordBatch) -> Result<(), ArrowError>`
* `async fn write_batches(&mut self, batches: &[&RecordBatch]) -> Result<(), ArrowError>`
* `async fn finish(&mut self) -> Result<(), ArrowError>` (flush + `complete`)
* (Optional ergonomics) `async fn close(self) -> Result` to consume and return the underlying writer

4. **Implementation approach (buffer + flush, like Parquet)**

To avoid rewriting the encoder as “async”, reuse the existing synchronous encoding logic by staging into an in-memory buffer and then pushing the bytes into the async sink:
* On construction:
* Encode the OCF header via existing `AvroFormat::start_stream`, but write it into a `Vec` buffer
* Flush that buffer to `AsyncFileWriter::write(Bytes)`
* On `write`:
* For **OCF**:
* Encode the batch into a `Vec` using the existing `RecordEncoder`
* Apply block-level compression exactly as the sync writer does
* Append block framing (`count`, `block_size`, `block_bytes`, `sync_marker`) into a staging buffer
* Flush staging buffer to the async sink
* For **SOE stream**:
* Encode into a staging buffer and flush (preserving the existing prefix behavior)
* Reuse `WriterBuilder` settings (`with_compression`, `with_capacity`, `with_row_capacity`, `with_fingerprint_strategy`) to size buffers and control prefixes/compression.

5. **Interop and test plan (explicit requirement)**

Add tests demonstrating interoperability with the AsyncReader in #8930:
* Roundtrip test: `AsyncAvroWriter` → bytes/file/object_store → `AsyncAvroReader` yields the same `RecordBatch` values
* Also validate roundtrip with the existing sync `ReaderBuilder` to ensure backwards compatibility
* Cover both:
* uncompressed OCF
* compressed OCF (as supported by `CompressionCodec`)

6. **object_store writer adapter**

For symmetry with `parquet::arrow::async_writer::ParquetObjectWriter` and the AsyncReader’s `object_store` integration add an `AvroObjectWriter` (behind `feature = "object_store"`) implementing `AsyncFileWriter` and performing multipart upload.

Illustrative API sketch:

```rust
use tokio::fs::File;
use arrow_avro::writer::{WriterBuilder /*, AsyncAvroWriter */};

# async fn example(schema: arrow_schema::Schema, batch: arrow_array::RecordBatch) -> Result<(), arrow_schema::ArrowError> {
let file = File::create("out.avro").await.map_err(|e| arrow_schema::ArrowError::IoError(e.to_string(), e))?;

// via builder (preferred)
let mut w = WriterBuilder::new(schema)
.with_compression(None)
.with_capacity(64 * 1024)
.build_async::<_, arrow_avro::writer::AvroOcfFormat>(file)?;

// write and finish
w.write(&batch).await?;
w.finish().await?;
# Ok(())
# }
```

**Describe alternatives you've considered**

* **Use synchronous `AvroWriter` in `spawn_blocking`**

* Avoids blocking the async runtime directly, but complicates control flow and backpressure, and can still become a bottleneck at scale.
* **Write to `Vec` with sync `AvroWriter` and then upload**
* Requires buffering entire outputs, increasing memory usage and latency.

**Additional context**

Notes / open questions:
* OCF sync markers are generated randomly today for `AvroOcfFormat`. If we want any “byte-for-byte equality” tests between sync and async writers, we may need a way to inject a deterministic sync marker in tests, or rely purely on decode/roundtrip validation.
* Naming bikeshed: should the terminal method be `finish().await` (matching the current sync `Writer`) or `close().await` (matching Parquet’s async writer)? Either is fine, but consistent naming would be nice.

Contributor guide

Open the contributing guide

Research direction

Start by reading the existing arrow-avro writer and WriterBuilder, then compare parquet/src/arrow/async_writer and the AsyncReader work in PR #8930. The implementation would add arrow-avro/src/writer/async_writer.rs, async sink and writer APIs, and optional object_store support. Done means async OCF/SOE output round-trips through AsyncReader and the existing sync ReaderBuilder for uncompressed and supported compressed OCF.

Written by the indexing model from the issue text.

Assessment

Tech stack
rust
Domain
backend, data
Issue type
Feature
Difficulty
5/5
Estimated time
Over a week
Activity status
Stale
Clarity
Mostly clear
Newbie friendliness
25/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.