This is an automated email from the ASF dual-hosted git repository.
alamb pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/arrow-rs.git
The following commit(s) were added to refs/heads/main by this push:
new fd8ead5d3c arrow-avro: Deprecate object_store integration (#10484)
fd8ead5d3c is described below
commit fd8ead5d3cd17f85cf0455901372ff1fa51b2563
Author: Frederic Branczyk <[email protected]>
AuthorDate: Fri Jul 31 15:08:31 2026 +0200
arrow-avro: Deprecate object_store integration (#10484)
Deprecate AvroObjectReader in favor of implementing AsyncFileReader
directly (see the example on the trait docs and
arrow-avro/examples/object_store.rs). To reduce the code needed by
implementors I also added SpawnedReader. Still keeping object_store as a
dev-dependency so we can ensure compatibility and use it in testing.
# Which issue does this PR close?
- Part of #10308
# Rationale for this change
See the issue
# What changes are included in this PR?
Decoupling and removal of object_store dependencies in arrow-avro.
# Are these changes tested?
Yes, see unit tests
# Are there any user-facing changes?
No, but deprecation warnings that there will be.
@alamb
---
arrow-avro/Cargo.toml | 8 +
arrow-avro/README.md | 34 ++--
arrow-avro/examples/object_store.rs | 133 ++++++++++++++++
arrow-avro/src/lib.rs | 38 ++---
.../src/reader/async_reader/async_file_reader.rs | 52 +++++-
arrow-avro/src/reader/async_reader/mod.rs | 78 ++++++---
arrow-avro/src/reader/async_reader/spawn.rs | 176 +++++++++++++++++++++
arrow-avro/src/reader/async_reader/store.rs | 11 ++
arrow-avro/src/reader/mod.rs | 3 +-
9 files changed, 475 insertions(+), 58 deletions(-)
diff --git a/arrow-avro/Cargo.toml b/arrow-avro/Cargo.toml
index 0cc5af7010..a4c20f0553 100644
--- a/arrow-avro/Cargo.toml
+++ b/arrow-avro/Cargo.toml
@@ -48,6 +48,9 @@ avro_custom_types = ["dep:arrow-select"]
# Enable async APIs
async = ["futures", "tokio"]
# Enable object_store integration
+# Deprecated: implement `AsyncFileReader` directly instead, see the example on
+# the `AsyncFileReader` trait documentation and
`arrow-avro/examples/object_store.rs`.
+# This feature will be removed in a future release.
object_store = ["dep:object_store", "async"]
[dependencies]
@@ -98,6 +101,11 @@ once_cell = "1.21.3"
half = { version = "2.1", default-features = false }
tokio = { version = "1.0", default-features = false, features = ["macros",
"rt-multi-thread", "io-util", "fs"] }
+[[example]]
+name = "object_store"
+required-features = ["async"]
+path = "./examples/object_store.rs"
+
[[bench]]
name = "avro_reader"
harness = false
diff --git a/arrow-avro/README.md b/arrow-avro/README.md
index dbc1e1760e..2a9c25a421 100644
--- a/arrow-avro/README.md
+++ b/arrow-avro/README.md
@@ -105,25 +105,18 @@ fn main() -> anyhow::Result<()> {
See the crate docs for runnable SOE and Confluent round‑trip examples.
-### Async reading from object stores (`object_store` feature)
+### Async reading (`async` feature)
```rust,ignore
-use std::sync::Arc;
-use arrow_avro::reader::{AsyncAvroFileReader, AvroObjectReader};
+use arrow_avro::reader::AsyncAvroFileReader;
use futures::TryStreamExt;
-use object_store::ObjectStore;
-use object_store::local::LocalFileSystem;
-use object_store::path::Path;
#[tokio::main]
async fn main() -> anyhow::Result<()> {
- let store: Arc<dyn ObjectStore> = Arc::new(LocalFileSystem::new());
- let path = Path::from("data/example.avro");
-
- let meta = store.head(&path).await?;
- let reader = AvroObjectReader::new(store, path);
+ let file = tokio::fs::File::open("data/example.avro").await?;
+ let file_size = file.metadata().await?.len();
- let stream = AsyncAvroFileReader::builder(reader, meta.size, 1024)
+ let stream = AsyncAvroFileReader::builder(file, file_size, 1024)
.try_build()
.await?;
@@ -135,6 +128,13 @@ async fn main() -> anyhow::Result<()> {
}
```
+Any `AsyncFileReader` implementation can be used as the source, so object
+storage services (S3, GCS, Azure Blob, etc.) can be integrated by implementing
+`AsyncFileReader` on top of a client such as the [`object_store`] crate. See
+the example on the `AsyncFileReader` trait documentation.
+
+[`object_store`]: https://docs.rs/object_store/latest/object_store/
+
---
## Feature Flags (what they do and when to use them)
@@ -158,12 +158,12 @@ async fn main() -> anyhow::Result<()> {
* Only **OCF** uses these codecs (they compress per‑block). They do **not**
apply to raw Avro frames used by Confluent wire format or SOE. The crate’s
`compression` module is specifically for **OCF blocks**.
* `deflate` uses `flate2` with the `rust_backend` (no system zlib required).
-### Async & Object Store
+### Async
| Feature | Default | What it enables
| When to use
|
|----------------|--------:|-----------------------------------------------------------------------------|-------------------------------------------------------------------------------|
-| `async` | ⬜ | Async APIs for reading Avro via `futures` and
`tokio` | Enable for non-blocking async Avro reading with
`AsyncAvroFileReader`. |
-| `object_store` | ⬜ | Integration with `object_store` crate (implies
`async`) | Enable for reading Avro from cloud storage (S3,
GCS, Azure Blob, etc.). |
+| `async` | ⬜ | Async APIs for reading Avro via `futures` and
`tokio` | Enable for non-blocking async Avro reading with
`AsyncAvroFileReader`, including from cloud storage via a custom
`AsyncFileReader`. |
+| `object_store` | ⬜ | **Deprecated**: the deprecated `AvroObjectReader`
(implies `async`) | Do not enable in new code; implement
`AsyncFileReader` directly instead. Will be removed in a future release. |
### Schema fingerprints & custom logical type helpers
@@ -193,10 +193,10 @@ async fn main() -> anyhow::Result<()> {
```toml
arrow-avro = { version = "58", default-features = false, features =
["deflate", "snappy", "zstd"] }
```
-* Async reading from object stores (S3, GCS, etc.):
+* Async reading (including from object stores such as S3, GCS, etc.):
```toml
- arrow-avro = { version = "58", features = ["object_store"] }
+ arrow-avro = { version = "58", features = ["async"] }
```
* Fingerprint helpers:
diff --git a/arrow-avro/examples/object_store.rs
b/arrow-avro/examples/object_store.rs
new file mode 100644
index 0000000000..14c7b490be
--- /dev/null
+++ b/arrow-avro/examples/object_store.rs
@@ -0,0 +1,133 @@
+// 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.
+
+use arrow_array::{ArrayRef, Int64Array, RecordBatch};
+use arrow_avro::errors::AvroError;
+use arrow_avro::reader::{AsyncAvroFileReader, AsyncFileReader, SpawnedReader};
+use arrow_avro::writer::AvroWriter;
+use bytes::Bytes;
+use futures::future::BoxFuture;
+use futures::{FutureExt, TryStreamExt};
+use object_store::memory::InMemory;
+use object_store::path::Path;
+use object_store::{ObjectStore, ObjectStoreExt};
+use std::error::Error;
+use std::ops::Range;
+use std::sync::Arc;
+
+/// This example demonstrates reading Avro files from object storage via the
+/// [`object_store`] crate, without the deprecated `AvroObjectReader` type.
+///
+/// # Example Overview
+///
+/// 1. Writes an Avro Object Container File to an [`ObjectStore`]
+///
+/// 2. Reads it back with [`ObjectStoreReader`], a minimal [`AsyncFileReader`]
+/// implementation on top of an [`ObjectStore`] (equivalent of
+/// `AvroObjectReader`)
+///
+/// 3. Reads it again with the reader wrapped in a [`SpawnedReader`], which
+/// performs all I/O on a separate tokio runtime so that the runtime
+/// decoding Avro is not also driving the I/O (equivalent of
+/// `AvroObjectReader::with_runtime`)
+#[tokio::main]
+async fn main() -> Result<(), Box<dyn Error>> {
+ let store: Arc<dyn ObjectStore> = Arc::new(InMemory::new());
+ let path = Path::from("example.avro");
+
+ // 1. Write an Avro Object Container File to the store
+ let col = Arc::new(Int64Array::from_iter_values([1, 2, 3])) as ArrayRef;
+ let batch = RecordBatch::try_from_iter([("col", col)])?;
+
+ let mut writer = AvroWriter::new(Vec::new(),
batch.schema().as_ref().clone())?;
+ writer.write(&batch)?;
+ writer.finish()?;
+ store
+ .put(&path, Bytes::from(writer.into_inner()).into())
+ .await?;
+
+ // 2. Read it back with an `AsyncFileReader` implemented on `ObjectStore`.
+ // The builder requires the file size, which can be obtained via `head`
+ let file_size = store.head(&path).await?.size;
+
+ let reader = ObjectStoreReader::new(Arc::clone(&store), path.clone());
+ let stream = AsyncAvroFileReader::builder(reader, file_size, 1024)
+ .try_build()
+ .await?;
+ let read: Vec<RecordBatch> = stream.try_collect().await?;
+ assert_eq!(read, vec![batch.clone()]);
+ println!("read {} rows", read[0].num_rows());
+
+ // 3. Read again, performing the I/O on a dedicated runtime
+ let io_runtime = tokio::runtime::Builder::new_multi_thread()
+ .worker_threads(1)
+ .enable_all()
+ .build()
+ .expect("failed to build I/O runtime");
+
+ let reader = ObjectStoreReader::new(Arc::clone(&store), path);
+ let reader = SpawnedReader::new(reader, io_runtime.handle().clone());
+ let stream = AsyncAvroFileReader::builder(reader, file_size, 1024)
+ .try_build()
+ .await?;
+ let read: Vec<RecordBatch> = stream.try_collect().await?;
+ assert_eq!(read, vec![batch]);
+ println!("read {} rows via dedicated I/O runtime", read[0].num_rows());
+
+ io_runtime.shutdown_background();
+ Ok(())
+}
+
+/// An [`AsyncFileReader`] for a location in an [`ObjectStore`]
+///
+/// This mirrors the example on the [`AsyncFileReader`] trait documentation.
+#[derive(Clone, Debug)]
+struct ObjectStoreReader {
+ store: Arc<dyn ObjectStore>,
+ path: Path,
+}
+
+impl ObjectStoreReader {
+ fn new(store: Arc<dyn ObjectStore>, path: Path) -> Self {
+ Self { store, path }
+ }
+}
+
+impl AsyncFileReader for ObjectStoreReader {
+ fn get_bytes(&mut self, range: Range<u64>) -> BoxFuture<'_, Result<Bytes,
AvroError>> {
+ async move {
+ self.store
+ .get_range(&self.path, range)
+ .await
+ .map_err(|e| AvroError::General(e.to_string()))
+ }
+ .boxed()
+ }
+
+ fn get_byte_ranges(
+ &mut self,
+ ranges: Vec<Range<u64>>,
+ ) -> BoxFuture<'_, Result<Vec<Bytes>, AvroError>> {
+ async move {
+ self.store
+ .get_ranges(&self.path, &ranges)
+ .await
+ .map_err(|e| AvroError::General(e.to_string()))
+ }
+ .boxed()
+ }
+}
diff --git a/arrow-avro/src/lib.rs b/arrow-avro/src/lib.rs
index da451ea146..4314ef96b2 100644
--- a/arrow-avro/src/lib.rs
+++ b/arrow-avro/src/lib.rs
@@ -129,25 +129,22 @@
//! feature is enabled.
//!
//! [`AsyncAvroFileReader`] implements `Stream<Item = Result<RecordBatch,
ArrowError>>`,
-//! allowing efficient async streaming of record batches. When the
`object_store` feature
-//! is enabled, [`AvroObjectReader`] provides integration with object storage
services
-//! such as S3 via the [object_store] crate.
+//! allowing efficient async streaming of record batches. Any
[`AsyncFileReader`]
+//! can be used as the source; there is a built-in implementation for types
+//! implementing `AsyncRead + AsyncSeek` (such as `tokio::fs::File`), and
object
+//! storage services such as S3 can be integrated by implementing
+//! [`AsyncFileReader`] on top of a client such as the [object_store] crate
+//! (see the example on the trait documentation).
//!
//! ```ignore
-//! use std::sync::Arc;
-//! use arrow_avro::reader::{AsyncAvroFileReader, AvroObjectReader};
+//! use arrow_avro::reader::AsyncAvroFileReader;
//! use futures::TryStreamExt;
-//! use object_store::ObjectStore;
-//! use object_store::local::LocalFileSystem;
-//! use object_store::path::Path;
//!
//! # async fn example() -> Result<(), Box<dyn std::error::Error>> {
-//! let store: Arc<dyn ObjectStore> = Arc::new(LocalFileSystem::new());
-//! let path = Path::from("data/example.avro");
-//! let meta = store.head(&path).await?;
+//! let file = tokio::fs::File::open("data/example.avro").await?;
+//! let file_size = file.metadata().await?.len();
//!
-//! let reader = AvroObjectReader::new(store, path);
-//! let stream = AsyncAvroFileReader::builder(reader, meta.size, 1024)
+//! let stream = AsyncAvroFileReader::builder(file, file_size, 1024)
//! .try_build()
//! .await?;
//!
@@ -163,15 +160,15 @@
//! ### Modules
//!
//! - [`reader`]: read Avro (OCF, SOE, Confluent) into Arrow `RecordBatch`es.
-//! - With the `async` feature: [`AsyncAvroFileReader`] for async streaming
reads.
-//! - With the `object_store` feature: [`AvroObjectReader`] for reading from
cloud storage.
+//! - With the `async` feature: [`AsyncAvroFileReader`] for async streaming
reads,
+//! from any [`AsyncFileReader`] source including cloud object storage.
//! - [`writer`]: write Arrow `RecordBatch`es as Avro (OCF, SOE, Confluent,
Apicurio).
//! - [`schema`]: Avro schema parsing / fingerprints / registries.
//! - [`compression`]: codecs used for **OCF block compression** (i.e.,
Deflate, Snappy, Zstandard, BZip2, and XZ).
//! - [`codec`]: internal Avro-Arrow type conversion and row decode/encode
plans.
//!
//! [`AsyncAvroFileReader`]: reader::AsyncAvroFileReader
-//! [`AvroObjectReader`]: reader::AvroObjectReader
+//! [`AsyncFileReader`]: reader::AsyncFileReader
//!
//! ### Features
//!
@@ -182,10 +179,13 @@
//! - `bzip2` — enable BZip2 block compression.
//! - `xz` — enable XZ/LZMA block compression.
//!
-//! **Async & Object Store (opt‑in)**
+//! **Async (opt‑in)**
//! - `async` — enable async APIs for reading Avro (`AsyncAvroFileReader`,
`AsyncFileReader` trait).
-//! - `object_store` — enable integration with the [`object_store`] crate for
reading Avro
-//! from cloud storage (S3, GCS, Azure Blob, etc.) via `AvroObjectReader`.
Implies `async`.
+//! Cloud storage (S3, GCS, Azure Blob, etc.) can be integrated by
implementing
+//! `AsyncFileReader` on top of a client such as the [`object_store`] crate.
+//! - `object_store` (**deprecated**): enables the deprecated
`AvroObjectReader`.
+//! Implement `AsyncFileReader` directly instead (see above). Implies
`async`.
+//! This feature will be removed in a future release.
//!
//! **Schema fingerprints & helpers (opt‑in)**
//! - `md5` — enable MD5 writer‑schema fingerprints.
diff --git a/arrow-avro/src/reader/async_reader/async_file_reader.rs
b/arrow-avro/src/reader/async_reader/async_file_reader.rs
index 1257a2f3dd..8e69d4611c 100644
--- a/arrow-avro/src/reader/async_reader/async_file_reader.rs
+++ b/arrow-avro/src/reader/async_reader/async_file_reader.rs
@@ -30,10 +30,56 @@ use tokio::io::{AsyncRead, AsyncReadExt, AsyncSeek,
AsyncSeekExt};
/// 1. There is a default implementation for types that implement [`AsyncRead`]
/// and [`AsyncSeek`], for example [`tokio::fs::File`].
///
-/// 2. [`super::AvroObjectReader`], available when the `object_store` crate
feature
-/// is enabled, implements this interface for [`ObjectStore`].
+/// 2. Implementations for remote storage, such as the `object_store` crate,
+/// can implement this interface directly, typically by pairing a store
+/// handle with an object path and delegating [`Self::get_bytes`] and
+/// [`Self::get_byte_ranges`] to ranged reads. [`super::SpawnedReader`] can
+/// wrap such a reader to perform its I/O on a dedicated tokio runtime.
///
-/// [`ObjectStore`]: object_store::ObjectStore
+/// # Example: implementing `AsyncFileReader` for the `object_store` crate
+///
+/// ```no_run
+/// # use std::ops::Range;
+/// # use std::sync::Arc;
+/// use arrow_avro::errors::AvroError;
+/// use arrow_avro::reader::AsyncFileReader;
+/// use bytes::Bytes;
+/// use futures::FutureExt;
+/// use futures::future::BoxFuture;
+/// use object_store::path::Path;
+/// use object_store::{ObjectStore, ObjectStoreExt};
+///
+/// #[derive(Clone, Debug)]
+/// struct ObjectStoreReader {
+/// store: Arc<dyn ObjectStore>,
+/// path: Path,
+/// }
+///
+/// impl AsyncFileReader for ObjectStoreReader {
+/// fn get_bytes(&mut self, range: Range<u64>) -> BoxFuture<'_,
Result<Bytes, AvroError>> {
+/// async move {
+/// self.store
+/// .get_range(&self.path, range)
+/// .await
+/// .map_err(|e| AvroError::General(e.to_string()))
+/// }
+/// .boxed()
+/// }
+///
+/// fn get_byte_ranges(
+/// &mut self,
+/// ranges: Vec<Range<u64>>,
+/// ) -> BoxFuture<'_, Result<Vec<Bytes>, AvroError>> {
+/// async move {
+/// self.store
+/// .get_ranges(&self.path, &ranges)
+/// .await
+/// .map_err(|e| AvroError::General(e.to_string()))
+/// }
+/// .boxed()
+/// }
+/// }
+/// ```
///
/// [`tokio::fs::File`]: https://docs.rs/tokio/latest/tokio/fs/struct.File.html
pub trait AsyncFileReader: Send {
diff --git a/arrow-avro/src/reader/async_reader/mod.rs
b/arrow-avro/src/reader/async_reader/mod.rs
index 43ff931a05..0d4603e762 100644
--- a/arrow-avro/src/reader/async_reader/mod.rs
+++ b/arrow-avro/src/reader/async_reader/mod.rs
@@ -35,14 +35,17 @@ use std::task::{Context, Poll};
mod async_file_reader;
mod builder;
+mod spawn;
pub use async_file_reader::AsyncFileReader;
pub use builder::{ReaderBuilder, read_header_info};
+pub use spawn::SpawnedReader;
#[cfg(feature = "object_store")]
mod store;
use crate::errors::AvroError;
+#[allow(deprecated)]
#[cfg(feature = "object_store")]
pub use store::AvroObjectReader;
@@ -543,7 +546,7 @@ impl<R: AsyncFileReader + Unpin + 'static> Stream for
AsyncAvroFileReader<R> {
}
}
-#[cfg(all(test, feature = "object_store"))]
+#[cfg(test)]
mod tests {
use super::*;
use crate::codec::Tz;
@@ -561,6 +564,45 @@ mod tests {
use std::collections::HashMap;
use std::sync::Arc;
+ /// An [`AsyncFileReader`] reading via an [`ObjectStore`], mirroring the
+ /// example on the [`AsyncFileReader`] trait documentation
+ #[derive(Clone, Debug)]
+ struct ObjectStoreReader {
+ store: Arc<dyn ObjectStore>,
+ path: Path,
+ }
+
+ impl ObjectStoreReader {
+ fn new(store: Arc<dyn ObjectStore>, path: Path) -> Self {
+ Self { store, path }
+ }
+ }
+
+ impl AsyncFileReader for ObjectStoreReader {
+ fn get_bytes(&mut self, range: Range<u64>) -> BoxFuture<'_,
Result<Bytes, AvroError>> {
+ async move {
+ self.store
+ .get_range(&self.path, range)
+ .await
+ .map_err(|e| AvroError::General(e.to_string()))
+ }
+ .boxed()
+ }
+
+ fn get_byte_ranges(
+ &mut self,
+ ranges: Vec<Range<u64>>,
+ ) -> BoxFuture<'_, Result<Vec<Bytes>, AvroError>> {
+ async move {
+ self.store
+ .get_ranges(&self.path, &ranges)
+ .await
+ .map_err(|e| AvroError::General(e.to_string()))
+ }
+ .boxed()
+ }
+ }
+
fn arrow_test_data(file: &str) -> String {
let base =
std::env::var("ARROW_TEST_DATA").unwrap_or_else(|_|
"../testing/data".to_string());
@@ -956,7 +998,7 @@ mod tests {
let file_size = store.head(&location).await.unwrap().size;
- let file_reader = AvroObjectReader::new(store, location);
+ let file_reader = ObjectStoreReader::new(store, location);
let mut builder = AsyncAvroFileReader::builder(file_reader, file_size,
batch_size);
if let Some(s) = schema {
@@ -1208,7 +1250,7 @@ mod tests {
let file_size = store.head(&location).await.unwrap().size;
- let file_reader = AvroObjectReader::new(store, location);
+ let file_reader = ObjectStoreReader::new(store, location);
let schema = get_alltypes_schema();
let reader_schema = AvroSchema::try_from(schema.as_ref()).unwrap();
let reader = AsyncAvroFileReader::builder(
@@ -1236,7 +1278,7 @@ mod tests {
let file_size = store.head(&location).await.unwrap().size;
- let file_reader = AvroObjectReader::new(store, location);
+ let file_reader = ObjectStoreReader::new(store, location);
let schema = get_alltypes_schema();
let reader_schema = AvroSchema::try_from(schema.as_ref()).unwrap();
let reader = AsyncAvroFileReader::builder(file_reader, file_size, 1)
@@ -1299,7 +1341,7 @@ mod tests {
let file_size = store.head(&location).await.unwrap().size;
- let mut file_reader = AvroObjectReader::new(store, location);
+ let mut file_reader = ObjectStoreReader::new(store, location);
let header_info = read_header_info(&mut file_reader, file_size, None)
.await
@@ -1392,7 +1434,7 @@ mod tests {
let location = Path::from_filesystem_path(&file_path).unwrap();
let file_size = store.head(&location).await.unwrap().size;
- let file_reader = AvroObjectReader::new(store, location);
+ let file_reader = ObjectStoreReader::new(store, location);
let reader = AsyncAvroFileReader::builder(file_reader, file_size, 2)
.try_build()
.await
@@ -1683,7 +1725,7 @@ mod tests {
let location = Path::from_filesystem_path(&file).unwrap();
let file_size = store.head(&location).await.unwrap().size;
- let file_reader = AvroObjectReader::new(store, location);
+ let file_reader = ObjectStoreReader::new(store, location);
let expected_schema = get_alltypes_schema()
.as_ref()
.clone()
@@ -1710,7 +1752,7 @@ mod tests {
let location = Path::from_filesystem_path(&file).unwrap();
let file_size = store.head(&location).await.unwrap().size;
- let file_reader = AvroObjectReader::new(store, location);
+ let file_reader = ObjectStoreReader::new(store, location);
let schema = get_alltypes_schema()
.project(&[0, 1, 7])
.unwrap()
@@ -1740,7 +1782,7 @@ mod tests {
let location = Path::from_filesystem_path(&file).unwrap();
let file_size = store.head(&location).await.unwrap().size;
- let file_reader = AvroObjectReader::new(store, location);
+ let file_reader = ObjectStoreReader::new(store, location);
// The schema produced by the reader should match the expected schema,
// attaching Avro type name metadata to fields of record and list
types.
@@ -1770,7 +1812,7 @@ mod tests {
let location = Path::from_filesystem_path(&file).unwrap();
let file_size = store.head(&location).await.unwrap().size;
- let file_reader = AvroObjectReader::new(store, location);
+ let file_reader = ObjectStoreReader::new(store, location);
let schema = get_alltypes_schema();
let reader_schema = AvroSchema::try_from(schema.as_ref()).unwrap();
@@ -1797,7 +1839,7 @@ mod tests {
let location = Path::from_filesystem_path(&file).unwrap();
let file_size = store.head(&location).await.unwrap().size;
- let file_reader = AvroObjectReader::new(store, location);
+ let file_reader = ObjectStoreReader::new(store, location);
let schema = get_alltypes_schema();
let reader_schema = AvroSchema::try_from(schema.as_ref()).unwrap();
@@ -1823,7 +1865,7 @@ mod tests {
let location = Path::from_filesystem_path(&file).unwrap();
let file_size = store.head(&location).await.unwrap().size;
- let file_reader = AvroObjectReader::new(store, location);
+ let file_reader = ObjectStoreReader::new(store, location);
let schema = get_alltypes_schema_with_tz("UTC");
let reader_schema = AvroSchema::try_from(schema.as_ref()).unwrap();
@@ -1860,7 +1902,7 @@ mod tests {
let location = Path::from_filesystem_path(&file).unwrap();
let file_size = store.head(&location).await.unwrap().size;
- let file_reader = AvroObjectReader::new(store, location);
+ let file_reader = ObjectStoreReader::new(store, location);
let reader = AsyncAvroFileReader::builder(file_reader, file_size, 1024)
.with_utf8_view(true)
@@ -1891,7 +1933,7 @@ mod tests {
let location = Path::from_filesystem_path(&file).unwrap();
let file_size = store.head(&location).await.unwrap().size;
- let file_reader = AvroObjectReader::new(store, location);
+ let file_reader = ObjectStoreReader::new(store, location);
let reader = AsyncAvroFileReader::builder(file_reader, file_size, 1024)
.with_utf8_view(false)
@@ -1922,7 +1964,7 @@ mod tests {
let location = Path::from_filesystem_path(&file).unwrap();
let file_size = store.head(&location).await.unwrap().size;
- let file_reader = AvroObjectReader::new(store, location);
+ let file_reader = ObjectStoreReader::new(store, location);
// Without strict mode, this should succeed
let reader = AsyncAvroFileReader::builder(file_reader, file_size, 1024)
@@ -1945,7 +1987,7 @@ mod tests {
let location = Path::from_filesystem_path(&file).unwrap();
let file_size = store.head(&location).await.unwrap().size;
- let file_reader = AvroObjectReader::new(store, location);
+ let file_reader = ObjectStoreReader::new(store, location);
// With strict mode, this should fail because of ['T', 'null'] unions
let result = AsyncAvroFileReader::builder(file_reader, file_size, 1024)
@@ -1974,7 +2016,7 @@ mod tests {
let location = Path::from_filesystem_path(&file).unwrap();
let file_size = store.head(&location).await.unwrap().size;
- let file_reader = AvroObjectReader::new(store, location);
+ let file_reader = ObjectStoreReader::new(store, location);
// With strict mode, properly ordered unions should still work
let reader = AsyncAvroFileReader::builder(file_reader, file_size, 1024)
@@ -1996,7 +2038,7 @@ mod tests {
let location = Path::from_filesystem_path(&file).unwrap();
let file_size = store.head(&location).await.unwrap().size;
- let file_reader = AvroObjectReader::new(store, location);
+ let file_reader = ObjectStoreReader::new(store, location);
let reader = AsyncAvroFileReader::builder(file_reader, file_size, 2)
.with_header_size_hint(128)
diff --git a/arrow-avro/src/reader/async_reader/spawn.rs
b/arrow-avro/src/reader/async_reader/spawn.rs
new file mode 100644
index 0000000000..c867e235bf
--- /dev/null
+++ b/arrow-avro/src/reader/async_reader/spawn.rs
@@ -0,0 +1,176 @@
+// 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.
+
+use std::future::Future;
+use std::ops::Range;
+
+use bytes::Bytes;
+use futures::future::BoxFuture;
+use futures::{FutureExt, TryFutureExt};
+use tokio::runtime::Handle;
+
+use crate::errors::AvroError;
+use crate::reader::async_reader::AsyncFileReader;
+
+/// An [`AsyncFileReader`] that performs I/O on a separate tokio runtime.
+///
+/// Tokio is a cooperative scheduler, and relies on tasks yielding in a timely
+/// manner to service IO. Therefore, running IO and CPU-bound tasks, such as
+/// avro decoding, on the same tokio runtime can lead to degraded throughput,
+/// dropped connections and other issues. For more information see [here].
+///
+/// This wrapper spawns each operation of the inner reader onto the provided
+/// runtime [`Handle`], so that the runtime driving the avro decoding does not
+/// also drive the I/O.
+///
+/// The inner reader must be [`Clone`] (typically an `Arc`'d handle to some
+/// shared resource) as each spawned task requires a `'static` copy of it.
+///
+/// [here]:
https://www.influxdata.com/blog/using-rustlangs-async-tokio-runtime-for-cpu-bound-tasks/
+#[derive(Clone, Debug)]
+pub struct SpawnedReader<R> {
+ inner: R,
+ handle: Handle,
+}
+
+impl<R> SpawnedReader<R> {
+ /// Creates a new [`SpawnedReader`] that performs the I/O of `inner` on
`handle`
+ pub fn new(inner: R, handle: Handle) -> Self {
+ Self { inner, handle }
+ }
+
+ /// Returns the inner reader
+ pub fn into_inner(self) -> R {
+ self.inner
+ }
+}
+
+/// Spawns `fut` on `handle`, propagating panics and mapping task cancellation
+/// to [`AvroError::External`]
+fn spawn<T>(
+ handle: &Handle,
+ fut: impl Future<Output = Result<T, AvroError>> + Send + 'static,
+) -> BoxFuture<'static, Result<T, AvroError>>
+where
+ T: Send + 'static,
+{
+ handle
+ .spawn(fut)
+ .map_ok_or_else(
+ |e| match e.try_into_panic() {
+ Err(e) => Err(AvroError::External(Box::new(e))),
+ Ok(p) => std::panic::resume_unwind(p),
+ },
+ |res| res,
+ )
+ .boxed()
+}
+
+impl<R> AsyncFileReader for SpawnedReader<R>
+where
+ R: AsyncFileReader + Clone + Send + 'static,
+{
+ fn get_bytes(&mut self, range: Range<u64>) -> BoxFuture<'_, Result<Bytes,
AvroError>> {
+ let mut inner = self.inner.clone();
+ spawn(&self.handle, async move { inner.get_bytes(range).await })
+ }
+
+ fn get_byte_ranges(
+ &mut self,
+ ranges: Vec<Range<u64>>,
+ ) -> BoxFuture<'_, Result<Vec<Bytes>, AvroError>> {
+ let mut inner = self.inner.clone();
+ spawn(
+ &self.handle,
+ async move { inner.get_byte_ranges(ranges).await },
+ )
+ }
+}
+
+#[cfg(test)]
+mod tests {
+ use super::*;
+ use std::sync::{Arc, Mutex};
+ use std::thread::ThreadId;
+
+ /// An in-memory [`AsyncFileReader`] that records the thread each request
ran on
+ #[derive(Clone)]
+ struct InMemoryReader {
+ data: Bytes,
+ threads: Arc<Mutex<Vec<ThreadId>>>,
+ }
+
+ impl AsyncFileReader for InMemoryReader {
+ fn get_bytes(&mut self, range: Range<u64>) -> BoxFuture<'_,
Result<Bytes, AvroError>> {
+ self.threads
+ .lock()
+ .unwrap()
+ .push(std::thread::current().id());
+ let data = self.data.slice(range.start as usize..range.end as
usize);
+ futures::future::ready(Ok(data)).boxed()
+ }
+ }
+
+ #[tokio::test]
+ async fn test_spawned_reader() {
+ let rt = tokio::runtime::Builder::new_multi_thread()
+ .worker_threads(1)
+ .build()
+ .unwrap();
+
+ let inner = InMemoryReader {
+ data: Bytes::from_static(b"hello world"),
+ threads: Default::default(),
+ };
+ let threads = inner.threads.clone();
+ let mut reader = SpawnedReader::new(inner, rt.handle().clone());
+
+ let bytes = reader.get_bytes(0..5).await.unwrap();
+ assert_eq!(bytes.as_ref(), b"hello");
+
+ let ranges = reader.get_byte_ranges(vec![0..5, 6..11]).await.unwrap();
+ assert_eq!(ranges[1].as_ref(), b"world");
+
+ // All I/O must have run on the spawned runtime, not the current one
+ let current_id = std::thread::current().id();
+ let threads = threads.lock().unwrap();
+ assert!(!threads.is_empty());
+ assert!(threads.iter().all(|id| *id != current_id));
+
+ // Runtimes have to be dropped in blocking contexts
+ tokio::runtime::Handle::current().spawn_blocking(move || drop(rt));
+ }
+
+ #[tokio::test]
+ async fn test_spawned_reader_fails_on_shutdown_runtime() {
+ let rt = tokio::runtime::Builder::new_multi_thread()
+ .worker_threads(1)
+ .build()
+ .unwrap();
+
+ let inner = InMemoryReader {
+ data: Bytes::from_static(b"hello world"),
+ threads: Default::default(),
+ };
+ let mut reader = SpawnedReader::new(inner, rt.handle().clone());
+
+ rt.shutdown_background();
+
+ let err = reader.get_bytes(0..1).await.unwrap_err().to_string();
+ assert!(err.contains("was cancelled"), "{err}");
+ }
+}
diff --git a/arrow-avro/src/reader/async_reader/store.rs
b/arrow-avro/src/reader/async_reader/store.rs
index 44a4abf1a2..9b2c126645 100644
--- a/arrow-avro/src/reader/async_reader/store.rs
+++ b/arrow-avro/src/reader/async_reader/store.rs
@@ -29,12 +29,18 @@ use std::sync::Arc;
use tokio::runtime::Handle;
/// An implementation of an AsyncFileReader using the [`ObjectStore`] API.
+#[deprecated(
+ since = "59.2.0",
+ note = "Implement `AsyncFileReader` directly instead; see the example on
the `AsyncFileReader` trait documentation and
`arrow-avro/examples/object_store.rs`. Use `SpawnedReader` to perform I/O on a
dedicated runtime."
+)]
+#[derive(Clone, Debug)]
pub struct AvroObjectReader {
store: Arc<dyn ObjectStore>,
path: Path,
runtime: Option<Handle>,
}
+#[allow(deprecated)]
impl AvroObjectReader {
/// Creates a new [`Self`] from a store implementation and file location.
pub fn new(store: Arc<dyn ObjectStore>, path: Path) -> Self {
@@ -53,6 +59,10 @@ impl AvroObjectReader {
/// other issues. For more information see [here].
///
/// [here]:
https://www.influxdata.com/blog/using-rustlangs-async-tokio-runtime-for-cpu-bound-tasks/
+ #[deprecated(
+ since = "59.2.0",
+ note = "Wrap the reader in a `SpawnedReader` instead, e.g.
`SpawnedReader::new(reader, handle)`"
+ )]
pub fn with_runtime(self, handle: Handle) -> Self {
Self {
runtime: Some(handle),
@@ -90,6 +100,7 @@ impl AvroObjectReader {
}
}
+#[allow(deprecated)]
impl AsyncFileReader for AvroObjectReader {
fn get_bytes(&mut self, range: Range<u64>) -> BoxFuture<'_, Result<Bytes,
AvroError>> {
self.spawn(|store, path| async move { store.get_range(path,
range).await }.boxed())
diff --git a/arrow-avro/src/reader/mod.rs b/arrow-avro/src/reader/mod.rs
index 6fb75422a4..03730e6045 100644
--- a/arrow-avro/src/reader/mod.rs
+++ b/arrow-avro/src/reader/mod.rs
@@ -504,10 +504,11 @@ pub mod async_reader;
pub use header::{HeaderInfo, read_header_info};
+#[allow(deprecated)]
#[cfg(feature = "object_store")]
pub use async_reader::AvroObjectReader;
#[cfg(feature = "async")]
-pub use async_reader::{AsyncAvroFileReader, AsyncFileReader};
+pub use async_reader::{AsyncAvroFileReader, AsyncFileReader, SpawnedReader};
fn is_incomplete_data(err: &AvroError) -> bool {
matches!(