Skip to content
Draft
Show file tree
Hide file tree
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
95 changes: 87 additions & 8 deletions datafusion/datasource-parquet/src/metadata.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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
Expand Down Expand Up @@ -770,15 +774,90 @@ fn has_any_exact_match(
}

/// Wrapper to implement [`FileMetadata`] for [`ParquetMetaData`].
pub struct CachedParquetMetaData(Arc<ParquetMetaData>);
pub struct CachedParquetMetaData {
metadata: Arc<ParquetMetaData>,
/// 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<ArrowReaderMetadata>,
/// 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<Option<(usize, ArrowReaderMetadata)>>,
}

impl CachedParquetMetaData {
pub fn new(metadata: Arc<ParquetMetaData>) -> Self {
Self(metadata)
Self {
metadata,
arrow_reader_metadata: OnceLock::new(),
coerced_arm: std::sync::Mutex::new(None),
}
}

pub fn parquet_metadata(&self) -> &Arc<ParquetMetaData> {
&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<ArrowReaderMetadata> {
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)
}
}

Expand All @@ -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<String, String> {
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())])
}
}
Expand Down
128 changes: 97 additions & 31 deletions datafusion/datasource-parquet/src/reader.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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<ArrowReaderMetadata>> {
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<SchemaRef>|
-> parquet::errors::Result<Option<ArrowReaderMetadata>> {
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::<CachedParquetMetaData>()
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 {
Expand All @@ -333,31 +423,7 @@ impl Drop for CachedParquetFileReader {
}
}

/// Wrapper to implement [`FileMetadata`] for [`ParquetMetaData`].
pub struct CachedParquetMetaData(Arc<ParquetMetaData>);

impl CachedParquetMetaData {
pub fn new(metadata: Arc<ParquetMetaData>) -> Self {
Self(metadata)
}

pub fn parquet_metadata(&self) -> &Arc<ParquetMetaData> {
&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<String, String> {
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;
Loading