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
134 changes: 132 additions & 2 deletions datafusion/datasource-parquet/src/metadata.rs
Original file line number Diff line number Diff line change
Expand Up @@ -45,8 +45,8 @@ 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,
FileMetaData, PageIndexPolicy, ParquetMetaData, ParquetMetaDataPushDecoder,
ParquetMetaDataReader, RowGroupMetaData, SortingColumn,
};
use parquet::file::statistics::Statistics as ParquetStatistics;
use parquet::schema::types::{ColumnDescriptor, SchemaDescriptor};
Expand Down Expand Up @@ -112,6 +112,35 @@ pub(crate) fn has_untrusted_byte_array_stats<'a>(
})
}

/// Whether a missing Parquet row-group `null_count` can be treated as exactly
/// zero for a file written by the writer recorded in `created_by`.
///
/// parquet-rs before 53.1.0 (fixed in apache/arrow-rs#6490) did not record a
/// `null_count` when it was zero, so for those writers a missing count is
/// exactly zero. Files written by DataFusion < 42.1.0 could link such a
/// parquet-rs, so their missing counts are exactly zero as well.
///
/// For every other writer a missing `null_count` is unknown. Treating it as
/// zero would let `IS NULL` / `COUNT` pruning and limit pruning return the
/// wrong results, so those files must be handled conservatively.
pub(crate) fn missing_null_counts_are_zero(file_metadata: &FileMetaData) -> bool {
let Some((writer, version)) = file_metadata
.created_by()
.and_then(|s| s.split_once(" version "))
else {
return false;
};
let mut parts = version.split(['.', ' ', '-']).map(str::parse::<u64>);
let (Some(Ok(major)), Some(Ok(minor))) = (parts.next(), parts.next()) else {
return false;
};
match writer {
"parquet-rs" => (major, minor) < (53, 1),
"datafusion" => (major, minor) < (42, 1),
_ => false,
}
}

