From 160a16d46e3914dc42a3dcee54769f35a114ab88 Mon Sep 17 00:00:00 2001 From: Frederic Branczyk Date: Thu, 16 Jul 2026 14:29:47 +0200 Subject: [PATCH 1/2] parquet: Deprecate object_store integration 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`. --- 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.rs | 2 +- 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, 746 insertions(+), 57 deletions(-) create mode 100644 parquet/examples/object_store.rs create mode 100644 parquet/src/arrow/async_reader/spawn.rs diff --git a/parquet/Cargo.toml b/parquet/Cargo.toml index 1d44c585210b..561e6ff48e20 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"] @@ -249,7 +257,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 039829f1b975..f411b9684fa5 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, + path: object_store::path::Path, +} + +impl AsyncFileReader for ObjectStoreReader { + fn get_bytes( + &mut self, + range: std::ops::Range, + ) -> BoxFuture<'_, Result> { + self.store + .get_range(&self.path, range) + .map_err(to_parquet_err) + .boxed() + } + + fn get_byte_ranges( + &mut self, + ranges: Vec>, + ) -> BoxFuture<'_, Result, 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, 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> { + 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 000000000000..65c59cac31ad --- /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`] (replaces `ParquetObjectWriter`). +/// +/// 2. Reads it back with [`ObjectStoreReader`], a minimal [`AsyncFileReader`] +/// implementation on top of an [`ObjectStore`] (replaces +/// `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 (replaces +/// `ParquetObjectReader::with_runtime`). +/// +/// [`AsyncFileWriter`]: parquet::arrow::async_writer::AsyncFileWriter +/// [`AsyncWrite`]: tokio::io::AsyncWrite +#[tokio::main] +async fn main() -> Result<()> { + let store: Arc = 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 = 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 = 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, + path: Path, +} + +impl ObjectStoreReader { + fn new(store: Arc, path: Path) -> Self { + Self { store, path } + } +} + +impl AsyncFileReader for ObjectStoreReader { + fn get_bytes(&mut self, range: Range) -> BoxFuture<'_, Result> { + self.store + .get_range(&self.path, range) + .map_err(to_parquet_err) + .boxed() + } + + fn get_byte_ranges(&mut self, ranges: Vec>) -> BoxFuture<'_, Result>> { + 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>> { + 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> { + 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 814b8250508e..3a9844ca6482 100644 --- a/parquet/src/arrow/arrow_reader/mod.rs +++ b/parquet/src/arrow/arrow_reader/mod.rs @@ -841,6 +841,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.rs b/parquet/src/arrow/arrow_reader/selection.rs index 2ddf812f9c39..f6a9dd853e03 100644 --- a/parquet/src/arrow/arrow_reader/selection.rs +++ b/parquet/src/arrow/arrow_reader/selection.rs @@ -200,7 +200,7 @@ 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 pub fn scan_ranges(&self, page_locations: &[PageLocation]) -> Vec> { let mut ranges: Vec> = vec![]; let mut row_offset = 0; diff --git a/parquet/src/arrow/async_reader/mod.rs b/parquet/src/arrow/async_reader/mod.rs index 5a0083b7164d..1d276df2a611 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, +/// path: Path, +/// } +/// +/// impl AsyncFileReader for ObjectStoreReader { +/// fn get_bytes(&mut self, range: Range) -> BoxFuture<'_, Result> { +/// self.store +/// .get_range(&self.path, range) +/// .map_err(to_parquet_err) +/// .boxed() +/// } +/// +/// fn get_byte_ranges(&mut self, ranges: Vec>) -> BoxFuture<'_, Result>> { +/// 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>> { +/// 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> { +/// 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 AsyncFileReader for T { options: Option<&'a ArrowReaderOptions>, ) -> BoxFuture<'a, Result>> { 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>> { - 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 000000000000..927e11e5c517 --- /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 { + inner: R, + handle: Handle, +} + +impl SpawnedReader { + /// 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( + handle: &Handle, + fut: impl Future> + Send + 'static, +) -> BoxFuture<'static, Result> +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 AsyncFileReader for SpawnedReader +where + R: AsyncFileReader + Clone + Send + 'static, +{ + fn get_bytes(&mut self, range: Range) -> BoxFuture<'_, Result> { + let mut inner = self.inner.clone(); + spawn(&self.handle, async move { inner.get_bytes(range).await }) + } + + fn get_byte_ranges(&mut self, ranges: Vec>) -> BoxFuture<'_, Result>> { + 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>> { + let mut inner = self.inner.clone(); + let options = options.cloned(); + spawn(&self.handle, async move { + inner.get_metadata(options.as_ref()).await + }) + } +} + +impl MetadataSuffixFetch for &mut SpawnedReader +where + R: AsyncFileReader + Clone + Send + 'static, + for<'a> &'a mut R: MetadataSuffixFetch, +{ + fn fetch_suffix(&mut self, suffix: usize) -> BoxFuture<'_, Result> { + 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>>, + } + + impl InMemoryReader { + fn new(data: Bytes) -> Self { + Self { + data, + threads: Default::default(), + } + } + } + + impl AsyncFileReader for InMemoryReader { + fn get_bytes(&mut self, range: Range) -> BoxFuture<'_, Result> { + 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>> { + 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 d47ca744d8f6..525b39a6ce23 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, @@ -62,6 +66,7 @@ pub struct ParquetObjectReader { runtime: Option, } +#[allow(deprecated)] impl ParquetObjectReader { /// Creates a new [`ParquetObjectReader`] for the provided [`ObjectStore`] and [`Path`]. pub fn new(store: Arc, 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> { 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) -> BoxFuture<'_, Result> { 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 a050ef77c46c..d9124417d5e5 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 698248e619db..fcc6c4b31b07 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 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 55eebfa3708d..3acdb61884c6 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 f86ab59bf755..a2ec2ed29768 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, -) { - 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, + pub path: Path, + } + + impl AsyncFileReader for ObjectStoreReader { + fn get_bytes(&mut self, range: Range) -> BoxFuture<'_, Result> { + self.store + .get_range(&self.path, range) + .map_err(to_parquet_err) + .boxed() + } + + fn get_byte_ranges( + &mut self, + ranges: Vec>, + ) -> BoxFuture<'_, Result>> { + async move { + self.store + .get_ranges(&self.path, &ranges) + .await + .map_err(to_parquet_err) + } + .boxed() + } - (meta, Arc::new(store) as Arc) + fn get_metadata<'a>( + &'a mut self, + options: Option<&'a ArrowReaderOptions>, + ) -> BoxFuture<'a, Result>> { + 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> { + 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 From 9e4035a94e943605ada40d9bd6b4ebef23742292 Mon Sep 17 00:00:00 2001 From: Frederic Branczyk Date: Tue, 28 Jul 2026 16:56:40 +0200 Subject: [PATCH 2/2] parquet: Adjust comments according to PR review --- parquet/examples/object_store.rs | 6 +++--- parquet/src/arrow/arrow_reader/selection.rs | 4 +++- 2 files changed, 6 insertions(+), 4 deletions(-) diff --git a/parquet/examples/object_store.rs b/parquet/examples/object_store.rs index 65c59cac31ad..47c1532733dd 100644 --- a/parquet/examples/object_store.rs +++ b/parquet/examples/object_store.rs @@ -40,15 +40,15 @@ use std::sync::Arc; /// 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`] (replaces `ParquetObjectWriter`). +/// implementing [`AsyncWrite`] (equivalent of `ParquetObjectWriter`). /// /// 2. Reads it back with [`ObjectStoreReader`], a minimal [`AsyncFileReader`] -/// implementation on top of an [`ObjectStore`] (replaces +/// 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 (replaces +/// decoding Parquet is not also driving the I/O (equivalent of /// `ParquetObjectReader::with_runtime`). /// /// [`AsyncFileWriter`]: parquet::arrow::async_writer::AsyncFileWriter diff --git a/parquet/src/arrow/arrow_reader/selection.rs b/parquet/src/arrow/arrow_reader/selection.rs index f6a9dd853e03..2820398a02b7 100644 --- a/parquet/src/arrow/arrow_reader/selection.rs +++ b/parquet/src/arrow/arrow_reader/selection.rs @@ -200,7 +200,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` in the `object_store` crate + /// 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> { let mut ranges: Vec> = vec![]; let mut row_offset = 0;