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
229 changes: 178 additions & 51 deletions datafusion/datasource-parquet/src/opener/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -556,10 +556,7 @@ impl ParquetOpenState {
}
ParquetOpenState::PruneWithStatistics(prepared) => {
let prepared_row_groups = (*prepared).prune_row_groups()?;
if should_load_page_index(
prepared_row_groups.prepared.page_pruning_predicate.as_ref(),
&prepared_row_groups.row_groups,
) {
if prepared_row_groups.should_load_page_index() {
Ok(ParquetOpenState::LoadPageIndex(
prepared_row_groups.load_page_index().boxed(),
))
Expand Down Expand Up @@ -1153,6 +1150,50 @@ impl FiltersPreparedParquetOpen {
}

impl RowGroupsPrunedParquetOpen {
/// Returns true if the reader would benefit from a page index load, given
/// the current pruning predicate and row group access plan.
///
/// The page index is used for data page pruning, and it is only useful
/// when:
///
/// 1. There is at least one row group that may have filtered rows
/// (if it is fully matched we know no rows will be filtered)
///
/// 2. There is a page index for at least one predicate column (some
/// parquet writers do not write the page index).
fn should_load_page_index(&self) -> bool {
let Some(page_pruning_predicate) = self.prepared.page_pruning_predicate.as_ref()
else {
return false;
};
let row_groups = &self.row_groups;
let fully_matched = row_groups.is_fully_matched();
// if all row groups are fully matched, nothing can be pruned
if row_groups.row_group_indexes().all(|idx| fully_matched[idx]) {
return false;
}

// Check the file's footer metadata to see if a page index was written
// for at least one predicate column in a surviving row group.
//
// Note: offsets are recorded in the footer, so we can determine if a
// page index exists before attempting to read it.
let parquet_metadata = self.prepared.loaded.reader_metadata.metadata();
let arrow_schema = &self.prepared.loaded.prepared.physical_file_schema;
let parquet_schema = parquet_metadata.file_metadata().schema_descr();
page_pruning_predicate.predicate_column_names().any(|name| {
let Some((leaf_idx, _)) = parquet_column(parquet_schema, arrow_schema, name)
else {
return false;
};
row_groups.row_group_indexes().any(|rg_idx| {
let column = parquet_metadata.row_group(rg_idx).column(leaf_idx);
column.column_index_offset().is_some()
&& column.offset_index_offset().is_some()
})
})
}

/// Load the page index if pruning requires it and metadata did not include it.
async fn load_page_index(mut self) -> Result<Self> {
self.prepared.loaded.reader_metadata = load_page_index(
Expand Down Expand Up @@ -1651,22 +1692,6 @@ pub(crate) fn build_pruning_predicates(
.build(Arc::clone(predicate))
}

/// Returns true if the page index must be loaded for page-level pruning.
///
/// The page index can only prune when at least one surviving row group is not
/// fully matched by row-group statistics alone.
fn should_load_page_index(
page_pruning_predicate: Option<&Arc<PagePruningAccessPlanFilter>>,
row_groups: &RowGroupAccessPlanFilter,
) -> bool {
page_pruning_predicate.is_some_and(|_| {
let fully_matched = row_groups.is_fully_matched();
row_groups
.row_group_indexes()
.any(|idx| !fully_matched[idx])
})
}

/// Returns a `ArrowReaderMetadata` with the page index loaded, loading
/// it from the underlying `AsyncFileReader` if necessary.
async fn load_page_index<T: AsyncFileReader>(
Expand Down Expand Up @@ -1719,7 +1744,7 @@ mod test {
CachedFileMetadataEntry, FileMetadataCache,
};
use datafusion_execution::cache::default_cache::DefaultCache;
use datafusion_expr::{col, lit};
use datafusion_expr::{Expr, col, lit};
use datafusion_physical_expr::{
PhysicalExpr,
expressions::{Column, DynamicFilterPhysicalExpr, Literal},
Expand All @@ -1734,10 +1759,10 @@ mod test {
use futures::StreamExt;
use futures::stream::BoxStream;
use object_store::{ObjectStore, ObjectStoreExt, memory::InMemory, path::Path};
use parquet::arrow::ArrowWriter;
use parquet::file::metadata::ColumnChunkMetaData;
use parquet::arrow::{ArrowSchemaConverter, ArrowWriter};
use parquet::file::metadata::{ColumnChunkMetaData, FileMetaData, ParquetMetaData};
use parquet::file::properties::WriterProperties;
use parquet::schema::types::{SchemaDescPtr, SchemaDescriptor};
use parquet::schema::types::SchemaDescPtr;
use std::collections::VecDeque;
use std::sync::Arc;

Expand Down Expand Up @@ -1834,19 +1859,121 @@ mod test {
.collect()
}

#[test]
fn should_load_page_index_checks_predicate_columns() {
// "a" has page index offsets recorded in the footer, "b" does not
let metadata = page_index_metadata(&[("a", true), ("b", false)], 1);

// predicate on "a": the file has a page index for it, so load it
assert!(should_load_page_index(
metadata.clone(),
Some(col("a").gt(lit(50i32))),
ParquetAccessPlan::new_all(1),
));

// predicate on "b": no page index for that column, so skip the load
assert!(!should_load_page_index(
metadata,
Some(col("b").gt(lit(50i32))),
ParquetAccessPlan::new_all(1),
));
}

fn test_schema_descr() -> SchemaDescPtr {
use parquet::basic::{LogicalType, Type as PhysicalType};
use parquet::schema::types::Type as SchemaType;
let schema = Schema::new(vec![Field::new("a", DataType::Utf8, false)]);
Arc::new(ArrowSchemaConverter::new().convert(&schema).unwrap())
}

/// Metadata for a file of Int32 `columns`, where each `(name,
/// has_page_index)` entry controls whether the footer records page index
/// offsets for that column.
fn page_index_metadata(
columns: &[(&str, bool)],
num_row_groups: usize,
) -> ParquetMetaData {
let arrow_schema = Schema::new(
columns
.iter()
.map(|(name, _)| Field::new(*name, DataType::Int32, false))
.collect::<Vec<_>>(),
);
let schema_descr =
Arc::new(ArrowSchemaConverter::new().convert(&arrow_schema).unwrap());

let row_groups = (0..num_row_groups)
.map(|_| {
let columns = columns
.iter()
.enumerate()
.map(|(idx, (_, has_page_index))| {
let mut builder =
ColumnChunkMetaData::builder(schema_descr.column(idx))
.set_num_values(10);
if *has_page_index {
builder = builder
.set_column_index_offset(Some(100))
.set_column_index_length(Some(10))
.set_offset_index_offset(Some(110))
.set_offset_index_length(Some(10));
}
builder.build().unwrap()
})
.collect();
RowGroupMetaData::builder(Arc::clone(&schema_descr))
.set_num_rows(10)
.set_column_metadata(columns)
.build()
.unwrap()
})
.collect();
let file_metadata =
FileMetaData::new(1, 10, None, None, Arc::clone(&schema_descr), None);
ParquetMetaData::new(file_metadata, row_groups)
}

/// Reports [`RowGroupsPrunedParquetOpen::should_load_page_index`] for
/// hand-built parquet `metadata` (no I/O), an optional predicate, and a
/// row group access plan.
fn should_load_page_index(
metadata: ParquetMetaData,
predicate: Option<Expr>,
plan: ParquetAccessPlan,
) -> bool {
use crate::RowGroupAccessPlanFilter;
use parquet::arrow::parquet_to_arrow_schema;

let field = SchemaType::primitive_type_builder("a", PhysicalType::BYTE_ARRAY)
.with_logical_type(Some(LogicalType::String))
.build()
.unwrap();
let schema = SchemaType::group_type_builder("schema")
.with_fields(vec![Arc::new(field)])
.build()
.unwrap();
Arc::new(SchemaDescriptor::new(Arc::new(schema)))
let arrow_schema: SchemaRef = Arc::new(
parquet_to_arrow_schema(metadata.file_metadata().schema_descr(), None)
.unwrap(),
);
let page_pruning_predicate = predicate.map(|expr| {
let predicate = logical2physical(&expr, &arrow_schema);
build_page_pruning_predicate(&predicate, &arrow_schema)
});

let store: Arc<dyn ObjectStore> = Arc::new(InMemory::new());
let morselizer = ParquetMorselizerBuilder::new()
.with_store(store)
.with_schema(Arc::clone(&arrow_schema))
.build();
let file = PartitionedFile::new("test.parquet".to_string(), 100);
let prepared = morselizer.prepare_open_file(file).unwrap();
let options = ArrowReaderOptions::new();
let reader_metadata =
ArrowReaderMetadata::try_new(Arc::new(metadata), options.clone()).unwrap();
let open = RowGroupsPrunedParquetOpen {
prepared: FiltersPreparedParquetOpen {
loaded: MetadataLoadedParquetOpen {
prepared,
reader_metadata,
options,
},
pruning_predicate: None,
page_pruning_predicate,
},
row_groups: RowGroupAccessPlanFilter::new(plan),
};
open.should_load_page_index()
}

impl ParquetMorselizerBuilder {
Expand Down Expand Up @@ -3034,31 +3161,31 @@ mod test {

#[test]
fn should_load_page_index_without_predicate() {
use crate::RowGroupAccessPlanFilter;
let row_groups = RowGroupAccessPlanFilter::new(ParquetAccessPlan::new_all(2));
assert!(!should_load_page_index(None, &row_groups));
assert!(!should_load_page_index(
page_index_metadata(&[("a", true)], 2),
None,
ParquetAccessPlan::new_all(2),
));
}

#[test]
fn should_load_page_index_when_surviving_row_groups_not_fully_matched() {
use crate::RowGroupAccessPlanFilter;
let schema = Arc::new(Schema::new(vec![Field::new("a", DataType::Int32, false)]));
let predicate = logical2physical(&col("a").gt(lit(50i32)), &schema);
let page_predicate = build_page_pruning_predicate(&predicate, &schema);
let row_groups = RowGroupAccessPlanFilter::new(ParquetAccessPlan::new_all(2));
assert!(should_load_page_index(Some(&page_predicate), &row_groups));
assert!(should_load_page_index(
page_index_metadata(&[("a", true)], 2),
Some(col("a").gt(lit(50i32))),
ParquetAccessPlan::new_all(2),
));
}

#[test]
fn should_load_page_index_when_all_surviving_row_groups_fully_matched() {
use crate::RowGroupAccessPlanFilter;
let schema = Arc::new(Schema::new(vec![Field::new("a", DataType::Int32, false)]));
let predicate = logical2physical(&col("a").is_not_null(), &schema);
let page_predicate = build_page_pruning_predicate(&predicate, &schema);
let mut plan = ParquetAccessPlan::new_all(1);
plan.mark_fully_matched(0);
let row_groups = RowGroupAccessPlanFilter::new(plan);
assert!(!should_load_page_index(Some(&page_predicate), &row_groups));
assert!(!should_load_page_index(
page_index_metadata(&[("a", true)], 1),
Some(col("a").is_not_null()),
plan,
));
}

#[tokio::test]
Expand Down Expand Up @@ -3692,7 +3819,7 @@ mod test {
async fn build_pushdown_morselizer(
store: &Arc<dyn ObjectStore>,
path: &str,
predicate_expr: datafusion_expr::Expr,
predicate_expr: Expr,
pushdown_filters: bool,
) -> Result<(ParquetMorselizer, PartitionedFile)> {
let (file_schema, data_size) = write_grouped_file(store, path, 1, 5).await;
Expand Down
Loading
Loading