diff --git a/crates/paimon/src/arrow/format/parquet.rs b/crates/paimon/src/arrow/format/parquet.rs index 5ce567bfd..e07e541d4 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,9 +611,9 @@ impl FormatFileReader for ParquetFormatReader { batch_size: Option, row_selection: Option>, ) -> crate::Result { - let shared_reader: Arc = reader.into(); - 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,7 +2526,62 @@ fn page_index_policy_tag(policy: PageIndexPolicy) -> u8 { } } -/// ArrowFileReader is a wrapper around a FileRead that impls parquets AsyncFileReader. +/// One Parquet file input shared by metadata and row-group readers. +/// +/// 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, + prefetched_file: OnceLock, +} + +impl ParquetInput { + fn new(source: Arc, file_size: u64) -> Self { + Self { + source, + file_size, + prefetched_file: OnceLock::new(), + } + } + + async fn read(&self, range: Range) -> crate::Result { + 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)), + })?; + 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 > file.len() { + return Err(Error::DataInvalid { + message: format!( + "Parquet range {}..{} exceeds file size {}", + range.start, range.end, self.file_size + ), + source: None, + }); + } + return Ok(file.slice(start..end)); + } + + 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.prefetched_file.set(bytes.clone()); + } + Ok(bytes) + } +} + +/// Arrow adapter over the shared input of one Parquet file. /// /// # TODO /// @@ -2542,9 +2591,9 @@ fn page_index_policy_tag(policy: PageIndexPolicy) -> u8 { /// - `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 { - file_size: u64, - r: Arc, + input: Arc, metadata_cache_enabled: bool, } @@ -2580,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, } } @@ -2592,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()) })) @@ -2605,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); @@ -2650,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? @@ -2662,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()) }) }) @@ -2752,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 { @@ -3418,10 +3472,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 +3502,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 +3516,7 @@ mod tests { file_size, &fields, None, - Some(32), + Some(8192), None, ) .await @@ -3597,10 +3655,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 +3686,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 +3699,7 @@ mod tests { file_size, &fields, None, - Some(32), + Some(8192), None, ) .await @@ -3648,7 +3710,7 @@ mod tests { file_size, &fields, None, - Some(32), + Some(8192), None, ) .await @@ -4841,6 +4903,65 @@ 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 mut reader = super::ArrowFileReader::new(data.len() as u64, Arc::new(tracker.clone())); + + 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(