diff --git a/datafusion/datasource-parquet/src/metadata.rs b/datafusion/datasource-parquet/src/metadata.rs index d3831766a42ab..512395ec220e4 100644 --- a/datafusion/datasource-parquet/src/metadata.rs +++ b/datafusion/datasource-parquet/src/metadata.rs @@ -40,7 +40,11 @@ use object_store::path::Path; use object_store::{ObjectMeta, ObjectStore}; use parquet::DecodeResult; use parquet::arrow::arrow_reader::statistics::StatisticsConverter; -use parquet::arrow::{parquet_column, parquet_to_arrow_schema}; +use parquet::arrow::arrow_reader::{ArrowReaderMetadata, ArrowReaderOptions}; +use parquet::arrow::{ + ProjectionMask, parquet_column, parquet_to_arrow_schema, + parquet_to_arrow_schema_and_field_levels, +}; use parquet::file::metadata::{ PageIndexPolicy, ParquetMetaData, ParquetMetaDataPushDecoder, RowGroupMetaData, SortingColumn, @@ -49,7 +53,7 @@ use parquet::file::statistics::Statistics as ParquetStatistics; use parquet::schema::types::SchemaDescriptor; use std::any::Any; use std::collections::HashMap; -use std::sync::Arc; +use std::sync::{Arc, OnceLock}; /// Minimum fraction of row groups that must report NDV statistics for the /// merged result to be `Inexact` rather than `Absent`, as the estimate @@ -770,15 +774,90 @@ fn has_any_exact_match( } /// Wrapper to implement [`FileMetadata`] for [`ParquetMetaData`]. -pub struct CachedParquetMetaData(Arc); +pub struct CachedParquetMetaData { + metadata: Arc, + /// Lazily-built [`ArrowReaderMetadata`] for this file. Constructing + /// this walks every leaf in the parquet schema (the "field-levels" + /// step), which is `O(N_columns)` work per file. Caching it lets + /// subsequent reader builds for the same file just `Clone` the result + /// (an `Arc` bump for the parquet metadata, the arrow schema, and the + /// dremel-level info) instead of redoing the walk. + arrow_reader_metadata: OnceLock, + /// Single-slot cache for [`ArrowReaderMetadata`] built with a supplied + /// schema (e.g. after `apply_file_schema_type_coercions` produced one). + /// Keyed by the supplied schema's `Arc` pointer — different supplied + /// schemas miss and overwrite the slot. In typical workloads (one query + /// touches every file with the same coerced table schema) every file + /// pays this rebuild at most once per session. + coerced_arm: std::sync::Mutex>, +} impl CachedParquetMetaData { pub fn new(metadata: Arc) -> Self { - Self(metadata) + Self { + metadata, + arrow_reader_metadata: OnceLock::new(), + coerced_arm: std::sync::Mutex::new(None), + } } pub fn parquet_metadata(&self) -> &Arc { - &self.0 + &self.metadata + } + + /// Return the cached [`ArrowReaderMetadata`] for this file, building it + /// on first use. + /// + /// `ArrowReaderMetadata` stores its schema and field-levels behind + /// `Arc`s, so [`ArrowReaderMetadata::clone`] is cheap and is what + /// callers should use to consume the cached value. + pub fn arrow_reader_metadata(&self) -> Result<&ArrowReaderMetadata> { + if let Some(v) = self.arrow_reader_metadata.get() { + return Ok(v); + } + let file_meta = self.metadata.file_metadata(); + let (schema, levels) = parquet_to_arrow_schema_and_field_levels( + file_meta.schema_descr(), + ProjectionMask::all(), + file_meta.key_value_metadata(), + )?; + let arm = ArrowReaderMetadata::from_field_levels( + Arc::clone(&self.metadata), + Arc::new(schema), + levels, + ); + // Race: if another thread also computed it, theirs wins. + let _ = self.arrow_reader_metadata.set(arm); + Ok(self.arrow_reader_metadata.get().expect("just set")) + } + + /// Get-or-build an [`ArrowReaderMetadata`] whose arrow schema is + /// `supplied_schema`. The result is memoised in a single-slot cache + /// keyed by `Arc::as_ptr(&supplied_schema)` — repeat calls with the + /// same schema (the common case: every file in a query coerces to the + /// same table schema) return a cheap [`ArrowReaderMetadata::clone`] + /// instead of re-walking the parquet schema. + pub fn coerced_arrow_reader_metadata( + &self, + supplied_schema: SchemaRef, + options: ArrowReaderOptions, + ) -> Result { + let key = Arc::as_ptr(&supplied_schema) as usize; + { + let guard = self.coerced_arm.lock().unwrap(); + if let Some((cached_key, cached)) = guard.as_ref() + && *cached_key == key + { + return Ok(cached.clone()); + } + } + let arm = ArrowReaderMetadata::try_new( + Arc::clone(&self.metadata), + options.with_schema(Arc::clone(&supplied_schema)), + )?; + let mut guard = self.coerced_arm.lock().unwrap(); + *guard = Some((key, arm.clone())); + Ok(arm) } } @@ -788,12 +867,12 @@ impl FileMetadata for CachedParquetMetaData { } fn memory_size(&self) -> usize { - self.0.memory_size() + self.metadata.memory_size() } fn extra_info(&self) -> HashMap { - let page_index = - self.0.column_index().is_some() && self.0.offset_index().is_some(); + let page_index = self.metadata.column_index().is_some() + && self.metadata.offset_index().is_some(); HashMap::from([("page_index".to_owned(), page_index.to_string())]) } } diff --git a/datafusion/datasource-parquet/src/reader.rs b/datafusion/datasource-parquet/src/reader.rs index 482bf8dced4f8..ed91d89396e83 100644 --- a/datafusion/datasource-parquet/src/reader.rs +++ b/datafusion/datasource-parquet/src/reader.rs @@ -28,11 +28,10 @@ use datafusion_physical_plan::metrics::ExecutionPlanMetricsSet; use futures::FutureExt; use futures::future::BoxFuture; use object_store::ObjectStore; -use parquet::arrow::arrow_reader::ArrowReaderOptions; +use arrow::datatypes::SchemaRef; +use parquet::arrow::arrow_reader::{ArrowReaderMetadata, ArrowReaderOptions}; use parquet::arrow::async_reader::{AsyncFileReader, ParquetObjectReader}; use parquet::file::metadata::ParquetMetaData; -use std::any::Any; -use std::collections::HashMap; use std::fmt::Debug; use std::ops::Range; use std::sync::Arc; @@ -319,6 +318,97 @@ impl AsyncFileReader for CachedParquetFileReader { } .boxed() } + + fn get_arrow_reader_metadata<'a>( + &'a mut self, + options: ArrowReaderOptions, + ) -> BoxFuture<'a, parquet::errors::Result> { + let object_meta = self.partitioned_file.object_meta.clone(); + let metadata_cache = Arc::clone(&self.metadata_cache); + + async move { + // We can serve from cache when no embedded-schema-skip and no + // virtual columns are requested. A `supplied_schema` is OK — + // the wrapper has a separate cache slot for post-coercion + // builds keyed by the supplied schema's Arc identity. + let can_use_cache = + !options.skip_arrow_metadata() && options.virtual_columns().is_empty(); + let supplied = options.supplied_schema().cloned(); + + #[cfg(feature = "parquet_encryption")] + let file_decryption_properties = + options.file_decryption_properties().map(Arc::clone); + #[cfg(not(feature = "parquet_encryption"))] + let file_decryption_properties = None; + + let try_serve_from_cache = |options: ArrowReaderOptions, + supplied: Option| + -> parquet::errors::Result> { + if !can_use_cache { + return Ok(None); + } + let Some(cached) = metadata_cache.get(&object_meta.location) else { + return Ok(None); + }; + if !cached.is_valid_for(&object_meta) { + return Ok(None); + } + let Some(cached_parquet) = cached + .file_metadata + .as_any() + .downcast_ref::() + else { + return Ok(None); + }; + let arm = if let Some(schema) = supplied { + cached_parquet + .coerced_arrow_reader_metadata(schema, options) + .map_err(|e| { + parquet::errors::ParquetError::General(format!( + "Failed to build coerced arrow reader metadata for {}: {e}", + object_meta.location, + )) + })? + } else { + cached_parquet + .arrow_reader_metadata() + .map_err(|e| { + parquet::errors::ParquetError::General(format!( + "Failed to build arrow reader metadata for {}: {e}", + object_meta.location, + )) + })? + .clone() + }; + Ok(Some(arm)) + }; + + // Fast path: cache hit (already-fetched metadata). + if let Some(arm) = try_serve_from_cache(options.clone(), supplied.clone())? { + return Ok(arm); + } + + // Slow path: fetch + cache the metadata, then retry. + let metadata = DFParquetMetadata::new(&self.store, &object_meta) + .with_decryption_properties(file_decryption_properties) + .with_file_metadata_cache(Some(Arc::clone(&metadata_cache))) + .with_metadata_size_hint(self.metadata_size_hint) + .fetch_metadata() + .await + .map_err(|e| { + parquet::errors::ParquetError::General(format!( + "Failed to fetch metadata for file {}: {e}", + object_meta.location, + )) + })?; + if let Some(arm) = try_serve_from_cache(options.clone(), supplied)? { + return Ok(arm); + } + + ArrowReaderMetadata::try_new(metadata, options) + } + .boxed() + } } impl Drop for CachedParquetFileReader { @@ -333,31 +423,7 @@ impl Drop for CachedParquetFileReader { } } -/// Wrapper to implement [`FileMetadata`] for [`ParquetMetaData`]. -pub struct CachedParquetMetaData(Arc); - -impl CachedParquetMetaData { - pub fn new(metadata: Arc) -> Self { - Self(metadata) - } - - pub fn parquet_metadata(&self) -> &Arc { - &self.0 - } -} - -impl FileMetadata for CachedParquetMetaData { - fn as_any(&self) -> &dyn Any { - self - } - - fn memory_size(&self) -> usize { - self.0.memory_size() - } - - fn extra_info(&self) -> HashMap { - let page_index = - self.0.column_index().is_some() && self.0.offset_index().is_some(); - HashMap::from([("page_index".to_owned(), page_index.to_string())]) - } -} +// `CachedParquetMetaData` lives in `crate::metadata` (where it carries the +// lazily-built, cached `ArrowReaderMetadata`). Re-export it here so the +// `pub use reader::*` in `mod.rs` keeps exposing it at the crate root. +pub use crate::metadata::CachedParquetMetaData;