Skip to content
Open
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
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,7 @@ use arrow::datatypes::{DataType, Field, Schema, SchemaRef};
use criterion::{BatchSize, BenchmarkId, Criterion, criterion_group, criterion_main};
use datafusion_datasource_parquet::metadata::DFParquetMetadata;
use parquet::arrow::ArrowSchemaConverter;
use parquet::basic::ColumnOrder;
use parquet::data_type::ByteArray;
use parquet::file::metadata::{
ColumnChunkMetaData, FileMetaData, ParquetMetaData, RowGroupMetaData,
Expand Down Expand Up @@ -169,13 +170,19 @@ fn make_synthetic_metadata(
})
.collect::<Vec<_>>();

// Model a modern writer so byte-array bounds have a trustworthy order.
let column_orders = schema_descr
.columns()
.iter()
.map(|column| ColumnOrder::TYPE_DEFINED_ORDER(column.sort_order()))
.collect();
let file_metadata = FileMetaData::new(
1,
(spec.row_groups * ROWS_PER_GROUP) as i64,
Some("datafusion parquet metadata benchmark".to_string()),
None,
schema_descr,
None,
Some(column_orders),
);

ParquetMetaData::new(file_metadata, row_groups)
Expand Down
68 changes: 67 additions & 1 deletion datafusion/datasource-parquet/src/metadata.rs
Original file line number Diff line number Diff line change
Expand Up @@ -43,12 +43,13 @@ 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::basic::{ColumnOrder, SortOrder, Type as PhysicalType};
use parquet::file::metadata::{
PageIndexPolicy, ParquetMetaData, ParquetMetaDataPushDecoder, ParquetMetaDataReader,
RowGroupMetaData, SortingColumn,
};
use parquet::file::statistics::Statistics as ParquetStatistics;
use parquet::schema::types::SchemaDescriptor;
use parquet::schema::types::{ColumnDescriptor, SchemaDescriptor};
use std::any::Any;
use std::sync::Arc;

Expand All @@ -57,6 +58,54 @@ use std::sync::Arc;
/// would be too unreliable otherwise.
const PARTIAL_NDV_THRESHOLD: f64 = 0.75;

fn requires_unsigned_byte_array_order(column: &ColumnDescriptor) -> bool {
matches!(
column.physical_type(),
PhysicalType::BYTE_ARRAY | PhysicalType::FIXED_LEN_BYTE_ARRAY
) && column.sort_order() != SortOrder::SIGNED
}

/// Whether byte-array bounds lack a recognized unsigned comparison order.
///
/// The deprecated Parquet `min`/`max` fields use signed comparison, unlike
/// Arrow's string and binary comparisons. Even the modern bounds cannot be
/// interpreted without the corresponding footer `column_orders` entry.
/// Signed logical types, such as decimals, retain their existing behavior.
pub(crate) fn has_untrusted_byte_array_order(
parquet_schema: &SchemaDescriptor,
column_orders: Option<&[ColumnOrder]>,
parquet_column_index: usize,
) -> bool {
let column = parquet_schema.column(parquet_column_index);
requires_unsigned_byte_array_order(&column)
&& (column.sort_order() != SortOrder::UNSIGNED
|| column_orders
.and_then(|orders| orders.get(parquet_column_index))
.copied()
!= Some(ColumnOrder::TYPE_DEFINED_ORDER(SortOrder::UNSIGNED)))
}

/// Whether any row group provides byte-array bounds in the deprecated signed
/// order, or the column's logical order is undefined.
pub(crate) fn has_untrusted_byte_array_stats<'a>(
parquet_schema: &SchemaDescriptor,
parquet_column_index: Option<usize>,
row_groups: impl IntoIterator<Item = &'a RowGroupMetaData>,
) -> bool {
parquet_column_index.is_some_and(|index| {
let column = parquet_schema.column(index);
requires_unsigned_byte_array_order(&column)
&& (column.sort_order() != SortOrder::UNSIGNED
|| row_groups.into_iter().any(|group| {
group.column(index).statistics().is_some_and(|stats| {
stats.is_min_max_deprecated()
&& (stats.min_bytes_opt().is_some()
|| stats.max_bytes_opt().is_some())
})
}))
})
}

