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 34c7c2b9e0 parquet: Remove explicit `object_store` integration (#10354)
34c7c2b9e0 is described below
commit 34c7c2b9e0cf8ca0e4a595bc9bca8aea8b3e3fa9
Author: Frederic Branczyk <[email protected]>
AuthorDate: Fri Jul 31 14:49:16 2026 +0200
parquet: Remove explicit `object_store` integration (#10354)
Deprecate ParquetObjectReader/ParquetObjectWriter in favor of
implementing AsyncFileReader directly (with an example on the trait
docs) and passing an AsyncWrite such as object_store's BufWriter to
AsyncArrowWriter. To reduce the code needed by implementors I also added
SpawnedReader and ParquetMetaDataReader::with_arrow_reader_options. In
the future we will only need to keep object_store as a dev-dependency to
we can ensure compatibility and use it in benchmarking but remove it
from being a full runtime dependency all together.
Also added a full example in `parquet/examples/object_store.rs`.
# Which issue does this PR close?
- Part of #10308
If we agree on the approach on this, I basically already have the same
thing ready for avro, but wanted to keep the PR smaller so we can align
on the direction before doing it all at once.
# What changes are included in this PR?
Deprecate the object_store dependency and add some extra helpers to
allow users to replace it with very few lines of code (see the example
where it only takes ~60 lines to replace what would take hundreds
without the helpers).
# Are these changes tested?
Yes
# Are there any user-facing changes?
Not yet, just deprecations.
# Additional notes
Full disclaimer: Since we don't use parquet, I don't know this code
super well, so I used Claude Code a bit to help me in creating this
patch. That said, I've fully read the AI-generated parts and understand
what they do, and they are how I would have written it, and how I
interpreted the prior discussion in #10308.
@alamb
---
parquet/Cargo.toml | 12 +-
parquet/benches/arrow_reader_clickbench.rs | 85 ++++++++-
parquet/examples/object_store.rs | 168 ++++++++++++++++++
parquet/src/arrow/arrow_reader/mod.rs | 38 ++++
parquet/src/arrow/arrow_reader/selection/mod.rs | 4 +-
parquet/src/arrow/async_reader/mod.rs | 117 +++++++++---
parquet/src/arrow/async_reader/spawn.rs | 227 ++++++++++++++++++++++++
parquet/src/arrow/async_reader/store.rs | 12 ++
parquet/src/arrow/async_writer/mod.rs | 6 +-
parquet/src/arrow/async_writer/store.rs | 19 +-
parquet/src/lib.rs | 11 +-
parquet/tests/encryption/encryption_async.rs | 106 +++++++++--
12 files changed, 748 insertions(+), 57 deletions(-)
diff --git a/parquet/Cargo.toml b/parquet/Cargo.toml
index 0ae0182263..181f5466aa 100644
--- a/parquet/Cargo.toml
+++ b/parquet/Cargo.toml
@@ -89,7 +89,7 @@ zstd = { version = "0.13", default-features = false }
serde_json = { version = "1.0", features = ["std"], default-features = false }
arrow = { workspace = true, features = ["ipc", "test_utils", "prettyprint",
"json"] }
arrow-cast = { workspace = true }
-tokio = { version = "1.0", default-features = false, features = ["macros",
"rt-multi-thread", "io-util", "fs"] }
+tokio = { version = "1.0", default-features = false, features = ["macros",
"rt-multi-thread", "io-util", "fs", "sync"] }
rand = { version = "0.9", default-features = false, features = ["std",
"std_rng", "thread_rng"] }
object_store = { workspace = true, features = ["azure", "fs"] }
sysinfo = { version = "0.39.6", default-features = false, features =
["system"] }
@@ -116,6 +116,9 @@ experimental = ["variant_experimental"]
# Enable async APIs
async = ["futures", "tokio"]
# Enable object_store integration
+# Deprecated: implement `AsyncFileReader` directly instead, see the example on
+# the `AsyncFileReader` trait documentation and
`parquet/examples/object_store.rs`.
+# This feature will be removed in a future release.
object_store = ["dep:object_store", "async"]
# Group Zstd dependencies
zstd = ["dep:zstd"]
@@ -157,6 +160,11 @@ name = "read_with_rowgroup"
required-features = ["arrow", "async"]
path = "./examples/read_with_rowgroup.rs"
+[[example]]
+name = "object_store"
+required-features = ["arrow", "async"]
+path = "./examples/object_store.rs"
+
[[test]]
name = "arrow_writer_layout"
required-features = ["arrow"]
@@ -254,7 +262,7 @@ harness = false
[[bench]]
name = "arrow_reader_clickbench"
-required-features = ["arrow", "async", "object_store"]
+required-features = ["arrow", "async"]
harness = false
[[bench]]
diff --git a/parquet/benches/arrow_reader_clickbench.rs
b/parquet/benches/arrow_reader_clickbench.rs
index 039829f1b9..f411b9684f 100644
--- a/parquet/benches/arrow_reader_clickbench.rs
+++ b/parquet/benches/arrow_reader_clickbench.rs
@@ -35,16 +35,21 @@ use arrow::compute::{like, nlike, or};
use arrow_array::types::{Int16Type, Int32Type, Int64Type};
use arrow_array::{ArrayRef, ArrowPrimitiveType, BooleanArray, PrimitiveArray,
StringViewArray};
use arrow_schema::{ArrowError, DataType, Schema};
+use bytes::Bytes;
use criterion::{Criterion, criterion_group, criterion_main};
-use futures::StreamExt;
+use futures::future::BoxFuture;
+use futures::{FutureExt, StreamExt, TryFutureExt};
use object_store::local::LocalFileSystem;
+use object_store::{GetOptions, GetRange, ObjectStore, ObjectStoreExt};
use parquet::arrow::arrow_reader::{
ArrowPredicate, ArrowPredicateFn, ArrowReaderMetadata, ArrowReaderOptions,
ParquetRecordBatchReaderBuilder, RowFilter,
};
-use parquet::arrow::async_reader::ParquetObjectReader;
+use parquet::arrow::async_reader::{AsyncFileReader, MetadataSuffixFetch};
use parquet::arrow::{ParquetRecordBatchStreamBuilder, ProjectionMask};
+use parquet::errors::ParquetError;
use parquet::file::metadata::PageIndexPolicy;
+use parquet::file::metadata::{ParquetMetaData, ParquetMetaDataReader};
use parquet::schema::types::SchemaDescriptor;
use std::fmt::{Display, Formatter};
use std::path::{Path, PathBuf};
@@ -101,6 +106,75 @@ criterion_group!(
);
criterion_main!(benches);
+fn to_parquet_err(e: object_store::Error) -> ParquetError {
+ ParquetError::External(Box::new(e))
+}
+
+/// An [`AsyncFileReader`] reading via an [`ObjectStore`], mirroring the
+/// example on the [`AsyncFileReader`] trait documentation
+#[derive(Clone)]
+struct ObjectStoreReader {
+ store: Arc<dyn ObjectStore>,
+ path: object_store::path::Path,
+}
+
+impl AsyncFileReader for ObjectStoreReader {
+ fn get_bytes(
+ &mut self,
+ range: std::ops::Range<u64>,
+ ) -> BoxFuture<'_, Result<Bytes, ParquetError>> {
+ self.store
+ .get_range(&self.path, range)
+ .map_err(to_parquet_err)
+ .boxed()
+ }
+
+ fn get_byte_ranges(
+ &mut self,
+ ranges: Vec<std::ops::Range<u64>>,
+ ) -> BoxFuture<'_, Result<Vec<Bytes>, ParquetError>> {
+ async move {
+ self.store
+ .get_ranges(&self.path, &ranges)
+ .await
+ .map_err(to_parquet_err)
+ }
+ .boxed()
+ }
+
+ fn get_metadata<'a>(
+ &'a mut self,
+ options: Option<&'a ArrowReaderOptions>,
+ ) -> BoxFuture<'a, Result<Arc<ParquetMetaData>, ParquetError>> {
+ async move {
+ let metadata = ParquetMetaDataReader::new()
+ .with_arrow_reader_options(options)
+ .load_via_suffix_and_finish(self)
+ .await?;
+ Ok(Arc::new(metadata))
+ }
+ .boxed()
+ }
+}
+
+impl MetadataSuffixFetch for &mut ObjectStoreReader {
+ fn fetch_suffix(&mut self, suffix: usize) -> BoxFuture<'_, Result<Bytes,
ParquetError>> {
+ let options = GetOptions {
+ range: Some(GetRange::Suffix(suffix as u64)),
+ ..Default::default()
+ };
+ async move {
+ let resp = self
+ .store
+ .get_opts(&self.path, options)
+ .await
+ .map_err(to_parquet_err)?;
+ resp.bytes().await.map_err(to_parquet_err)
+ }
+ .boxed()
+ }
+}
+
/// Predicate Function.
///
/// Functions are invoked with the requested array and return a
[`BooleanArray`]
@@ -736,7 +810,7 @@ impl ReadTest {
self.check_row_count(row_count);
}
- /// Run the filter and projection using the async `ObjectStore` reader
+ /// Like [`Self::run_async`] but reading via an [`ObjectStore`] based
reader
async fn run_async_object_store(&self) {
let hits_path = hits_1();
let parent = hits_path.parent().unwrap();
@@ -744,7 +818,10 @@ impl ReadTest {
let store =
Arc::new(LocalFileSystem::new_with_prefix(parent).unwrap());
let location = object_store::path::Path::from(file_name);
- let reader = ParquetObjectReader::new(store, location);
+ let reader = ObjectStoreReader {
+ store,
+ path: location,
+ };
// setup the reader
let mut stream = ParquetRecordBatchStreamBuilder::new_with_metadata(
diff --git a/parquet/examples/object_store.rs b/parquet/examples/object_store.rs
new file mode 100644
index 0000000000..47c1532733
--- /dev/null
+++ b/parquet/examples/object_store.rs
@@ -0,0 +1,168 @@
+// 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 bytes::Bytes;
+use futures::future::BoxFuture;
+use futures::{FutureExt, TryFutureExt, TryStreamExt};
+use object_store::buffered::BufWriter;
+use object_store::memory::InMemory;
+use object_store::path::Path;
+use object_store::{GetOptions, GetRange, ObjectStore, ObjectStoreExt};
+use parquet::arrow::arrow_reader::ArrowReaderOptions;
+use parquet::arrow::async_reader::{AsyncFileReader, MetadataSuffixFetch,
SpawnedReader};
+use parquet::arrow::{AsyncArrowWriter, ParquetRecordBatchStreamBuilder};
+use parquet::errors::{ParquetError, Result};
+use parquet::file::metadata::{ParquetMetaData, ParquetMetaDataReader};
+use std::ops::Range;
+use std::sync::Arc;
+
+/// This example demonstrates reading and writing Parquet files on object
+/// storage via the [`object_store`] crate, without the deprecated
+/// `ParquetObjectReader` and `ParquetObjectWriter` types.
+///
+/// # Example Overview
+///
+/// 1. Writes a Parquet file to an [`ObjectStore`] by passing an
+/// [`object_store::buffered::BufWriter`] directly to [`AsyncArrowWriter`],
+/// via the blanket [`AsyncFileWriter`] implementation for types
+/// implementing [`AsyncWrite`] (equivalent of `ParquetObjectWriter`).
+///
+/// 2. Reads it back with [`ObjectStoreReader`], a minimal [`AsyncFileReader`]
+/// implementation on top of an [`ObjectStore`] (equivalent of
+/// `ParquetObjectReader`).
+///
+/// 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 Parquet is not also driving the I/O (equivalent of
+/// `ParquetObjectReader::with_runtime`).
+///
+/// [`AsyncFileWriter`]: parquet::arrow::async_writer::AsyncFileWriter
+/// [`AsyncWrite`]: tokio::io::AsyncWrite
+#[tokio::main]
+async fn main() -> Result<()> {
+ let store: Arc<dyn ObjectStore> = Arc::new(InMemory::new());
+ let path = Path::from("example.parquet");
+
+ // 1. Write a Parquet file: a `BufWriter` implements `AsyncWrite` and can
+ // therefore be passed to `AsyncArrowWriter` directly.
+ let col = Arc::new(Int64Array::from_iter_values([1, 2, 3])) as ArrayRef;
+ let batch = RecordBatch::try_from_iter([("col", col)]).unwrap();
+
+ let writer = BufWriter::new(Arc::clone(&store), path.clone());
+ let mut writer = AsyncArrowWriter::try_new(writer, batch.schema(), None)?;
+ writer.write(&batch).await?;
+ writer.close().await?;
+
+ // 2. Read it back with an `AsyncFileReader` implemented on `ObjectStore`.
+ let reader = ObjectStoreReader::new(Arc::clone(&store), path.clone());
+ let builder = ParquetRecordBatchStreamBuilder::new(reader).await?;
+ let read: Vec<RecordBatch> = builder.build()?.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 builder = ParquetRecordBatchStreamBuilder::new(reader).await?;
+ let read: Vec<RecordBatch> = builder.build()?.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(())
+}
+
+fn to_parquet_err(e: object_store::Error) -> ParquetError {
+ ParquetError::External(Box::new(e))
+}
+
+/// 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>>
{
+ self.store
+ .get_range(&self.path, range)
+ .map_err(to_parquet_err)
+ .boxed()
+ }
+
+ fn get_byte_ranges(&mut self, ranges: Vec<Range<u64>>) -> BoxFuture<'_,
Result<Vec<Bytes>>> {
+ async move {
+ self.store
+ .get_ranges(&self.path, &ranges)
+ .await
+ .map_err(to_parquet_err)
+ }
+ .boxed()
+ }
+
+ /// Loads the metadata, respecting the provided [`ArrowReaderOptions`] via
+ /// [`ParquetMetaDataReader::with_arrow_reader_options`]
+ fn get_metadata<'a>(
+ &'a mut self,
+ options: Option<&'a ArrowReaderOptions>,
+ ) -> BoxFuture<'a, Result<Arc<ParquetMetaData>>> {
+ async move {
+ let metadata = ParquetMetaDataReader::new()
+ .with_arrow_reader_options(options)
+ .load_via_suffix_and_finish(self)
+ .await?;
+ Ok(Arc::new(metadata))
+ }
+ .boxed()
+ }
+}
+
+/// Supports fetching the Parquet footer without knowing the file size upfront,
+/// via suffix range requests
+impl MetadataSuffixFetch for &mut ObjectStoreReader {
+ fn fetch_suffix(&mut self, suffix: usize) -> BoxFuture<'_, Result<Bytes>> {
+ let options = GetOptions {
+ range: Some(GetRange::Suffix(suffix as u64)),
+ ..Default::default()
+ };
+ async move {
+ let resp = self
+ .store
+ .get_opts(&self.path, options)
+ .await
+ .map_err(to_parquet_err)?;
+ resp.bytes().await.map_err(to_parquet_err)
+ }
+ .boxed()
+ }
+}
diff --git a/parquet/src/arrow/arrow_reader/mod.rs
b/parquet/src/arrow/arrow_reader/mod.rs
index e4a7b3d135..2517e892fc 100644
--- a/parquet/src/arrow/arrow_reader/mod.rs
+++ b/parquet/src/arrow/arrow_reader/mod.rs
@@ -845,6 +845,44 @@ impl ArrowReaderOptions {
}
}
+impl ParquetMetaDataReader {
+ /// Applies the metadata related settings from [`ArrowReaderOptions`],
+ /// such as the [`ParquetMetaDataOptions`], decryption properties, and
+ /// [`PageIndexPolicy`] to this reader.
+ ///
+ /// The page index policies are only applied if at least one of them is not
+ /// [`PageIndexPolicy::Skip`], so policies previously configured on this
+ /// reader (e.g. from a preload setting) are preserved when the options do
+ /// not request the page index.
+ ///
+ /// This encodes the canonical way to construct a `ParquetMetaDataReader`
+ /// inside `AsyncFileReader::get_metadata` (available with the `async`
+ /// feature), so implementations outside this crate do not need to
+ /// duplicate it.
+ pub fn with_arrow_reader_options(mut self, options:
Option<&ArrowReaderOptions>) -> Self {
+ let Some(options) = options else { return self };
+
+ self =
self.with_metadata_options(Some(options.metadata_options().clone()));
+
+ #[cfg(feature = "encryption")]
+ {
+ self = self.with_decryption_properties(
+ options.file_decryption_properties.as_ref().map(Arc::clone),
+ );
+ }
+
+ if options.column_index_policy() != PageIndexPolicy::Skip
+ || options.offset_index_policy() != PageIndexPolicy::Skip
+ {
+ self = self
+ .with_column_index_policy(options.column_index_policy())
+ .with_offset_index_policy(options.offset_index_policy());
+ }
+
+ self
+ }
+}
+
/// The metadata necessary to construct a [`ArrowReaderBuilder`]
///
/// Note this structure is cheaply clone-able as it consists of several arcs.
diff --git a/parquet/src/arrow/arrow_reader/selection/mod.rs
b/parquet/src/arrow/arrow_reader/selection/mod.rs
index 41a59c048e..ff8d5e9a61 100644
--- a/parquet/src/arrow/arrow_reader/selection/mod.rs
+++ b/parquet/src/arrow/arrow_reader/selection/mod.rs
@@ -363,7 +363,9 @@ impl RowSelection {
///
/// Note: this method does not make any effort to combine consecutive
ranges, nor coalesce
/// ranges that are close together. This is instead delegated to the IO
subsystem to optimise,
- /// e.g. [`ObjectStore::get_ranges`](object_store::ObjectStore::get_ranges)
+ /// e.g. `ObjectStore::get_ranges` in the [`object_store`] crate
+ ///
+ /// [`object_store`]: https://crates.io/crates/object_store
pub fn scan_ranges(&self, page_locations: &[PageLocation]) ->
Vec<Range<u64>> {
match &self.inner {
RowSelectionInner::Selectors(selectors) => {
diff --git a/parquet/src/arrow/async_reader/mod.rs
b/parquet/src/arrow/async_reader/mod.rs
index 5a0083b716..1d276df2a6 100644
--- a/parquet/src/arrow/async_reader/mod.rs
+++ b/parquet/src/arrow/async_reader/mod.rs
@@ -50,11 +50,15 @@ use crate::file::metadata::{ParquetMetaData,
ParquetMetaDataReader};
mod metadata;
pub use metadata::*;
+mod spawn;
+pub use spawn::SpawnedReader;
+
#[cfg(feature = "object_store")]
mod store;
use crate::DecodeResult;
use crate::arrow::push_decoder::{ParquetPushDecoder,
ParquetPushDecoderBuilder, PushDecoderInput};
+#[allow(deprecated)]
#[cfg(feature = "object_store")]
pub use store::*;
@@ -65,10 +69,94 @@ pub use store::*;
/// 1. There is a default implementation for types that implement [`AsyncRead`]
/// and [`AsyncSeek`], for example [`tokio::fs::File`].
///
-/// 2. [`ParquetObjectReader`], 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. [`SpawnedReader`] can wrap
+/// such a reader to perform its I/O on a dedicated runtime, and
+/// [`ParquetMetaDataReader::with_arrow_reader_options`] simplifies
+/// implementing [`Self::get_metadata`].
+///
+/// # Example: implementing `AsyncFileReader` for the `object_store` crate
+///
+/// ```no_run
+/// # use std::ops::Range;
+/// # use std::sync::Arc;
+/// use bytes::Bytes;
+/// use futures::future::BoxFuture;
+/// use futures::{FutureExt, TryFutureExt};
+/// use object_store::path::Path;
+/// use object_store::{GetOptions, GetRange, ObjectStore, ObjectStoreExt};
+/// use parquet::arrow::arrow_reader::ArrowReaderOptions;
+/// use parquet::arrow::async_reader::{AsyncFileReader, MetadataSuffixFetch};
+/// use parquet::errors::{ParquetError, Result};
+/// use parquet::file::metadata::{ParquetMetaData, ParquetMetaDataReader};
+///
+/// fn to_parquet_err(e: object_store::Error) -> ParquetError {
+/// ParquetError::External(Box::new(e))
+/// }
+///
+/// #[derive(Clone)]
+/// struct ObjectStoreReader {
+/// store: Arc<dyn ObjectStore>,
+/// path: Path,
+/// }
+///
+/// impl AsyncFileReader for ObjectStoreReader {
+/// fn get_bytes(&mut self, range: Range<u64>) -> BoxFuture<'_,
Result<Bytes>> {
+/// self.store
+/// .get_range(&self.path, range)
+/// .map_err(to_parquet_err)
+/// .boxed()
+/// }
+///
+/// fn get_byte_ranges(&mut self, ranges: Vec<Range<u64>>) ->
BoxFuture<'_, Result<Vec<Bytes>>> {
+/// async move {
+/// self.store
+/// .get_ranges(&self.path, &ranges)
+/// .await
+/// .map_err(to_parquet_err)
+/// }
+/// .boxed()
+/// }
///
-/// [`ObjectStore`]: object_store::ObjectStore
+/// fn get_metadata<'a>(
+/// &'a mut self,
+/// options: Option<&'a ArrowReaderOptions>,
+/// ) -> BoxFuture<'a, Result<Arc<ParquetMetaData>>> {
+/// async move {
+/// let metadata = ParquetMetaDataReader::new()
+/// .with_arrow_reader_options(options)
+/// .load_via_suffix_and_finish(self)
+/// .await?;
+/// Ok(Arc::new(metadata))
+/// }
+/// .boxed()
+/// }
+/// }
+///
+/// /// Supports fetching the parquet footer without knowing the file size,
+/// /// via suffix range requests
+/// impl MetadataSuffixFetch for &mut ObjectStoreReader {
+/// fn fetch_suffix(&mut self, suffix: usize) -> BoxFuture<'_,
Result<Bytes>> {
+/// let options = GetOptions {
+/// range: Some(GetRange::Suffix(suffix as u64)),
+/// ..Default::default()
+/// };
+/// async move {
+/// let resp = self
+/// .store
+/// .get_opts(&self.path, options)
+/// .await
+/// .map_err(to_parquet_err)?;
+/// resp.bytes().await.map_err(to_parquet_err)
+/// }
+/// .boxed()
+/// }
+/// }
+/// ```
+///
+/// [`ParquetMetaDataReader::with_arrow_reader_options`]:
crate::file::metadata::ParquetMetaDataReader::with_arrow_reader_options
///
/// [`tokio::fs::File`]: https://docs.rs/tokio/latest/tokio/fs/struct.File.html
pub trait AsyncFileReader: Send {
@@ -164,21 +252,7 @@ impl<T: AsyncRead + AsyncSeek + Unpin + Send>
AsyncFileReader for T {
options: Option<&'a ArrowReaderOptions>,
) -> BoxFuture<'a, Result<Arc<ParquetMetaData>>> {
async move {
- let metadata_opts = options.map(|o| o.metadata_options().clone());
- let mut metadata_reader =
-
ParquetMetaDataReader::new().with_metadata_options(metadata_opts);
-
- if let Some(opts) = options {
- metadata_reader = metadata_reader
- .with_column_index_policy(opts.column_index_policy())
- .with_offset_index_policy(opts.offset_index_policy());
- }
-
- #[cfg(feature = "encryption")]
- let metadata_reader = metadata_reader.with_decryption_properties(
- options.and_then(|o|
o.file_decryption_properties.as_ref().map(Arc::clone)),
- );
-
+ let metadata_reader =
ParquetMetaDataReader::new().with_arrow_reader_options(options);
let parquet_metadata =
metadata_reader.load_via_suffix_and_finish(self).await?;
Ok(Arc::new(parquet_metadata))
}
@@ -922,12 +996,7 @@ mod tests {
&'a mut self,
options: Option<&'a ArrowReaderOptions>,
) -> BoxFuture<'a, Result<Arc<ParquetMetaData>>> {
- let mut metadata_reader = ParquetMetaDataReader::new();
- if let Some(opts) = options {
- metadata_reader = metadata_reader
- .with_column_index_policy(opts.column_index_policy())
- .with_offset_index_policy(opts.offset_index_policy());
- }
+ let metadata_reader =
ParquetMetaDataReader::new().with_arrow_reader_options(options);
self.metadata = Some(Arc::new(
metadata_reader.parse_and_finish(&self.data).unwrap(),
));
diff --git a/parquet/src/arrow/async_reader/spawn.rs
b/parquet/src/arrow/async_reader/spawn.rs
new file mode 100644
index 0000000000..927e11e5c5
--- /dev/null
+++ b/parquet/src/arrow/async_reader/spawn.rs
@@ -0,0 +1,227 @@
+// 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 std::sync::Arc;
+
+use bytes::Bytes;
+use futures::future::BoxFuture;
+use futures::{FutureExt, TryFutureExt};
+use tokio::runtime::Handle;
+
+use crate::arrow::arrow_reader::ArrowReaderOptions;
+use crate::arrow::async_reader::{AsyncFileReader, MetadataSuffixFetch};
+use crate::errors::{ParquetError, Result};
+use crate::file::metadata::ParquetMetaData;
+
+/// 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
+/// parquet 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 parquet decoding does
+/// not also drive the I/O.
+///
+/// Note that [`Self::get_metadata`] spawns the entire metadata load, so the
+/// footer is also decoded on the provided runtime.
+///
+/// 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 [`ParquetError::External`]
+fn spawn<T>(
+ handle: &Handle,
+ fut: impl Future<Output = Result<T>> + Send + 'static,
+) -> BoxFuture<'static, Result<T>>
+where
+ T: Send + 'static,
+{
+ handle
+ .spawn(fut)
+ .map_ok_or_else(
+ |e| match e.try_into_panic() {
+ Err(e) => Err(ParquetError::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>>
{
+ 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>>> {
+ let mut inner = self.inner.clone();
+ spawn(
+ &self.handle,
+ async move { inner.get_byte_ranges(ranges).await },
+ )
+ }
+
+ fn get_metadata<'a>(
+ &'a mut self,
+ options: Option<&'a ArrowReaderOptions>,
+ ) -> BoxFuture<'a, Result<Arc<ParquetMetaData>>> {
+ let mut inner = self.inner.clone();
+ let options = options.cloned();
+ spawn(&self.handle, async move {
+ inner.get_metadata(options.as_ref()).await
+ })
+ }
+}
+
+impl<R> MetadataSuffixFetch for &mut SpawnedReader<R>
+where
+ R: AsyncFileReader + Clone + Send + 'static,
+ for<'a> &'a mut R: MetadataSuffixFetch,
+{
+ fn fetch_suffix(&mut self, suffix: usize) -> BoxFuture<'_, Result<Bytes>> {
+ let mut inner = self.inner.clone();
+ spawn(&self.handle, async move {
+ (&mut inner).fetch_suffix(suffix).await
+ })
+ }
+}
+
+#[cfg(test)]
+mod tests {
+ use super::*;
+ use crate::arrow::ParquetRecordBatchStreamBuilder;
+ use crate::file::metadata::ParquetMetaDataReader;
+ use futures::TryStreamExt;
+ use std::thread::ThreadId;
+
+ /// An in-memory [`AsyncFileReader`] that records the thread each request
ran on
+ #[derive(Clone)]
+ struct InMemoryReader {
+ data: Bytes,
+ threads: Arc<std::sync::Mutex<Vec<ThreadId>>>,
+ }
+
+ impl InMemoryReader {
+ fn new(data: Bytes) -> Self {
+ Self {
+ data,
+ threads: Default::default(),
+ }
+ }
+ }
+
+ impl AsyncFileReader for InMemoryReader {
+ fn get_bytes(&mut self, range: Range<u64>) -> BoxFuture<'_,
Result<Bytes>> {
+ 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()
+ }
+
+ fn get_metadata<'a>(
+ &'a mut self,
+ options: Option<&'a ArrowReaderOptions>,
+ ) -> BoxFuture<'a, Result<Arc<ParquetMetaData>>> {
+ self.threads
+ .lock()
+ .unwrap()
+ .push(std::thread::current().id());
+ let metadata = ParquetMetaDataReader::new()
+ .with_arrow_reader_options(options)
+ .parse_and_finish(&self.data);
+ futures::future::ready(metadata.map(Arc::new)).boxed()
+ }
+ }
+
+ #[tokio::test]
+ async fn test_spawned_reader() {
+ let testdata = arrow::util::test_util::parquet_test_data();
+ let path = format!("{testdata}/alltypes_plain.parquet");
+ let data = Bytes::from(std::fs::read(path).unwrap());
+
+ let rt = tokio::runtime::Builder::new_multi_thread()
+ .worker_threads(1)
+ .build()
+ .unwrap();
+
+ let inner = InMemoryReader::new(data);
+ let threads = inner.threads.clone();
+ let reader = SpawnedReader::new(inner, rt.handle().clone());
+
+ let builder =
ParquetRecordBatchStreamBuilder::new(reader).await.unwrap();
+ let batches: Vec<_> =
builder.build().unwrap().try_collect().await.unwrap();
+
+ assert_eq!(batches.len(), 1);
+ assert_eq!(batches[0].num_rows(), 8);
+
+ // 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::new(Bytes::from_static(b"PAR1"));
+ 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/parquet/src/arrow/async_reader/store.rs
b/parquet/src/arrow/async_reader/store.rs
index d47ca744d8..525b39a6ce 100644
--- a/parquet/src/arrow/async_reader/store.rs
+++ b/parquet/src/arrow/async_reader/store.rs
@@ -51,6 +51,10 @@ use tokio::runtime::Handle;
/// print_parquet_metadata(&mut stdout(), builder.metadata());
/// # }
/// ```
+#[deprecated(
+ since = "59.2.0",
+ note = "Implement `AsyncFileReader` directly instead; see the example on
the `AsyncFileReader` trait documentation and
`parquet/examples/object_store.rs`. Use `SpawnedReader` to perform I/O on a
dedicated runtime."
+)]
#[derive(Clone, Debug)]
pub struct ParquetObjectReader {
store: Arc<dyn ObjectStore>,
@@ -62,6 +66,7 @@ pub struct ParquetObjectReader {
runtime: Option<Handle>,
}
+#[allow(deprecated)]
impl ParquetObjectReader {
/// Creates a new [`ParquetObjectReader`] for the provided [`ObjectStore`]
and [`Path`].
pub fn new(store: Arc<dyn ObjectStore>, path: Path) -> Self {
@@ -133,6 +138,10 @@ impl ParquetObjectReader {
/// 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),
@@ -168,6 +177,7 @@ impl ParquetObjectReader {
}
}
+#[allow(deprecated)]
impl MetadataSuffixFetch for &mut ParquetObjectReader {
fn fetch_suffix(&mut self, suffix: usize) -> BoxFuture<'_, Result<Bytes>> {
let options = GetOptions {
@@ -184,6 +194,7 @@ impl MetadataSuffixFetch for &mut ParquetObjectReader {
}
}
+#[allow(deprecated)]
impl AsyncFileReader for ParquetObjectReader {
fn get_bytes(&mut self, range: Range<u64>) -> BoxFuture<'_, Result<Bytes>>
{
self.spawn(|store, path| store.get_range(path, range).boxed())
@@ -247,6 +258,7 @@ impl AsyncFileReader for ParquetObjectReader {
}
#[cfg(test)]
+#[allow(deprecated)]
mod tests {
use crate::arrow::async_reader::ArrowReaderOptions;
use crate::file::metadata::PageIndexPolicy;
diff --git a/parquet/src/arrow/async_writer/mod.rs
b/parquet/src/arrow/async_writer/mod.rs
index a050ef77c4..d9124417d5 100644
--- a/parquet/src/arrow/async_writer/mod.rs
+++ b/parquet/src/arrow/async_writer/mod.rs
@@ -53,10 +53,14 @@
//! # }
//! ```
//!
-//! [`object_store`] provides it's native implementation of
[`AsyncFileWriter`] by [`ParquetObjectWriter`].
+//! There is a blanket implementation of [`AsyncFileWriter`] for all types that
+//! implement [`AsyncWrite`], so writers such as `tokio::fs::File`, or
+//! `object_store::buffered::BufWriter` for writing to object storage, can be
+//! passed to [`AsyncArrowWriter`] directly.
#[cfg(feature = "object_store")]
mod store;
+#[allow(deprecated)]
#[cfg(feature = "object_store")]
pub use store::*;
diff --git a/parquet/src/arrow/async_writer/store.rs
b/parquet/src/arrow/async_writer/store.rs
index 698248e619..fcc6c4b31b 100644
--- a/parquet/src/arrow/async_writer/store.rs
+++ b/parquet/src/arrow/async_writer/store.rs
@@ -28,15 +28,19 @@ use tokio::io::AsyncWriteExt;
/// [`ParquetObjectWriter`] for writing to parquet to [`ObjectStore`]
///
+/// This type is deprecated: [`BufWriter`] implements [`AsyncWrite`] and can
+/// therefore be passed to [`AsyncArrowWriter`] directly via the blanket
+/// [`AsyncFileWriter`] implementation for [`AsyncWrite`] types:
+///
/// ```
/// # use arrow_array::{ArrayRef, Int64Array, RecordBatch};
+/// # use object_store::buffered::BufWriter;
/// # use object_store::memory::InMemory;
/// # use object_store::path::Path;
/// # use object_store::{ObjectStore, ObjectStoreExt};
/// # use std::sync::Arc;
///
/// # use parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder;
-/// # use parquet::arrow::async_writer::ParquetObjectWriter;
/// # use parquet::arrow::AsyncArrowWriter;
///
/// # #[tokio::main(flavor="current_thread")]
@@ -46,7 +50,7 @@ use tokio::io::AsyncWriteExt;
/// let col = Arc::new(Int64Array::from_iter_values([1, 2, 3])) as
ArrayRef;
/// let to_write = RecordBatch::try_from_iter([("col", col)]).unwrap();
///
-/// let object_store_writer = ParquetObjectWriter::new(store.clone(),
Path::from("test"));
+/// let object_store_writer = BufWriter::new(store.clone(),
Path::from("test"));
/// let mut writer =
/// AsyncArrowWriter::try_new(object_store_writer, to_write.schema(),
None).unwrap();
/// writer.write(&to_write).await.unwrap();
@@ -68,11 +72,19 @@ use tokio::io::AsyncWriteExt;
/// assert_eq!(to_write, read);
/// # }
/// ```
+///
+/// [`AsyncWrite`]: tokio::io::AsyncWrite
+/// [`AsyncArrowWriter`]: crate::arrow::async_writer::AsyncArrowWriter
+#[deprecated(
+ since = "59.2.0",
+ note = "Pass an `object_store::buffered::BufWriter` to `AsyncArrowWriter`
directly instead; see `parquet/examples/object_store.rs`"
+)]
#[derive(Debug)]
pub struct ParquetObjectWriter {
w: BufWriter,
}
+#[allow(deprecated)]
impl ParquetObjectWriter {
/// Create a new [`ParquetObjectWriter`] that writes to the specified path
in the given store.
///
@@ -92,6 +104,7 @@ impl ParquetObjectWriter {
}
}
+#[allow(deprecated)]
impl AsyncFileWriter for ParquetObjectWriter {
fn write(&mut self, bs: Bytes) -> BoxFuture<'_, Result<()>> {
Box::pin(async {
@@ -111,12 +124,14 @@ impl AsyncFileWriter for ParquetObjectWriter {
})
}
}
+#[allow(deprecated)]
impl From<BufWriter> for ParquetObjectWriter {
fn from(w: BufWriter) -> Self {
Self::from_buf_writer(w)
}
}
#[cfg(test)]
+#[allow(deprecated)]
mod tests {
use arrow_array::{ArrayRef, Int64Array, RecordBatch};
use object_store::memory::InMemory;
diff --git a/parquet/src/lib.rs b/parquet/src/lib.rs
index 55eebfa370..3acdb61884 100644
--- a/parquet/src/lib.rs
+++ b/parquet/src/lib.rs
@@ -101,15 +101,18 @@
//! read and write [`RecordBatch`]es asynchronously.
//!
//! Most users will use [`AsyncArrowWriter`] for writing and
[`ParquetRecordBatchStreamBuilder`]
-//! for reading. When the `object_store` feature is enabled,
[`ParquetObjectReader`]
-//! provides efficient integration with object storage services such as S3 via
the [object_store]
-//! crate, automatically optimizing IO based on any predicates or projections
provided.
+//! for reading, automatically optimizing IO based on any predicates or
projections provided.
+//! Object storage services such as S3 can be integrated by implementing
+//! [`AsyncFileReader`] on top of a client such as the [object_store] crate,
+//! or by passing a writer implementing [`AsyncWrite`] (such as
+//! `object_store::buffered::BufWriter`) to [`AsyncArrowWriter`].
//!
//! [`async_reader`]: arrow::async_reader
//! [`async_writer`]: arrow::async_writer
//! [`AsyncArrowWriter`]: arrow::async_writer::AsyncArrowWriter
+//! [`AsyncFileReader`]: arrow::async_reader::AsyncFileReader
+//! [`AsyncWrite`]: https://docs.rs/tokio/latest/tokio/io/trait.AsyncWrite.html
//! [`ParquetRecordBatchStreamBuilder`]:
arrow::async_reader::ParquetRecordBatchStreamBuilder
-//! [`ParquetObjectReader`]: arrow::async_reader::ParquetObjectReader
//!
//! ## Variant Logical Type (`variant_experimental` feature)
//!
diff --git a/parquet/tests/encryption/encryption_async.rs
b/parquet/tests/encryption/encryption_async.rs
index f86ab59bf7..a2ec2ed297 100644
--- a/parquet/tests/encryption/encryption_async.rs
+++ b/parquet/tests/encryption/encryption_async.rs
@@ -420,32 +420,100 @@ async fn test_write_non_uniform_encryption() {
.await;
}
-#[cfg(feature = "object_store")]
-async fn get_encrypted_meta_store() -> (
- object_store::ObjectMeta,
- std::sync::Arc<dyn object_store::ObjectStore>,
-) {
- use object_store::local::LocalFileSystem;
+/// An [`AsyncFileReader`] reading via an [`ObjectStore`], mirroring the
+/// example on the [`AsyncFileReader`] trait documentation
+///
+/// [`AsyncFileReader`]: parquet::arrow::async_reader::AsyncFileReader
+/// [`ObjectStore`]: object_store::ObjectStore
+mod object_store_reader {
+ use bytes::Bytes;
+ use futures::future::BoxFuture;
+ use futures::{FutureExt, TryFutureExt};
use object_store::path::Path;
- use object_store::{ObjectStore, ObjectStoreExt};
-
+ use object_store::{GetOptions, GetRange, ObjectStore, ObjectStoreExt};
+ use parquet::arrow::arrow_reader::ArrowReaderOptions;
+ use parquet::arrow::async_reader::{AsyncFileReader, MetadataSuffixFetch};
+ use parquet::errors::{ParquetError, Result};
+ use parquet::file::metadata::{ParquetMetaData, ParquetMetaDataReader};
+ use std::ops::Range;
use std::sync::Arc;
- let test_data = arrow::util::test_util::parquet_test_data();
- let store = LocalFileSystem::new_with_prefix(test_data).unwrap();
- let meta = store
- .head(&Path::from("uniform_encryption.parquet.encrypted"))
- .await
- .unwrap();
+ fn to_parquet_err(e: object_store::Error) -> ParquetError {
+ ParquetError::External(Box::new(e))
+ }
+
+ #[derive(Clone)]
+ pub struct ObjectStoreReader {
+ pub store: Arc<dyn ObjectStore>,
+ pub path: Path,
+ }
+
+ impl AsyncFileReader for ObjectStoreReader {
+ fn get_bytes(&mut self, range: Range<u64>) -> BoxFuture<'_,
Result<Bytes>> {
+ self.store
+ .get_range(&self.path, range)
+ .map_err(to_parquet_err)
+ .boxed()
+ }
+
+ fn get_byte_ranges(
+ &mut self,
+ ranges: Vec<Range<u64>>,
+ ) -> BoxFuture<'_, Result<Vec<Bytes>>> {
+ async move {
+ self.store
+ .get_ranges(&self.path, &ranges)
+ .await
+ .map_err(to_parquet_err)
+ }
+ .boxed()
+ }
- (meta, Arc::new(store) as Arc<dyn ObjectStore>)
+ fn get_metadata<'a>(
+ &'a mut self,
+ options: Option<&'a ArrowReaderOptions>,
+ ) -> BoxFuture<'a, Result<Arc<ParquetMetaData>>> {
+ async move {
+ let metadata = ParquetMetaDataReader::new()
+ .with_arrow_reader_options(options)
+ .load_via_suffix_and_finish(self)
+ .await?;
+ Ok(Arc::new(metadata))
+ }
+ .boxed()
+ }
+ }
+
+ impl MetadataSuffixFetch for &mut ObjectStoreReader {
+ fn fetch_suffix(&mut self, suffix: usize) -> BoxFuture<'_,
Result<Bytes>> {
+ let options = GetOptions {
+ range: Some(GetRange::Suffix(suffix as u64)),
+ ..Default::default()
+ };
+ async move {
+ let resp = self
+ .store
+ .get_opts(&self.path, options)
+ .await
+ .map_err(to_parquet_err)?;
+ resp.bytes().await.map_err(to_parquet_err)
+ }
+ .boxed()
+ }
+ }
}
#[tokio::test]
-#[cfg(feature = "object_store")]
async fn test_read_encrypted_file_from_object_store() {
- use parquet::arrow::async_reader::{AsyncFileReader, ParquetObjectReader};
- let (meta, store) = get_encrypted_meta_store().await;
+ use object_store::local::LocalFileSystem;
+ use object_store::path::Path;
+ use object_store_reader::ObjectStoreReader;
+ use parquet::arrow::async_reader::AsyncFileReader;
+ use std::sync::Arc;
+
+ let test_data = arrow::util::test_util::parquet_test_data();
+ let store = Arc::new(LocalFileSystem::new_with_prefix(test_data).unwrap());
+ let path = Path::from("uniform_encryption.parquet.encrypted");
let key_code: &[u8] = "0123456789012345".as_bytes();
let decryption_properties =
FileDecryptionProperties::builder(key_code.to_vec())
@@ -453,7 +521,7 @@ async fn test_read_encrypted_file_from_object_store() {
.unwrap();
let options =
ArrowReaderOptions::new().with_file_decryption_properties(decryption_properties);
- let mut reader = ParquetObjectReader::new(store,
meta.location).with_file_size(meta.size);
+ let mut reader = ObjectStoreReader { store, path };
let metadata = reader.get_metadata(Some(&options)).await.unwrap();
let builder = ParquetRecordBatchStreamBuilder::new_with_options(reader,
options)
.await