From 101669d89ef46e7e10b4b8f379ac56edf8eae42c Mon Sep 17 00:00:00 2001 From: Andrew Lamb Date: Mon, 20 Jul 2026 15:20:36 -0400 Subject: [PATCH 1/2] chore: Deprecate ParquetObject{Reader,Writer}, add object_store example --- parquet/Cargo.toml | 5 + parquet/benches/arrow_reader_clickbench.rs | 2 + parquet/examples/object_store.rs | 697 +++++++++++++++++++ parquet/src/arrow/async_reader/store.rs | 17 + parquet/src/arrow/async_writer/store.rs | 17 + parquet/tests/encryption/encryption_async.rs | 1 + 6 files changed, 739 insertions(+) create mode 100644 parquet/examples/object_store.rs diff --git a/parquet/Cargo.toml b/parquet/Cargo.toml index 1d44c585210b..54d441c6595a 100644 --- a/parquet/Cargo.toml +++ b/parquet/Cargo.toml @@ -157,6 +157,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"] diff --git a/parquet/benches/arrow_reader_clickbench.rs b/parquet/benches/arrow_reader_clickbench.rs index 039829f1b975..1b70ce60ad93 100644 --- a/parquet/benches/arrow_reader_clickbench.rs +++ b/parquet/benches/arrow_reader_clickbench.rs @@ -42,6 +42,7 @@ use parquet::arrow::arrow_reader::{ ArrowPredicate, ArrowPredicateFn, ArrowReaderMetadata, ArrowReaderOptions, ParquetRecordBatchReaderBuilder, RowFilter, }; +#[allow(deprecated)] use parquet::arrow::async_reader::ParquetObjectReader; use parquet::arrow::{ParquetRecordBatchStreamBuilder, ProjectionMask}; use parquet::file::metadata::PageIndexPolicy; @@ -737,6 +738,7 @@ impl ReadTest { } /// Run the filter and projection using the async `ObjectStore` reader + #[allow(deprecated)] async fn run_async_object_store(&self) { let hits_path = hits_1(); let parent = hits_path.parent().unwrap(); diff --git a/parquet/examples/object_store.rs b/parquet/examples/object_store.rs new file mode 100644 index 000000000000..aa9ea01844f6 --- /dev/null +++ b/parquet/examples/object_store.rs @@ -0,0 +1,697 @@ +// 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. + +//! Example [`ParquetObjectReader`] and [`ParquetObjectWriter`] that +//! read and write Parquet files to object storage using [`ObjectStore`]. +//! +//! These structures used to be a direct part of the `parquet` crate, but will +//! be moved to this example to avoid a dependency on `object_store` and simplify +//! the dependency chain. +//! +//! # Upgrade Guide from Parquet 59 and earlier: +//! +//! To upgrade, simply copy/paste the code from this file into your own crate +//! and update the imports to point to your crate instead of +//! `parquet::arrow::async_reader` and `parquet::arrow::async_writer`. + +use std::{ops::Range, sync::Arc}; + +use bytes::Bytes; +use futures::{FutureExt, TryFutureExt, future::BoxFuture}; +use object_store::ObjectStoreExt; +use object_store::buffered::BufWriter; +use object_store::{GetOptions, GetRange}; +use object_store::{ObjectStore, path::Path}; +use parquet::arrow::arrow_reader::ArrowReaderOptions; +use parquet::arrow::async_reader::{AsyncFileReader, MetadataSuffixFetch}; +use parquet::arrow::async_writer::AsyncFileWriter; +use parquet::errors::{ParquetError, Result}; +use parquet::file::metadata::{PageIndexPolicy, ParquetMetaData, ParquetMetaDataReader}; +use tokio::io::AsyncWriteExt; +use tokio::runtime::Handle; + +/// Reads Parquet files in object storage using [`ObjectStore`]. +/// +/// ```no_run +/// # use std::io::stdout; +/// # use std::sync::Arc; +/// # use object_store::azure::MicrosoftAzureBuilder; +/// # use object_store::{ObjectStore, ObjectStoreExt}; +/// # use object_store::path::Path; +/// # use parquet::arrow::async_reader::ParquetObjectReader; +/// # use parquet::arrow::ParquetRecordBatchStreamBuilder; +/// # use parquet::schema::printer::print_parquet_metadata; +/// # async fn run() { +/// // Populate configuration from environment +/// let storage_container = Arc::new(MicrosoftAzureBuilder::from_env().build().unwrap()); +/// let location = Path::from("path/to/blob.parquet"); +/// let meta = storage_container.head(&location).await.unwrap(); +/// println!("Found Blob with {}B at {}", meta.size, meta.location); +/// +/// // Show Parquet metadata +/// let reader = ParquetObjectReader::new(storage_container, meta.location).with_file_size(meta.size); +/// let builder = ParquetRecordBatchStreamBuilder::new(reader).await.unwrap(); +/// print_parquet_metadata(&mut stdout(), builder.metadata()); +/// # } +/// ``` +#[derive(Clone, Debug)] +pub struct ParquetObjectReader { + store: Arc, + path: Path, + file_size: Option, + metadata_size_hint: Option, + preload_column_index: bool, + preload_offset_index: bool, + runtime: Option, +} + +impl ParquetObjectReader { + /// Creates a new [`ParquetObjectReader`] for the provided [`ObjectStore`] and [`Path`]. + pub fn new(store: Arc, path: Path) -> Self { + Self { + store, + path, + file_size: None, + metadata_size_hint: None, + preload_column_index: false, + preload_offset_index: false, + runtime: None, + } + } + + /// Provide a hint as to the size of the parquet file's footer, + /// see [`ParquetMetaDataReader::with_prefetch_hint`] + pub fn with_footer_size_hint(self, hint: usize) -> Self { + Self { + metadata_size_hint: Some(hint), + ..self + } + } + + /// Provide the byte size of this file. + /// + /// If provided, the file size will ensure that only bounded range requests are used. If file + /// size is not provided, the reader will use suffix range requests to fetch the metadata. + /// + /// Providing this size up front is an important optimization to avoid extra calls when the + /// underlying store does not support suffix range requests. + /// + /// The file size can be obtained using [`ObjectStore::list`] or [`ObjectStoreExt::head`]. + pub fn with_file_size(self, file_size: u64) -> Self { + Self { + file_size: Some(file_size), + ..self + } + } + + /// Whether to load the Column Index as part of [`Self::get_metadata`] + /// + /// Note: This setting may be overridden by [`ArrowReaderOptions`] `page_index_policy`. + /// If `page_index_policy` is `Optional` or `Required`, it will take precedence + /// over this preload flag. When it is `Skip` (default), this flag is used. + pub fn with_preload_column_index(self, preload_column_index: bool) -> Self { + Self { + preload_column_index, + ..self + } + } + + /// Whether to load the Offset Index as part of [`Self::get_metadata`] + /// + /// Note: This setting may be overridden by [`ArrowReaderOptions`] `page_index_policy`. + /// If `page_index_policy` is `Optional` or `Required`, it will take precedence + /// over this preload flag. When it is `Skip` (default), this flag is used. + pub fn with_preload_offset_index(self, preload_offset_index: bool) -> Self { + Self { + preload_offset_index, + ..self + } + } + + /// Perform IO on the provided 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]. + /// + /// [here]: https://www.influxdata.com/blog/using-rustlangs-async-tokio-runtime-for-cpu-bound-tasks/ + pub fn with_runtime(self, handle: Handle) -> Self { + Self { + runtime: Some(handle), + ..self + } + } + + fn spawn(&self, f: F) -> BoxFuture<'_, Result> + where + F: for<'a> FnOnce(&'a Arc, &'a Path) -> BoxFuture<'a, Result> + + Send + + 'static, + O: Send + 'static, + E: Into + Send + 'static, + { + match &self.runtime { + Some(handle) => { + let path = self.path.clone(); + let store = Arc::clone(&self.store); + handle + .spawn(async move { f(&store, &path).await }) + .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.map_err(|e| e.into()), + ) + .boxed() + } + None => f(&self.store, &self.path).map_err(|e| e.into()).boxed(), + } + } +} + +impl MetadataSuffixFetch for &mut ParquetObjectReader { + fn fetch_suffix(&mut self, suffix: usize) -> BoxFuture<'_, Result> { + let options = GetOptions { + range: Some(GetRange::Suffix(suffix as u64)), + ..Default::default() + }; + self.spawn(|store, path| { + async move { + let resp = store + .get_opts(path, options) + .await + .map_err(to_parquet_err)?; + let bytes = resp.bytes().await.map_err(to_parquet_err)?; + Ok::<_, ParquetError>(bytes) + } + .boxed() + }) + } +} + +impl AsyncFileReader for ParquetObjectReader { + fn get_bytes(&mut self, range: Range) -> BoxFuture<'_, Result> { + self.spawn(|store, path| store.get_range(path, range).map_err(to_parquet_err).boxed()) + } + + fn get_byte_ranges(&mut self, ranges: Vec>) -> BoxFuture<'_, Result>> + where + Self: Send, + { + self.spawn(|store, path| { + async move { + store + .get_ranges(path, &ranges) + .await + .map_err(to_parquet_err) + } + .boxed() + }) + } + + // This method doesn't directly call `self.spawn` because all of the IO that is done down the + // line due to this method call is done through `self.get_bytes` and/or `self.get_byte_ranges`. + // When `self` is passed into `ParquetMetaDataReader::load_and_finish`, it treats it as + // an `impl MetadataFetch` and calls those methods to get data from it. Due to `Self`'s impl of + // `AsyncFileReader`, the calls to `MetadataFetch::fetch` are just delegated to + // `Self::get_bytes`. + fn get_metadata<'a>( + &'a mut self, + options: Option<&'a ArrowReaderOptions>, + ) -> BoxFuture<'a, Result>> { + Box::pin(async move { + let metadata_opts = options.map(|o| o.metadata_options().clone()); + let mut metadata = ParquetMetaDataReader::new() + .with_metadata_options(metadata_opts) + .with_column_index_policy(PageIndexPolicy::from(self.preload_column_index)) + .with_offset_index_policy(PageIndexPolicy::from(self.preload_offset_index)) + .with_prefetch_hint(self.metadata_size_hint); + + #[cfg(feature = "encryption")] + if let Some(options) = options { + metadata = metadata.with_decryption_properties( + options.file_decryption_properties().map(Arc::clone), + ); + } + + // Override page index policies from ArrowReaderOptions if specified and not Skip. + // When page_index_policy is Skip (default), use the reader's preload flags. + // When page_index_policy is Optional or Required, override the preload flags + // to ensure the specified policy takes precedence. + if let Some(options) = options { + if options.column_index_policy() != PageIndexPolicy::Skip + || options.offset_index_policy() != PageIndexPolicy::Skip + { + metadata = metadata + .with_column_index_policy(options.column_index_policy()) + .with_offset_index_policy(options.offset_index_policy()); + } + } + + let metadata = if let Some(file_size) = self.file_size { + metadata.load_and_finish(self, file_size).await? + } else { + metadata.load_via_suffix_and_finish(self).await? + }; + + Ok(Arc::new(metadata)) + }) + } +} + +/// [`ParquetObjectWriter`] for writing to parquet to [`ObjectStore`] +/// +/// ``` +/// # use arrow_array::{ArrayRef, Int64Array, RecordBatch}; +/// # 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")] +/// # async fn main() { +/// let store = Arc::new(InMemory::new()); +/// +/// 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 mut writer = +/// AsyncArrowWriter::try_new(object_store_writer, to_write.schema(), None).unwrap(); +/// writer.write(&to_write).await.unwrap(); +/// writer.close().await.unwrap(); +/// +/// let buffer = store +/// .get(&Path::from("test")) +/// .await +/// .unwrap() +/// .bytes() +/// .await +/// .unwrap(); +/// let mut reader = ParquetRecordBatchReaderBuilder::try_new(buffer) +/// .unwrap() +/// .build() +/// .unwrap(); +/// let read = reader.next().unwrap().unwrap(); +/// +/// assert_eq!(to_write, read); +/// # } +/// ``` +#[derive(Debug)] +pub struct ParquetObjectWriter { + w: BufWriter, +} + +impl ParquetObjectWriter { + /// Create a new [`ParquetObjectWriter`] that writes to the specified path in the given store. + /// + /// To configure the writer behavior, please build [`BufWriter`] and then use [`Self::from_buf_writer`] + pub fn new(store: Arc, path: Path) -> Self { + Self::from_buf_writer(BufWriter::new(store, path)) + } + + /// Construct a new ParquetObjectWriter via a existing BufWriter. + pub fn from_buf_writer(w: BufWriter) -> Self { + Self { w } + } + + /// Consume the writer and return the underlying BufWriter. + pub fn into_inner(self) -> BufWriter { + self.w + } +} + +impl AsyncFileWriter for ParquetObjectWriter { + fn write(&mut self, bs: Bytes) -> BoxFuture<'_, Result<()>> { + Box::pin(async { + self.w + .put(bs) + .await + .map_err(|err| ParquetError::External(Box::new(err))) + }) + } + + fn complete(&mut self) -> BoxFuture<'_, Result<()>> { + Box::pin(async { + self.w + .shutdown() + .await + .map_err(|err| ParquetError::External(Box::new(err))) + }) + } +} +impl From for ParquetObjectWriter { + fn from(w: BufWriter) -> Self { + Self::from_buf_writer(w) + } +} + +fn to_parquet_err(e: object_store::Error) -> ParquetError { + ParquetError::External(Box::new(e)) +} + +/// Writes a [`RecordBatch`] to an in-memory [`ObjectStore`] using +/// [`ParquetObjectWriter`] and reads it back using [`ParquetObjectReader`]. +/// +/// [`RecordBatch`]: arrow_array::RecordBatch +#[tokio::main(flavor = "current_thread")] +async fn main() -> Result<()> { + use arrow_array::{ArrayRef, Int64Array, RecordBatch}; + use futures::TryStreamExt; + use object_store::memory::InMemory; + use parquet::arrow::{AsyncArrowWriter, ParquetRecordBatchStreamBuilder}; + + let store = Arc::new(InMemory::new()); + let path = Path::from("example.parquet"); + + // Write a RecordBatch to the object store + let col = Arc::new(Int64Array::from_iter_values([1, 2, 3])) as ArrayRef; + let to_write = RecordBatch::try_from_iter([("col", col)])?; + let writer = ParquetObjectWriter::new(Arc::clone(&store) as _, path.clone()); + let mut writer = AsyncArrowWriter::try_new(writer, to_write.schema(), None)?; + writer.write(&to_write).await?; + writer.close().await?; + + // Read the RecordBatch back from the object store + let reader = ParquetObjectReader::new(store as _, path); + let builder = ParquetRecordBatchStreamBuilder::new(reader).await?; + let batches: Vec<_> = builder.build()?.try_collect().await?; + + assert_eq!(batches, vec![to_write]); + println!( + "Round tripped {} batches through the object store", + batches.len() + ); + Ok(()) +} + +#[cfg(test)] +mod tests { + use super::*; + use arrow::util::test_util::parquet_test_data; + use arrow_array::{ArrayRef, Int64Array, RecordBatch}; + use futures::FutureExt; + use futures::TryStreamExt; + use object_store::local::LocalFileSystem; + use object_store::memory::InMemory; + use object_store::path::Path; + use object_store::{ObjectMeta, ObjectStore, ObjectStoreExt}; + use parquet::arrow::AsyncArrowWriter; + use parquet::arrow::ParquetRecordBatchStreamBuilder; + use parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder; + use parquet::errors::ParquetError; + use parquet::file::metadata::PageIndexPolicy; + use std::sync::Arc; + use std::sync::atomic::{AtomicUsize, Ordering}; + + async fn get_meta_store() -> (ObjectMeta, Arc) { + let res = parquet_test_data(); + let store = LocalFileSystem::new_with_prefix(res).unwrap(); + + let meta = store + .head(&Path::from("alltypes_plain.parquet")) + .await + .unwrap(); + + (meta, Arc::new(store) as Arc) + } + + async fn get_meta_store_with_page_index() -> (ObjectMeta, Arc) { + let res = parquet_test_data(); + let store = LocalFileSystem::new_with_prefix(res).unwrap(); + + let meta = store + .head(&Path::from("alltypes_tiny_pages_plain.parquet")) + .await + .unwrap(); + + (meta, Arc::new(store) as Arc) + } + + #[tokio::test] + async fn test_simple() { + let (meta, store) = get_meta_store().await; + let object_reader = + ParquetObjectReader::new(store, meta.location).with_file_size(meta.size); + + let builder = ParquetRecordBatchStreamBuilder::new(object_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); + } + + #[tokio::test] + async fn test_simple_without_file_length() { + let (meta, store) = get_meta_store().await; + let object_reader = ParquetObjectReader::new(store, meta.location); + + let builder = ParquetRecordBatchStreamBuilder::new(object_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); + } + + #[tokio::test] + async fn test_not_found() { + let (mut meta, store) = get_meta_store().await; + meta.location = Path::from("I don't exist.parquet"); + + let object_reader = + ParquetObjectReader::new(store, meta.location).with_file_size(meta.size); + // Cannot use unwrap_err as ParquetRecordBatchStreamBuilder: !Debug + match ParquetRecordBatchStreamBuilder::new(object_reader).await { + Ok(_) => panic!("expected failure"), + Err(e) => { + let err = e.to_string(); + assert!(err.contains("I don't exist.parquet not found:"), "{err}",); + } + } + } + + #[tokio::test] + async fn test_runtime_is_used() { + let num_actions = Arc::new(AtomicUsize::new(0)); + + let (a1, a2) = (num_actions.clone(), num_actions.clone()); + let rt = tokio::runtime::Builder::new_multi_thread() + .on_thread_park(move || { + a1.fetch_add(1, Ordering::Relaxed); + }) + .on_thread_unpark(move || { + a2.fetch_add(1, Ordering::Relaxed); + }) + .build() + .unwrap(); + + let (meta, store) = get_meta_store().await; + + let initial_actions = num_actions.load(Ordering::Relaxed); + + let reader = ParquetObjectReader::new(store, meta.location) + .with_file_size(meta.size) + .with_runtime(rt.handle().clone()); + + let builder = ParquetRecordBatchStreamBuilder::new(reader).await.unwrap(); + let batches: Vec<_> = builder.build().unwrap().try_collect().await.unwrap(); + + // Just copied these assert_eqs from the `test_simple` above + assert_eq!(batches.len(), 1); + assert_eq!(batches[0].num_rows(), 8); + + assert!(num_actions.load(Ordering::Relaxed) - initial_actions > 0); + + // Runtimes have to be dropped in blocking contexts, so we need to move this one to a new + // blocking thread to drop it. + tokio::runtime::Handle::current().spawn_blocking(move || drop(rt)); + } + + /// Unit test that `ParquetObjectReader::spawn`spawns on the provided runtime + #[tokio::test] + async fn test_runtime_thread_id_different() { + let rt = tokio::runtime::Builder::new_multi_thread() + .worker_threads(1) + .build() + .unwrap(); + + let (meta, store) = get_meta_store().await; + + let reader = ParquetObjectReader::new(store, meta.location) + .with_file_size(meta.size) + .with_runtime(rt.handle().clone()); + + let current_id = std::thread::current().id(); + + let other_id = reader + .spawn(|_, _| async move { Ok::<_, ParquetError>(std::thread::current().id()) }.boxed()) + .await + .unwrap(); + + assert_ne!(current_id, other_id); + + tokio::runtime::Handle::current().spawn_blocking(move || drop(rt)); + } + + #[tokio::test] + async fn io_fails_on_shutdown_runtime() { + let rt = tokio::runtime::Builder::new_multi_thread() + .worker_threads(1) + .build() + .unwrap(); + + let (meta, store) = get_meta_store().await; + + let mut reader = ParquetObjectReader::new(store, meta.location) + .with_file_size(meta.size) + .with_runtime(rt.handle().clone()); + + rt.shutdown_background(); + + let err = reader.get_bytes(0..1).await.unwrap_err().to_string(); + + assert!(err.to_string().contains("was cancelled")); + } + + #[tokio::test] + async fn test_page_index_policy_skip_uses_preload_true() { + let (meta, store) = get_meta_store_with_page_index().await; + + // Create reader with preload flags set to true + let mut reader = ParquetObjectReader::new(store.clone(), meta.location.clone()) + .with_file_size(meta.size) + .with_preload_column_index(true) + .with_preload_offset_index(true); + + // Create options with page_index_policy set to Skip (default) + let options = ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Skip); + + // Get metadata - Skip means use reader's preload flags (true) + let metadata = reader.get_metadata(Some(&options)).await.unwrap(); + + // With preload=true, indexes should be loaded since the test file has them + assert!(metadata.column_index().is_some()); + } + + #[tokio::test] + async fn test_page_index_policy_optional_overrides_preload_false() { + let (meta, store) = get_meta_store_with_page_index().await; + + // Create reader with preload flags set to false + let mut reader = ParquetObjectReader::new(store.clone(), meta.location.clone()) + .with_file_size(meta.size) + .with_preload_column_index(false) + .with_preload_offset_index(false); + + // Create options with page_index_policy set to Optional + let options = ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Optional); + + // Get metadata - Optional overrides preload flags and attempts to load indexes + let metadata = reader.get_metadata(Some(&options)).await.unwrap(); + + // With Optional policy, it will TRY to load indexes but won't fail if they don't exist + // The test file has page indexes, so they will be some + assert!(metadata.column_index().is_some()); + } + + #[tokio::test] + async fn test_page_index_policy_optional_vs_skip() { + let (meta, store) = get_meta_store_with_page_index().await; + + // Test 1: preload=false + Skip policy -> uses preload flags (false) + let mut reader1 = ParquetObjectReader::new(store.clone(), meta.location.clone()) + .with_file_size(meta.size) + .with_preload_column_index(false) + .with_preload_offset_index(false); + + let options1 = ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Skip); + let metadata1 = reader1.get_metadata(Some(&options1)).await.unwrap(); + + // Test 2: preload=false + Optional policy -> overrides to try loading + let mut reader2 = ParquetObjectReader::new(store.clone(), meta.location.clone()) + .with_file_size(meta.size) + .with_preload_column_index(false) + .with_preload_offset_index(false); + + let options2 = ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Optional); + let metadata2 = reader2.get_metadata(Some(&options2)).await.unwrap(); + + // Both should succeed (no panic/error) + // metadata1 (Skip) uses preload=false -> Skip policy + // metadata2 (Optional) overrides preload=false -> Optional policy + assert!(metadata1.column_index().is_none()); + assert!(metadata2.column_index().is_some()); + } + + #[tokio::test] + async fn test_page_index_policy_no_options_uses_preload() { + let (meta, store) = get_meta_store_with_page_index().await; + + // Create reader with preload flags set to true + let mut reader = ParquetObjectReader::new(store, meta.location) + .with_file_size(meta.size) + .with_preload_column_index(true) + .with_preload_offset_index(true); + + // Get metadata without options - should use reader's preload flags + let metadata = reader.get_metadata(None).await.unwrap(); + + // With no options provided, preload flags (true) should be respected + // and converted to Optional policy internally (preload=true -> Optional) + // The test file has page indexes, so they will be some + assert!(metadata.column_index().is_some() && metadata.column_index().is_some()); + } + + #[tokio::test] + async fn test_async_writer() { + let store = Arc::new(InMemory::new()); + + 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 mut writer = + AsyncArrowWriter::try_new(object_store_writer, to_write.schema(), None).unwrap(); + writer.write(&to_write).await.unwrap(); + writer.close().await.unwrap(); + + let buffer = store + .get(&Path::from("test")) + .await + .unwrap() + .bytes() + .await + .unwrap(); + let mut reader = ParquetRecordBatchReaderBuilder::try_new(buffer) + .unwrap() + .build() + .unwrap(); + let read = reader.next().unwrap().unwrap(); + + assert_eq!(to_write, read); + } +} diff --git a/parquet/src/arrow/async_reader/store.rs b/parquet/src/arrow/async_reader/store.rs index d47ca744d8f6..196271b3c516 100644 --- a/parquet/src/arrow/async_reader/store.rs +++ b/parquet/src/arrow/async_reader/store.rs @@ -29,7 +29,16 @@ use object_store::{ObjectStore, path::Path}; use tokio::runtime::Handle; /// Reads Parquet files in object storage using [`ObjectStore`]. /// +/// This structure is deprecated and will be removed in a future release, in +/// order to remove the `parquet` crate's dependency on `object_store`. To +/// upgrade, copy the implementation from the [`object_store` example] into +/// your own codebase. See [#10308] for more details. +/// +/// [`object_store` example]: https://github.com/apache/arrow-rs/blob/main/parquet/examples/object_store.rs +/// [#10308]: https://github.com/apache/arrow-rs/issues/10308 +/// /// ```no_run +/// # #![allow(deprecated)] /// # use std::io::stdout; /// # use std::sync::Arc; /// # use object_store::azure::MicrosoftAzureBuilder; @@ -51,6 +60,10 @@ use tokio::runtime::Handle; /// print_parquet_metadata(&mut stdout(), builder.metadata()); /// # } /// ``` +#[deprecated( + since = "59.4.0", + note = "copy the implementation from https://github.com/apache/arrow-rs/blob/main/parquet/examples/object_store.rs into your own codebase" +)] #[derive(Clone, Debug)] pub struct ParquetObjectReader { store: Arc, @@ -62,6 +75,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 { @@ -168,6 +182,7 @@ impl ParquetObjectReader { } } +#[allow(deprecated)] impl MetadataSuffixFetch for &mut ParquetObjectReader { fn fetch_suffix(&mut self, suffix: usize) -> BoxFuture<'_, Result> { let options = GetOptions { @@ -184,6 +199,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 +263,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/store.rs b/parquet/src/arrow/async_writer/store.rs index 698248e619db..72692a14d757 100644 --- a/parquet/src/arrow/async_writer/store.rs +++ b/parquet/src/arrow/async_writer/store.rs @@ -28,7 +28,16 @@ use tokio::io::AsyncWriteExt; /// [`ParquetObjectWriter`] for writing to parquet to [`ObjectStore`] /// +/// This structure is deprecated and will be removed in a future release, in +/// order to remove the `parquet` crate's dependency on `object_store`. To +/// upgrade, copy the implementation from the [`object_store` example] into +/// your own codebase. See [#10308] for more details. +/// +/// [`object_store` example]: https://github.com/apache/arrow-rs/blob/main/parquet/examples/object_store.rs +/// [#10308]: https://github.com/apache/arrow-rs/issues/10308/// [`object_store` example]: https://github.com/apache/arrow-rs/blob/main/parquet/examples/object_store.rs +/// /// ``` +/// # #![allow(deprecated)] /// # use arrow_array::{ArrayRef, Int64Array, RecordBatch}; /// # use object_store::memory::InMemory; /// # use object_store::path::Path; @@ -68,11 +77,16 @@ use tokio::io::AsyncWriteExt; /// assert_eq!(to_write, read); /// # } /// ``` +#[deprecated( + since = "59.4.0", + note = "copy the implementation from https://github.com/apache/arrow-rs/blob/main/parquet/examples/object_store.rs into your own codebase" +)] #[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 +106,7 @@ impl ParquetObjectWriter { } } +#[allow(deprecated)] impl AsyncFileWriter for ParquetObjectWriter { fn write(&mut self, bs: Bytes) -> BoxFuture<'_, Result<()>> { Box::pin(async { @@ -111,12 +126,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/tests/encryption/encryption_async.rs b/parquet/tests/encryption/encryption_async.rs index f86ab59bf755..0666f1c74daa 100644 --- a/parquet/tests/encryption/encryption_async.rs +++ b/parquet/tests/encryption/encryption_async.rs @@ -443,6 +443,7 @@ async fn get_encrypted_meta_store() -> ( #[tokio::test] #[cfg(feature = "object_store")] +#[allow(deprecated)] async fn test_read_encrypted_file_from_object_store() { use parquet::arrow::async_reader::{AsyncFileReader, ParquetObjectReader}; let (meta, store) = get_encrypted_meta_store().await; From fd3fcafae1e02af1100ae591370eb57c74a9ca6f Mon Sep 17 00:00:00 2001 From: Andrew Lamb Date: Mon, 20 Jul 2026 16:01:12 -0400 Subject: [PATCH 2/2] Fix malformed doc link in ParquetObjectWriter deprecation docs Co-Authored-By: Claude Fable 5 --- parquet/src/arrow/async_writer/store.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/parquet/src/arrow/async_writer/store.rs b/parquet/src/arrow/async_writer/store.rs index 72692a14d757..df07fad5d08b 100644 --- a/parquet/src/arrow/async_writer/store.rs +++ b/parquet/src/arrow/async_writer/store.rs @@ -34,7 +34,7 @@ use tokio::io::AsyncWriteExt; /// your own codebase. See [#10308] for more details. /// /// [`object_store` example]: https://github.com/apache/arrow-rs/blob/main/parquet/examples/object_store.rs -/// [#10308]: https://github.com/apache/arrow-rs/issues/10308/// [`object_store` example]: https://github.com/apache/arrow-rs/blob/main/parquet/examples/object_store.rs +/// [#10308]: https://github.com/apache/arrow-rs/issues/10308 /// /// ``` /// # #![allow(deprecated)]