/// Handles fetching Parquet file schema, metadata and statistics
/// from object store.
///
Expand Down Expand Up @@ -496,6 +545,23 @@ impl<'a> DFParquetMetadata<'a> {
file_metadata.schema_descr(),
) {
Ok(stats_converter) => {
let parquet_index = stats_converter.parquet_column_index();
if parquet_index.is_some_and(|index| {
has_untrusted_byte_array_order(
file_metadata.schema_descr(),
file_metadata.column_orders().map(Vec::as_slice),
index,
)
}) || has_untrusted_byte_array_stats(
file_metadata.schema_descr(),
parquet_index,
row_groups_metadata,
) {
// The remaining row groups cannot establish bounds
// for the whole file. Keep unrelated statistics.
min_accs[idx] = None;
max_accs[idx] = None;
}
let mut accumulators = StatisticsAccumulators {
min_accs: &mut min_accs,
max_accs: &mut max_accs,
Expand Down
2 changes: 2 additions & 0 deletions datafusion/datasource-parquet/src/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,8 @@ mod schema_coercion;
mod sink;
mod sort;
pub mod source;
#[cfg(test)]
mod statistics_order_tests;
mod supported_predicates;
#[cfg(test)]
mod test_util;
Expand Down
5 changes: 2 additions & 3 deletions datafusion/datasource-parquet/src/opener/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1105,10 +1105,9 @@ impl FiltersPreparedParquetOpen {
// If there is a predicate that can be evaluated against the metadata
if let Some(predicate) = self.pruning_predicate.as_ref().map(|p| p.as_ref()) {
if prepared.enable_row_group_stats_pruning {
row_groups.prune_by_statistics(
row_groups.prune_by_statistics_with_metadata(
&prepared.physical_file_schema,
loaded.reader_metadata.parquet_schema(),
rg_metadata,
&file_metadata,
predicate,
&prepared.file_metrics,
);
Expand Down
15 changes: 15 additions & 0 deletions datafusion/datasource-parquet/src/page_filter.rs
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@ use std::sync::Arc;

use super::metrics::ParquetFileMetrics;
use crate::ParquetAccessPlan;
use crate::metadata::has_untrusted_byte_array_order;

use arrow::array::BooleanArray;
use arrow::{
Expand Down Expand Up @@ -490,6 +491,7 @@ struct PagesPruningStatistics<'a> {
column_index: &'a ParquetColumnIndex,
offset_index: &'a ParquetOffsetIndex,
page_offsets: &'a Vec<PageLocation>,
trusted_min_max: bool,
}

impl<'a> PagesPruningStatistics<'a> {
Expand Down Expand Up @@ -529,6 +531,12 @@ impl<'a> PagesPruningStatistics<'a> {
return None;
};
let page_offsets = offset_index_metadata.page_locations();
let file_metadata = parquet_metadata.file_metadata();
let trusted_min_max = !has_untrusted_byte_array_order(
file_metadata.schema_descr(),
file_metadata.column_orders().map(Vec::as_slice),
parquet_column_index,
);

Some(Self {
row_group_index,
Expand All @@ -537,6 +545,7 @@ impl<'a> PagesPruningStatistics<'a> {
column_index,
offset_index,
page_offsets,
trusted_min_max,
})
}

Expand All @@ -563,6 +572,9 @@ impl<'a> PagesPruningStatistics<'a> {
}
impl PruningStatistics for PagesPruningStatistics<'_> {
fn min_values(&self, _column: &datafusion_common::Column) -> Option<ArrayRef> {
if !self.trusted_min_max {
return None;
}
match self.converter.data_page_mins(
self.column_index,
self.offset_index,
Expand All @@ -577,6 +589,9 @@ impl PruningStatistics for PagesPruningStatistics<'_> {
}

fn max_values(&self, _column: &datafusion_common::Column) -> Option<ArrayRef> {
if !self.trusted_min_max {
return None;
}
match self.converter.data_page_maxes(
self.column_index,
self.offset_index,
Expand Down
5 changes: 5 additions & 0 deletions datafusion/datasource-parquet/src/push_decoder.rs
Original file line number Diff line number Diff line change
Expand Up @@ -211,6 +211,11 @@ impl RowGroupPruner {
.collect::<Vec<_>>();
let stats = RowGroupPruningStatistics {
parquet_schema: self.parquet_metadata.file_metadata().schema_descr(),
column_orders: self
.parquet_metadata
.file_metadata()
.column_orders()
.map(Vec::as_slice),
row_group_metadatas,
arrow_schema: self.arrow_schema.as_ref(),
// Match the existing static row-group pruning behavior: when a
Expand Down
Loading