From 292d245da3337c9bd705db60e14dd34c9ca3f65b Mon Sep 17 00:00:00 2001 From: xiaohongbo Date: Mon, 28 Sep 2026 15:45:20 -0700 Subject: [PATCH 1/2] [core] Reuse whole-file Parquet metadata prefetch --- crates/paimon/src/arrow/format/parquet.rs | 147 ++++++++++++++++++++-- 1 file changed, 138 insertions(+), 9 deletions(-) diff --git a/crates/paimon/src/arrow/format/parquet.rs b/crates/paimon/src/arrow/format/parquet.rs index 5ce567bfd..863a23c98 100644 --- a/crates/paimon/src/arrow/format/parquet.rs +++ b/crates/paimon/src/arrow/format/parquet.rs @@ -55,7 +55,7 @@ use parquet::file::statistics::Statistics as ParquetStatistics; use std::cmp::Ordering; use std::collections::HashMap; use std::ops::Range; -use std::sync::{Arc, Mutex}; +use std::sync::{Arc, Mutex, OnceLock}; use tokio::sync::{mpsc, oneshot}; pub(crate) struct ParquetFormatReader { @@ -611,7 +611,7 @@ impl FormatFileReader for ParquetFormatReader { batch_size: Option, row_selection: Option>, ) -> crate::Result { - let shared_reader: Arc = reader.into(); + let shared_reader = shared_parquet_file_reader(reader, file_size); let arrow_file_reader = ArrowFileReader::new(file_size, Arc::clone(&shared_reader)) .with_metadata_cache_enabled(self.metadata_cache_enabled); @@ -2560,6 +2560,85 @@ const METADATA_SIZE_HINT: usize = 512 * 1024; /// avoid excessive small IO requests whose per-request overhead dominates. const IO_BLOCK_SIZE: u64 = 4 * 1024 * 1024; +/// Reuses a whole-file metadata prefetch for subsequent data ranges. +/// +/// `ParquetMetaDataReader` reads the entire object when it is no larger than +/// [`METADATA_SIZE_HINT`]. Without this adapter, the decoded stream then asks +/// the object store for the selected data range again. Keep the already-read +/// bytes for the lifetime of this file reader and serve those ranges by slicing +/// the same [`Bytes`] allocation. +struct WholeFileCachingRead { + inner: Arc, + file_size: u64, + whole_file: OnceLock, +} + +impl WholeFileCachingRead { + fn new(inner: Arc, file_size: u64) -> Self { + Self { + inner, + file_size, + whole_file: OnceLock::new(), + } + } +} + +#[async_trait] +impl FileRead for WholeFileCachingRead { + async fn read(&self, range: Range) -> crate::Result { + if let Some(whole_file) = self.whole_file.get() { + let start = usize::try_from(range.start).map_err(|error| Error::DataInvalid { + message: format!("Parquet range start does not fit usize: {}", range.start), + source: Some(Box::new(error)), + })?; + let end = usize::try_from(range.end).map_err(|error| Error::DataInvalid { + message: format!("Parquet range end does not fit usize: {}", range.end), + source: Some(Box::new(error)), + })?; + if start > end || end > whole_file.len() { + return Err(Error::DataInvalid { + message: format!( + "Parquet range {}..{} exceeds file size {}", + range.start, range.end, self.file_size + ), + source: None, + }); + } + return Ok(whole_file.slice(start..end)); + } + + let bytes = self.inner.read(range.clone()).await?; + if range.start == 0 + && range.end == self.file_size + && u64::try_from(bytes.len()).ok() == Some(self.file_size) + { + let _ = self.whole_file.set(bytes.clone()); + } + Ok(bytes) + } + + fn cache_namespace(&self) -> Option { + self.inner.cache_namespace() + } + + fn cache_key(&self) -> Option<&str> { + self.inner.cache_key() + } + + fn file_format_metadata_cache(&self) -> Option<&(dyn std::any::Any + Send + Sync)> { + self.inner.file_format_metadata_cache() + } +} + +fn shared_parquet_file_reader(reader: Box, file_size: u64) -> Arc { + let reader: Arc = reader.into(); + if file_size <= METADATA_SIZE_HINT as u64 { + Arc::new(WholeFileCachingRead::new(reader, file_size)) + } else { + reader + } +} + /// Rows in a Parquet file, read from its footer alone. pub(crate) async fn read_row_count( reader: Box, @@ -3418,10 +3497,13 @@ mod tests { #[tokio::test] async fn test_parquet_reader_reads_row_groups_concurrently_in_order() { - const ROWS: i32 = 512; + // Keep this file above the whole-file metadata prefetch threshold so + // the test continues to exercise remote row-group concurrency. + const ROWS: i32 = 131_072; let schema = writer_arrow_schema(); let props = parquet::file::properties::WriterProperties::builder() - .set_max_row_group_row_count(Some(64)) + .set_max_row_group_row_count(Some(16_384)) + .set_compression(Compression::UNCOMPRESSED) .set_dictionary_enabled(false) .build(); let mut data = Vec::new(); @@ -3445,6 +3527,7 @@ mod tests { max_in_flight: Arc::clone(&max_in_flight), }; let file_size = file_reader.data.len() as u64; + assert!(file_size > super::METADATA_SIZE_HINT as u64); let fields = vec![DataField::new( 0, "id".to_string(), @@ -3458,7 +3541,7 @@ mod tests { file_size, &fields, None, - Some(32), + Some(8192), None, ) .await @@ -3597,10 +3680,13 @@ mod tests { #[tokio::test] async fn test_parquet_read_budget_is_shared_across_readers() { - const ROWS: i32 = 256; + // Keep this file above the whole-file metadata prefetch threshold so + // the test continues to observe shared remote-read concurrency. + const ROWS: i32 = 65_536; let schema = writer_arrow_schema(); let props = parquet::file::properties::WriterProperties::builder() - .set_max_row_group_row_count(Some(32)) + .set_max_row_group_row_count(Some(8192)) + .set_compression(Compression::UNCOMPRESSED) .set_dictionary_enabled(false) .build(); let mut data = Vec::new(); @@ -3625,6 +3711,7 @@ mod tests { max_in_flight: Arc::clone(&max_in_flight), }; let file_size = data.len() as u64; + assert!(file_size > super::METADATA_SIZE_HINT as u64); let fields = vec![DataField::new( 0, "id".to_string(), @@ -3637,7 +3724,7 @@ mod tests { file_size, &fields, None, - Some(32), + Some(8192), None, ) .await @@ -3648,7 +3735,7 @@ mod tests { file_size, &fields, None, - Some(32), + Some(8192), None, ) .await @@ -4841,6 +4928,48 @@ mod tests { } } + #[tokio::test] + async fn small_parquet_reuses_whole_file_metadata_prefetch() { + let data = Bytes::from( + write_multi_row_group_parquet(64, 64, EnabledStatistics::Chunk, false).await, + ); + assert!(data.len() <= super::METADATA_SIZE_HINT); + let tracker = TrackingFileRead::new(data.clone()); + let fields = vec![int_field("id"), int_field("value")]; + + let batches = ParquetFormatReader::default() + .read_batch_stream( + Box::new(tracker.clone()), + data.len() as u64, + &fields, + None, + None, + None, + ) + .await + .unwrap() + .try_collect::>() + .await + .unwrap(); + + assert_eq!(batches.iter().map(RecordBatch::num_rows).sum::(), 64); + assert_eq!(tracker.read_count(), 1); + assert_eq!(tracker.bytes_read(), data.len() as u64); + } + + #[tokio::test] + async fn whole_file_reuse_is_limited_to_metadata_prefetch_size() { + let data = Bytes::from(vec![0; super::METADATA_SIZE_HINT + 1]); + let tracker = TrackingFileRead::new(data.clone()); + let reader = + super::shared_parquet_file_reader(Box::new(tracker.clone()), data.len() as u64); + + reader.read(0..data.len() as u64).await.unwrap(); + reader.read(1..2).await.unwrap(); + + assert_eq!(tracker.read_count(), 2); + } + #[tokio::test] async fn parquet_metadata_cache_reuses_metadata_across_readers() { let data = Bytes::from( From 3cf8b02db607a45aa34665e72e31d1f11f133a4f Mon Sep 17 00:00:00 2001 From: xiaohongbo Date: Mon, 28 Sep 2026 16:09:16 -0700 Subject: [PATCH 2/2] [core] Share Parquet input across readers --- crates/paimon/src/arrow/format/parquet.rs | 194 +++++++++++----------- 1 file changed, 93 insertions(+), 101 deletions(-) diff --git a/crates/paimon/src/arrow/format/parquet.rs b/crates/paimon/src/arrow/format/parquet.rs index 863a23c98..e07e541d4 100644 --- a/crates/paimon/src/arrow/format/parquet.rs +++ b/crates/paimon/src/arrow/format/parquet.rs @@ -611,9 +611,9 @@ impl FormatFileReader for ParquetFormatReader { batch_size: Option, row_selection: Option>, ) -> crate::Result { - let shared_reader = shared_parquet_file_reader(reader, file_size); - let arrow_file_reader = ArrowFileReader::new(file_size, Arc::clone(&shared_reader)) + let arrow_file_reader = ArrowFileReader::new(file_size, reader.into()) .with_metadata_cache_enabled(self.metadata_cache_enabled); + let row_group_reader = arrow_file_reader.clone(); let empty_predicates = Vec::new(); let (preds, file_fields): (&[Predicate], &[DataField]) = match predicates { @@ -901,12 +901,11 @@ impl FormatFileReader for ParquetFormatReader { return; } }; - let row_group_reader = Arc::clone(&shared_reader); + let row_group_reader = row_group_reader.clone(); let row_group_metadata = reader_metadata.clone(); let row_group_mask = mask.clone(); tokio::spawn(read_row_group( row_group_reader, - file_size, row_group_metadata, row_group_mask, row_group_index, @@ -972,7 +971,7 @@ impl FormatFileReader for ParquetFormatReader { // merge input, nor wait for output buffers owned downstream. let _permit = budget.reserve_memory(estimated_bytes)?; let mut stream = build_row_group_stream( - Arc::clone(&shared_reader), file_size, metadata.clone(), mask.clone(), + row_group_reader.clone(), metadata.clone(), mask.clone(), index, batch_size, selection, shared_filters.clone(), )?; while let Some(batch) = stream.next().await { @@ -1075,8 +1074,7 @@ fn projected_row_group_bytes(row_group: &RowGroupMetaData, projection: &Projecti #[allow(clippy::too_many_arguments)] async fn read_row_group( - reader: Arc, - file_size: u64, + reader: ArrowFileReader, reader_metadata: ArrowReaderMetadata, projection: ProjectionMask, row_group_index: usize, @@ -1087,7 +1085,6 @@ async fn read_row_group( ) { let stream = match build_row_group_stream( reader, - file_size, reader_metadata, projection, row_group_index, @@ -1139,8 +1136,7 @@ impl ArrowPredicate for SharedParquetPredicate { #[allow(clippy::too_many_arguments)] fn build_row_group_stream( - reader: Arc, - file_size: u64, + reader: ArrowFileReader, metadata: ArrowReaderMetadata, projection: ProjectionMask, index: usize, @@ -1148,12 +1144,9 @@ fn build_row_group_stream( selection: Option, predicates: Vec, ) -> crate::Result> { - let mut builder = ParquetRecordBatchStreamBuilder::new_with_metadata( - ArrowFileReader::new(file_size, reader), - metadata, - ) - .with_projection(projection) - .with_row_groups(vec![index]); + let mut builder = ParquetRecordBatchStreamBuilder::new_with_metadata(reader, metadata) + .with_projection(projection) + .with_row_groups(vec![index]); if let Some(selection) = selection { builder = builder.with_row_selection(selection); } @@ -2498,13 +2491,14 @@ impl ParquetMetadataCache { let column_index = options.map_or(PageIndexPolicy::Skip, |o| o.column_index_policy()); let offset_index = options.map_or(PageIndexPolicy::Skip, |o| o.offset_index_policy()); let key = reader - .r + .input + .source .cache_key() .filter(|key| !key.is_empty()) .map(|file| ParquetMetadataCacheKey { file: file.to_string(), // Paimon data files are immutable; size separates supported replacements. - size: reader.file_size, + size: reader.input.file_size, column_index: page_index_policy_tag(column_index), offset_index: page_index_policy_tag(offset_index), }); @@ -2532,61 +2526,29 @@ fn page_index_policy_tag(policy: PageIndexPolicy) -> u8 { } } -/// ArrowFileReader is a wrapper around a FileRead that impls parquets AsyncFileReader. -/// -/// # TODO -/// -/// [ParquetObjectReader](https://docs.rs/parquet/latest/src/parquet/arrow/async_reader/store.rs.html#64) -/// contains the following hints to speed up metadata loading, similar to iceberg, we can consider adding them to this struct: -/// -/// - `metadata_size_hint`: Provide a hint as to the size of the parquet file's footer. -/// - `preload_column_index`: Load the Column Index as part of [`Self::get_metadata`]. -/// - `preload_offset_index`: Load the Offset Index as part of [`Self::get_metadata`]. -struct ArrowFileReader { - file_size: u64, - r: Arc, - metadata_cache_enabled: bool, -} - -/// coalesce threshold: 1 MiB. -const RANGE_COALESCE_BYTES: u64 = 1024 * 1024; -/// concurrent range fetches. -const RANGE_FETCH_CONCURRENCY: usize = 10; -/// metadata prefetch hint: 512 KiB. -const METADATA_SIZE_HINT: usize = 512 * 1024; -/// Minimum range size for splitting: 4 MiB. -/// The block size used for split alignment and as the minimum split -/// granularity. Ranges smaller than this will not be split further to -/// avoid excessive small IO requests whose per-request overhead dominates. -const IO_BLOCK_SIZE: u64 = 4 * 1024 * 1024; - -/// Reuses a whole-file metadata prefetch for subsequent data ranges. +/// One Parquet file input shared by metadata and row-group readers. /// -/// `ParquetMetaDataReader` reads the entire object when it is no larger than -/// [`METADATA_SIZE_HINT`]. Without this adapter, the decoded stream then asks -/// the object store for the selected data range again. Keep the already-read -/// bytes for the lifetime of this file reader and serve those ranges by slicing -/// the same [`Bytes`] allocation. -struct WholeFileCachingRead { - inner: Arc, +/// This mirrors Java's `ParquetFileReader` lifecycle: the footer and data reads +/// use the same file input. `ParquetMetaDataReader` fetches the entire object +/// when it is no larger than [`METADATA_SIZE_HINT`], so retain that fetch and +/// serve later data ranges by slicing the same [`Bytes`] allocation. +struct ParquetInput { + source: Arc, file_size: u64, - whole_file: OnceLock, + prefetched_file: OnceLock, } -impl WholeFileCachingRead { - fn new(inner: Arc, file_size: u64) -> Self { +impl ParquetInput { + fn new(source: Arc, file_size: u64) -> Self { Self { - inner, + source, file_size, - whole_file: OnceLock::new(), + prefetched_file: OnceLock::new(), } } -} -#[async_trait] -impl FileRead for WholeFileCachingRead { async fn read(&self, range: Range) -> crate::Result { - if let Some(whole_file) = self.whole_file.get() { + if let Some(file) = self.prefetched_file.get() { let start = usize::try_from(range.start).map_err(|error| Error::DataInvalid { message: format!("Parquet range start does not fit usize: {}", range.start), source: Some(Box::new(error)), @@ -2595,7 +2557,7 @@ impl FileRead for WholeFileCachingRead { message: format!("Parquet range end does not fit usize: {}", range.end), source: Some(Box::new(error)), })?; - if start > end || end > whole_file.len() { + if start > end || end > file.len() { return Err(Error::DataInvalid { message: format!( "Parquet range {}..{} exceeds file size {}", @@ -2604,41 +2566,49 @@ impl FileRead for WholeFileCachingRead { source: None, }); } - return Ok(whole_file.slice(start..end)); + return Ok(file.slice(start..end)); } - let bytes = self.inner.read(range.clone()).await?; - if range.start == 0 + let bytes = self.source.read(range.clone()).await?; + if self.file_size <= METADATA_SIZE_HINT as u64 + && range.start == 0 && range.end == self.file_size && u64::try_from(bytes.len()).ok() == Some(self.file_size) { - let _ = self.whole_file.set(bytes.clone()); + let _ = self.prefetched_file.set(bytes.clone()); } Ok(bytes) } - - fn cache_namespace(&self) -> Option { - self.inner.cache_namespace() - } - - fn cache_key(&self) -> Option<&str> { - self.inner.cache_key() - } - - fn file_format_metadata_cache(&self) -> Option<&(dyn std::any::Any + Send + Sync)> { - self.inner.file_format_metadata_cache() - } } -fn shared_parquet_file_reader(reader: Box, file_size: u64) -> Arc { - let reader: Arc = reader.into(); - if file_size <= METADATA_SIZE_HINT as u64 { - Arc::new(WholeFileCachingRead::new(reader, file_size)) - } else { - reader - } +/// Arrow adapter over the shared input of one Parquet file. +/// +/// # TODO +/// +/// [ParquetObjectReader](https://docs.rs/parquet/latest/src/parquet/arrow/async_reader/store.rs.html#64) +/// contains the following hints to speed up metadata loading, similar to iceberg, we can consider adding them to this struct: +/// +/// - `metadata_size_hint`: Provide a hint as to the size of the parquet file's footer. +/// - `preload_column_index`: Load the Column Index as part of [`Self::get_metadata`]. +/// - `preload_offset_index`: Load the Offset Index as part of [`Self::get_metadata`]. +#[derive(Clone)] +struct ArrowFileReader { + input: Arc, + metadata_cache_enabled: bool, } +/// coalesce threshold: 1 MiB. +const RANGE_COALESCE_BYTES: u64 = 1024 * 1024; +/// concurrent range fetches. +const RANGE_FETCH_CONCURRENCY: usize = 10; +/// metadata prefetch hint: 512 KiB. +const METADATA_SIZE_HINT: usize = 512 * 1024; +/// Minimum range size for splitting: 4 MiB. +/// The block size used for split alignment and as the minimum split +/// granularity. Ranges smaller than this will not be split further to +/// avoid excessive small IO requests whose per-request overhead dominates. +const IO_BLOCK_SIZE: u64 = 4 * 1024 * 1024; + /// Rows in a Parquet file, read from its footer alone. pub(crate) async fn read_row_count( reader: Box, @@ -2659,8 +2629,7 @@ pub(crate) async fn read_row_count( impl ArrowFileReader { fn new(file_size: u64, r: Arc) -> Self { Self { - file_size, - r, + input: Arc::new(ParquetInput::new(r, file_size)), metadata_cache_enabled: true, } } @@ -2671,7 +2640,7 @@ impl ArrowFileReader { } fn read_bytes(&mut self, range: Range) -> BoxFuture<'_, parquet::errors::Result> { - Box::pin(self.r.read(range.start..range.end).map_err(|err| { + Box::pin(self.input.read(range.start..range.end).map_err(|err| { let err_msg = format!("{err}"); parquet::errors::ParquetError::External(err_msg.into()) })) @@ -2684,7 +2653,7 @@ impl ArrowFileReader { let metadata_opts = options.map(|o| o.metadata_options().clone()); let column_index_policy = options.map(|o| o.column_index_policy()); let offset_index_policy = options.map(|o| o.offset_index_policy()); - let file_size = self.file_size; + let file_size = self.input.file_size; let mut reader = ParquetMetaDataReader::new() .with_prefetch_hint(Some(METADATA_SIZE_HINT)) .with_metadata_options(metadata_opts); @@ -2729,11 +2698,12 @@ impl AsyncFileReader for ArrowFileReader { let fetch_ranges = split_ranges_for_concurrency(coalesced, concurrency); // Fetch merged ranges concurrently. - let r = &self.r; + let input = &self.input; let fetched: Vec = if fetch_ranges.len() <= concurrency { // All ranges fit within the concurrency limit — fire them all at once. futures::future::try_join_all(fetch_ranges.iter().map(|range| { - r.read(range.clone()) + input + .read(range.clone()) .map_err(|e| parquet::errors::ParquetError::External(format!("{e}").into())) })) .await? @@ -2741,7 +2711,7 @@ impl AsyncFileReader for ArrowFileReader { // More ranges than concurrency slots — use buffered stream. futures::stream::iter(fetch_ranges.iter().cloned()) .map(|range| async move { - r.read(range).await.map_err(|e| { + input.read(range).await.map_err(|e| { parquet::errors::ParquetError::External(format!("{e}").into()) }) }) @@ -2831,9 +2801,14 @@ impl AsyncFileReader for ArrowFileReader { if !metadata_cache_enabled { return self.load_metadata(options.as_ref()).await; } - let Some(context) = self.r.file_format_metadata_cache().and_then(|cache| { - cache.downcast_ref::() - }) else { + let Some(context) = self + .input + .source + .file_format_metadata_cache() + .and_then(|cache| { + cache.downcast_ref::() + }) + else { return self.load_metadata(options.as_ref()).await; }; if context.max_bytes() == 0 { @@ -4961,15 +4936,32 @@ mod tests { async fn whole_file_reuse_is_limited_to_metadata_prefetch_size() { let data = Bytes::from(vec![0; super::METADATA_SIZE_HINT + 1]); let tracker = TrackingFileRead::new(data.clone()); - let reader = - super::shared_parquet_file_reader(Box::new(tracker.clone()), data.len() as u64); + let mut reader = super::ArrowFileReader::new(data.len() as u64, Arc::new(tracker.clone())); - reader.read(0..data.len() as u64).await.unwrap(); - reader.read(1..2).await.unwrap(); + reader.get_bytes(0..data.len() as u64).await.unwrap(); + reader.get_bytes(1..2).await.unwrap(); assert_eq!(tracker.read_count(), 2); } + #[tokio::test] + async fn cloned_parquet_readers_share_whole_file_prefetch() { + let data = Bytes::from(vec![0; super::METADATA_SIZE_HINT]); + let tracker = TrackingFileRead::new(data.clone()); + let mut metadata_reader = + super::ArrowFileReader::new(data.len() as u64, Arc::new(tracker.clone())); + let mut row_group_reader = metadata_reader.clone(); + + metadata_reader + .get_bytes(0..data.len() as u64) + .await + .unwrap(); + row_group_reader.get_bytes(1..2).await.unwrap(); + + assert_eq!(tracker.read_count(), 1); + assert_eq!(tracker.bytes_read(), data.len() as u64); + } + #[tokio::test] async fn parquet_metadata_cache_reuses_metadata_across_readers() { let data = Bytes::from(