Skip to content
Draft
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
199 changes: 160 additions & 39 deletions crates/paimon/src/arrow/format/parquet.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -611,9 +611,9 @@ impl FormatFileReader for ParquetFormatReader {
batch_size: Option<usize>,
row_selection: Option<Vec<RowRange>>,
) -> crate::Result<ArrowRecordBatchStream> {
let shared_reader: Arc<dyn FileRead> = 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 {
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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<dyn FileRead>,
file_size: u64,
reader: ArrowFileReader,
reader_metadata: ArrowReaderMetadata,
projection: ProjectionMask,
row_group_index: usize,
Expand All @@ -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,
Expand Down Expand Up @@ -1139,21 +1136,17 @@ impl ArrowPredicate for SharedParquetPredicate {

#[allow(clippy::too_many_arguments)]
fn build_row_group_stream(
reader: Arc<dyn FileRead>,
file_size: u64,
reader: ArrowFileReader,
metadata: ArrowReaderMetadata,
projection: ProjectionMask,
index: usize,
batch_size: Option<usize>,
selection: Option<RowSelection>,
predicates: Vec<SharedParquetPredicate>,
) -> crate::Result<parquet::arrow::async_reader::ParquetRecordBatchStream<ArrowFileReader>> {
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);
}
Expand Down Expand Up @@ -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),
});
Expand Down Expand Up @@ -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<dyn FileRead>,
file_size: u64,
prefetched_file: OnceLock<Bytes>,
}

impl ParquetInput {
fn new(source: Arc<dyn FileRead>, file_size: u64) -> Self {
Self {
source,
file_size,
prefetched_file: OnceLock::new(),
}
}

async fn read(&self, range: Range<u64>) -> crate::Result<Bytes> {
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
///
Expand All @@ -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<dyn FileRead>,
input: Arc<ParquetInput>,
metadata_cache_enabled: bool,
}

Expand Down Expand Up @@ -2580,8 +2629,7 @@ pub(crate) async fn read_row_count(
impl ArrowFileReader {
fn new(file_size: u64, r: Arc<dyn FileRead>) -> Self {
Self {
file_size,
r,
input: Arc::new(ParquetInput::new(r, file_size)),
metadata_cache_enabled: true,
}
}
Expand All @@ -2592,7 +2640,7 @@ impl ArrowFileReader {
}

fn read_bytes(&mut self, range: Range<u64>) -> BoxFuture<'_, parquet::errors::Result<Bytes>> {
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())
}))
Expand All @@ -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);
Expand Down Expand Up @@ -2650,19 +2698,20 @@ 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<Bytes> = 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?
} else {
// 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())
})
})
Expand Down Expand Up @@ -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::<crate::io::FileFormatMetadataCacheContext>()
}) else {
let Some(context) = self
.input
.source
.file_format_metadata_cache()
.and_then(|cache| {
cache.downcast_ref::<crate::io::FileFormatMetadataCacheContext>()
})
else {
return self.load_metadata(options.as_ref()).await;
};
if context.max_bytes() == 0 {
Expand Down Expand Up @@ -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();
Expand All @@ -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(),
Expand All @@ -3458,7 +3516,7 @@ mod tests {
file_size,
&fields,
None,
Some(32),
Some(8192),
None,
)
.await
Expand Down Expand Up @@ -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();
Expand All @@ -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(),
Expand All @@ -3637,7 +3699,7 @@ mod tests {
file_size,
&fields,
None,
Some(32),
Some(8192),
None,
)
.await
Expand All @@ -3648,7 +3710,7 @@ mod tests {
file_size,
&fields,
None,
Some(32),
Some(8192),
None,
)
.await
Expand Down Expand Up @@ -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::<Vec<_>>()
.await
.unwrap();

assert_eq!(batches.iter().map(RecordBatch::num_rows).sum::<usize>(), 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(
Expand Down
Loading