/// Handles fetching Parquet file schema, metadata and statistics
/// from object store.
///
Expand Down Expand Up @@ -519,6 +548,7 @@ impl<'a> DFParquetMetadata<'a> {
statistics.num_rows = Precision::Exact(num_rows);

let file_metadata = metadata.file_metadata();
let missing_null_counts_as_zero = missing_null_counts_are_zero(file_metadata);
let mut physical_file_schema = parquet_to_arrow_schema(
file_metadata.schema_descr(),
file_metadata.key_value_metadata(),
Expand Down Expand Up @@ -551,10 +581,23 @@ impl<'a> DFParquetMetadata<'a> {
file_metadata.schema_descr(),
) {
Ok(stats_converter) => {

// A missing null_count is only exactly zero for
// writers known to omit it (i.e. old parquet-rs).
// For every other writer a missing count is
// unknown, and must not surface as an exact zero
// in file statistics used by pruning, aggregates
// and sort pushdown.
let stats_converter = stats_converter
.with_missing_null_counts_as_zero(
missing_null_counts_as_zero,
);

// An omitted count must not become an exact zero in
// file statistics used for pruning and aggregates.
let stats_converter =
stats_converter.with_missing_null_counts_as_zero(false);

let parquet_index = stats_converter.parquet_column_index();
if parquet_index.is_some_and(|index| {
has_untrusted_min_max_order(
Expand Down Expand Up @@ -1192,6 +1235,93 @@ mod tests {
use arrow::array::Int32Array;
use arrow::compute::SortOptions;
use arrow::datatypes::Field;
use parquet::schema::types::Type as ParquetType;

/// Builds `FileMetaData` for a single-column INT32 schema with the given
/// `created_by` string.
fn file_metadata_with_created_by(created_by: Option<&str>) -> FileMetaData {
let schema = Arc::new(SchemaDescriptor::new(Arc::new(
ParquetType::group_type_builder("schema")
.with_fields(vec![Arc::new(
ParquetType::primitive_type_builder("a", PhysicalType::INT32)
.build()
.unwrap(),
)])
.build()
.unwrap(),
)));
FileMetaData::new(1, 0, created_by.map(str::to_string), None, schema, None)
}

#[test]
fn test_missing_null_counts_are_zero() {
// parquet-rs < 53.1.0 omitted null counts that are zero
for created_by in [
"parquet-rs version 5.1.0",
"parquet-rs version 52.0.1",
"parquet-rs version 53.0.0",
"parquet-rs version 53.0.0 (build abc)",
] {
assert!(
missing_null_counts_are_zero(&file_metadata_with_created_by(Some(
created_by
))),
"expected {created_by:?} missing counts to be treated as zero"
);
}
// parquet-rs >= 53.1.0 always records null counts
for created_by in [
"parquet-rs version 53.1.0",
"parquet-rs version 59.3.0",
"parquet-rs version 60.0.0",
] {
assert!(
!missing_null_counts_are_zero(&file_metadata_with_created_by(Some(
created_by
))),
"expected {created_by:?} missing counts to stay unknown"
);
}
// DataFusion < 42.1.0 may link a parquet-rs that omits zero null counts
for created_by in [
"datafusion version 5.1.0",
"datafusion version 41.1.2",
"datafusion version 42.0.0",
] {
assert!(
missing_null_counts_are_zero(&file_metadata_with_created_by(Some(
created_by
))),
"expected {created_by:?} missing counts to be treated as zero"
);
}
// DataFusion >= 42.1.0 always records null counts
for created_by in ["datafusion version 42.1.0", "datafusion version 55.1.0"] {
assert!(
!missing_null_counts_are_zero(&file_metadata_with_created_by(Some(
created_by
))),
"expected {created_by:?} missing counts to stay unknown"
);
}
// Unknown writers, unparsable versions and absent created_by are
// conservative: missing counts stay unknown
for created_by in [
None,
Some("parquet-mr version 1.13.1"),
Some("duckdb version 1.1.0"),
Some("custom string"),
Some("parquet-rs"),
Some("parquet-rs version not-a-version"),
Some("datafusion version 42"),
] {
let metadata = file_metadata_with_created_by(created_by);
assert!(
!missing_null_counts_are_zero(&metadata),
"expected {created_by:?} missing counts to stay unknown"
);
}
}

#[test]
fn test_lex_ordering_to_sorting_columns_uses_writer_schema() -> Result<()> {
Expand Down
17 changes: 11 additions & 6 deletions datafusion/datasource-parquet/src/push_decoder.rs
Original file line number Diff line number Diff line change
Expand Up @@ -62,6 +62,7 @@ use datafusion_pruning::{PruningPredicate, PruningPredicateBuilder};
use crate::ParquetFileMetrics;
use crate::access_plan::PreparedAccessPlan;
use crate::decoder_projection::DecoderProjection;
use crate::metadata::missing_null_counts_are_zero;
use crate::metrics::{ByteProgress, RowFilterSkippedFullyMatchedMetric};
use crate::row_filter::{
PrebuiltRowFilterCandidate, prebuild_row_filter_candidates, row_filter_from_prebuilt,
Expand Down Expand Up @@ -244,15 +245,19 @@ impl RowGroupPruner {
.iter()
.map(|&i| self.parquet_metadata.row_group(i))
.collect::<Vec<_>>();
let file_metadata = self.parquet_metadata.file_metadata();
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),
parquet_schema: file_metadata.schema_descr(),
column_orders: file_metadata.column_orders().map(Vec::as_slice),
row_group_metadatas,
arrow_schema: self.arrow_schema.as_ref(),

// Match the static row-group pruning behavior: a missing null count
// is exactly zero for old parquet-rs / DataFusion writers and
// unknown for everyone else. Runtime pruning only needs to prove a
// row group *cannot* contain matching rows, so this is sound.
missing_null_counts_as_zero: missing_null_counts_are_zero(file_metadata),

};

match pp.prune(&stats) {
Expand Down
50 changes: 44 additions & 6 deletions datafusion/datasource-parquet/src/row_group_filter.rs
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,10 @@ use std::sync::Arc;

use super::{ParquetAccessPlan, ParquetFileMetrics, RowGroupAccess};
use crate::bloom_filter::BloomFilterStatistics;
use crate::metadata::{has_untrusted_byte_array_stats, has_untrusted_min_max_order};
use crate::metadata::{
has_untrusted_byte_array_stats, has_untrusted_min_max_order,
missing_null_counts_are_zero,
};
use crate::pruning::build_inverted_predicate;
use arrow::array::{ArrayRef, BooleanArray, UInt64Array};
use arrow::compute::nullif;
Expand All @@ -31,7 +34,7 @@ use datafusion_datasource::FileRange;
use datafusion_pruning::PruningPredicate;
use parquet::arrow::arrow_reader::statistics::StatisticsConverter;
use parquet::basic::ColumnOrder;
use parquet::file::metadata::{ParquetMetaData, RowGroupMetaData};
use parquet::file::metadata::{FileMetaData, ParquetMetaData, RowGroupMetaData};
use parquet::schema::types::SchemaDescriptor;

/// Reduces the [`ParquetAccessPlan`] based on row group level metadata.
Expand Down Expand Up @@ -278,9 +281,9 @@ impl RowGroupAccessPlanFilter {
arrow_schema,
parquet_schema,
groups,
None,
predicate,
metrics,
None,
);
}

Expand All @@ -304,9 +307,9 @@ impl RowGroupAccessPlanFilter {
arrow_schema,
file_metadata.schema_descr(),
metadata.row_groups(),
file_metadata.column_orders().map(Vec::as_slice),
predicate,
metrics,
Some(file_metadata),
);
}

Expand All @@ -315,9 +318,9 @@ impl RowGroupAccessPlanFilter {
arrow_schema: &Schema,
parquet_schema: &SchemaDescriptor,
groups: &[RowGroupMetaData],
column_orders: Option<&[ColumnOrder]>,
predicate: &PruningPredicate,
metrics: &ParquetFileMetrics,
file_metadata: Option<&FileMetaData>,
) {
// scoped timer updates on drop
let _timer_guard = metrics.statistics_eval_time.timer();
Expand All @@ -330,11 +333,29 @@ impl RowGroupAccessPlanFilter {
.map(|&i| &groups[i])
.collect::<Vec<_>>();

// A missing null count is exactly zero for old parquet-rs writers, but
// unknown for everything else. When no footer metadata is available
// (this only happens in tests), fall back to the StatisticsConverter
// default that treats a missing count as zero to preserve the original
// pruning behavior.
let missing_null_counts_as_zero = file_metadata
.map(missing_null_counts_are_zero)
.unwrap_or(true);

// The footer also records the comparison order of string and binary
// bounds; `prune_by_statistics` has no footer and passes `None`.
let column_orders = file_metadata
.and_then(|metadata| metadata.column_orders())
.map(Vec::as_slice);

let pruning_stats = RowGroupPruningStatistics {
parquet_schema,
column_orders,
row_group_metadatas,
arrow_schema,

missing_null_counts_as_zero,

};

// try to prune the row groups in a single call
Expand All @@ -351,10 +372,17 @@ impl RowGroupAccessPlanFilter {
}
}

// Check if any of the matched row groups are fully contained by the predicate
// Fully matched row groups require a stronger proof: every row
// must pass the predicate. When no footer is available, fall
// back to false (not true) so that the fully-matched claim
// stays sound for limit pruning.
let fully_matched_flag = file_metadata
.map(missing_null_counts_are_zero)
.unwrap_or(false);
self.identify_fully_matched_row_groups(
&fully_contained_candidates_original_idx,
&pruning_stats,
fully_matched_flag,
groups,
predicate,
metrics,
Expand All @@ -380,6 +408,7 @@ impl RowGroupAccessPlanFilter {
&mut self,
candidate_row_group_indices: &[usize],
pruning_stats: &RowGroupPruningStatistics<'_>,
fully_matched_missing_null_counts_as_zero: bool,
groups: &[RowGroupMetaData],
predicate: &PruningPredicate,
metrics: &ParquetFileMetrics,
Expand All @@ -402,6 +431,15 @@ impl RowGroupAccessPlanFilter {
.map(|&i| &groups[i])
.collect::<Vec<_>>(),
arrow_schema,

// Fully matched row groups require a stronger proof: every row
// must pass the predicate. Use the flag derived from the footer
// metadata, which for old parquet-rs writers confirms that a
// missing count is genuinely zero. Without footer metadata, use
// false so the claim stays sound for limit pruning.
missing_null_counts_as_zero: fully_matched_missing_null_counts_as_zero,


};

let Ok(inverted_values) = inverted_predicate.prune(&inverted_pruning_stats)
Expand Down
Loading