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
106 changes: 59 additions & 47 deletions datafusion/core/tests/parquet/row_group_pruning.rs

Large diffs are not rendered by default.

15 changes: 5 additions & 10 deletions datafusion/core/tests/sql/explain_analyze.rs
Original file line number Diff line number Diff line change
Expand Up @@ -880,10 +880,7 @@ async fn parquet_explain_analyze() {

// should contain aggregated stats
assert_contains!(&formatted, "output_rows=8");
assert_contains!(
&formatted,
"row_groups_pruned_bloom_filter=1 total \u{2192} 1 matched"
);
assert_not_contains!(&formatted, "row_groups_pruned_bloom_filter");
assert_contains!(
&formatted,
"row_groups_pruned_statistics=1 total \u{2192} 1 matched"
Expand All @@ -896,16 +893,14 @@ async fn parquet_explain_analyze() {
// (file-> row-group -> page)
let i_file = formatted.find("files_ranges_pruned_statistics").unwrap();
let i_rowgroup_stat = formatted.find("row_groups_pruned_statistics").unwrap();
let i_rowgroup_bloomfilter =
formatted.find("row_groups_pruned_bloom_filter").unwrap();
let i_page_rows = formatted.find("page_index_rows_pruned").unwrap();
let i_page_pages = formatted.find("page_index_pages_pruned").unwrap();

assert!(
(i_file < i_rowgroup_stat)
&& (i_rowgroup_stat < i_rowgroup_bloomfilter)
&& (i_rowgroup_bloomfilter < i_page_pages && i_page_pages < i_page_rows),
"The parquet pruning metrics should be displayed in an order of: file range -> row group statistics -> row group bloom filter -> page index."
&& (i_rowgroup_stat < i_page_pages)
&& (i_page_pages < i_page_rows),
"The parquet pruning metrics should be displayed in an order of: file range -> row group statistics -> page index."
);
}

Expand Down Expand Up @@ -1179,7 +1174,7 @@ async fn parquet_explain_analyze_verbose() {
.to_string();

// should contain the raw per file stats (with the label)
assert_contains!(&formatted, "row_groups_pruned_bloom_filter{partition=0");
assert_not_contains!(&formatted, "row_groups_pruned_bloom_filter");
assert_contains!(&formatted, "row_groups_pruned_statistics{partition=0");
}

Expand Down
54 changes: 41 additions & 13 deletions datafusion/datasource-parquet/src/opener/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1336,26 +1336,14 @@ impl FiltersPreparedParquetOpen {
.row_groups_pruned_statistics
.add_matched(row_groups.remaining_row_group_count());
}

if !prepared.enable_bloom_filter || row_groups.is_empty() {
// Update metrics: bloom filter unavailable, so all row groups are
// matched (not pruned)
prepared
.file_metrics
.row_groups_pruned_bloom_filter
.add_matched(row_groups.remaining_row_group_count());
}
} else {
// Update metrics: no predicate, so all row groups are matched (not pruned)
// by statistics. Bloom pruning did not run, so its metrics are unchanged.
let remaining = row_groups.remaining_row_group_count();
prepared
.file_metrics
.row_groups_pruned_statistics
.add_matched(remaining);
prepared
.file_metrics
.row_groups_pruned_bloom_filter
.add_matched(remaining);
}

Ok(RowGroupsPrunedParquetOpen {
Expand Down Expand Up @@ -3413,6 +3401,32 @@ mod test {
.expect("row_groups_pruned_statistics metric is emitted")
}

/// Bloom pruning counters after a scan, as `(pruned, matched)`.
///
/// Direct `MetricsSet` lookup, not plan display: idle Bloom metrics are
/// omitted from displayed plans even when the counters remain registered.
///
/// Opening one file still registers this name more than once:
/// `prepare_open_file` creates `ParquetFileMetrics`, and each
/// `ParquetFileReaderFactory::create_reader` call (initial reader plus
/// Bloom replacement reader) creates another independent set. Sum the
/// counters the same way `MetricsSet::sum_by_name` aggregates pruning
/// metrics. Only the opener's set is incremented by Bloom pruning; the
/// reader copies stay at zero.
fn bloom_filter_pruning_metrics(metrics: &ExecutionPlanMetricsSet) -> (usize, usize) {
use datafusion_physical_plan::metrics::MetricValue;
match metrics
.clone_inner()
.sum_by_name("row_groups_pruned_bloom_filter")
{
Some(MetricValue::PruningMetrics {
pruning_metrics, ..
}) => (pruning_metrics.pruned(), pruning_metrics.matched()),
Some(_) => panic!("row_groups_pruned_bloom_filter is not a pruning metric"),
None => panic!("row_groups_pruned_bloom_filter metric is registered"),
}
}

#[tokio::test]
async fn test_prune_all_null_column_equality_from_file_statistics() {
// Regression: a column whose file statistics say every value is
Expand Down Expand Up @@ -4743,6 +4757,16 @@ mod test {
// contribute to this metric.
bloom_bytes.push(counter_metric_value(&metrics, "bytes_scanned"));
assert_eq!(collect_int32_values(stream).await, vec![1, 1, 1, 1]);
let (pruned, matched) = bloom_filter_pruning_metrics(&metrics);
assert_eq!(pruned, 0);
if stats_pruning {
// [1,1,1] is fully matched by statistics and skips Bloom.
// [0,1,2] still evaluates Bloom for `a = 1` and cannot prune.
assert_eq!(matched, 1);
} else {
// Statistics pruning is off, so both row groups evaluate Bloom.
assert_eq!(matched, 2);
}
}
assert!(bloom_bytes[1] > 0, "partial row group needs Bloom I/O");
assert!(
Expand Down Expand Up @@ -4781,6 +4805,10 @@ mod test {
.unwrap();
assert_eq!(counter_metric_value(&metrics, "bytes_scanned"), 0);
assert_eq!(collect_int32_values(stream).await, vec![1, 1, 1]);
let (pruned, matched) = bloom_filter_pruning_metrics(&metrics);
// Bloom reads and evaluation are skipped for an entirely fully matched file.
assert_eq!(pruned, 0);
assert_eq!(matched, 0);
}

#[test]
Expand Down
57 changes: 51 additions & 6 deletions datafusion/datasource-parquet/src/row_group_filter.rs
Original file line number Diff line number Diff line change
Expand Up @@ -472,7 +472,8 @@ impl RowGroupAccessPlanFilter {
}

if self.access_plan.is_fully_matched(idx) {
metrics.row_groups_pruned_bloom_filter.add_matched(1);
// Statistics already proved every row matches. Bloom did not
// evaluate this group, so do not record a Bloom outcome.
continue;
}

Expand All @@ -481,7 +482,6 @@ impl RowGroupAccessPlanFilter {
// evaluation in that case: it runs once per row group and can be expensive for wide
// predicates, a cost files written without bloom filters would pay for nothing.
if stats.is_empty() {
metrics.row_groups_pruned_bloom_filter.add_matched(1);
continue;
}

Expand All @@ -493,7 +493,7 @@ impl RowGroupAccessPlanFilter {
"Error evaluating row group predicate on bloom filter: {e}"
);
metrics.predicate_evaluation_errors.add(1);
false
continue;
}
};

Expand Down Expand Up @@ -622,9 +622,12 @@ mod tests {

use arrow::datatypes::DataType::Decimal128;
use arrow::datatypes::{DataType, Field};
use datafusion_expr::{cast, col, lit};
use datafusion_physical_expr::PhysicalExpr;
use datafusion_expr::{Operator, cast, col, lit};
use datafusion_physical_expr::planner::logical2physical;
use datafusion_physical_expr::{
PhysicalExpr,
expressions::{BinaryExpr, CastExpr, Column as PhysicalColumn, Literal},
};
use datafusion_physical_plan::metrics::ExecutionPlanMetricsSet;
use datafusion_pruning::PruningPredicateBuilder;
use parquet::arrow::ArrowSchemaConverter;
Expand Down Expand Up @@ -1640,7 +1643,49 @@ mod tests {
// the row group with a bloom filter is still evaluated and pruned.
assert_pruned(row_groups, ExpectedPruning::Some(vec![0]));
assert_eq!(metrics.row_groups_pruned_bloom_filter.pruned(), 1);
assert_eq!(metrics.row_groups_pruned_bloom_filter.matched(), 1);
assert_eq!(metrics.row_groups_pruned_bloom_filter.matched(), 0);
}

#[test]
fn bloom_filter_pruning_error_retains_row_group_without_match() {
let schema =
Arc::new(Schema::new(vec![Field::new("c1", DataType::Int32, false)]));
let cast_expr: Arc<dyn PhysicalExpr> = Arc::new(CastExpr::new(
Arc::new(PhysicalColumn::new("c1", 0)),
DataType::FixedSizeBinary(4),
None,
));
let literal: Arc<dyn PhysicalExpr> = Arc::new(Literal::new(
ScalarValue::FixedSizeBinary(4, Some(vec![0, 0, 0, 1])),
));
let expr = Arc::new(BinaryExpr::new(cast_expr, Operator::Eq, literal));
let pruning_predicate = PruningPredicateBuilder::new()
.with_file_schema(schema)
.try_build(expr)
.unwrap();

// A nonempty map bypasses the no-filter path. Evaluating the predicate
// then fails because Arrow cannot cast Int32 statistics to FixedSizeBinary.
let mut sbbf = Sbbf::new_with_ndv_fpp(10, 0.01).unwrap();
sbbf.insert(&1_i32);
let mut bloom_statistics = BloomFilterStatistics::new();
bloom_statistics.insert("loaded_filter", sbbf, PhysicalType::INT32, 4);
assert!(pruning_predicate.prune(&bloom_statistics).is_err());

let metrics = parquet_file_metrics();
let mut row_groups = RowGroupAccessPlanFilter::new(ParquetAccessPlan::new_all(1));
row_groups.prune_by_bloom_filters(
&pruning_predicate,
&metrics,
&[bloom_statistics],
);

// An evaluation failure is conservative: retain the row group, report
// the error separately, and do not call the result a Bloom match.
assert_pruned(row_groups, ExpectedPruning::Some(vec![0]));
assert_eq!(metrics.predicate_evaluation_errors.value(), 1);
assert_eq!(metrics.row_groups_pruned_bloom_filter.pruned(), 0);
assert_eq!(metrics.row_groups_pruned_bloom_filter.matched(), 0);
}

fn get_row_group_meta_data(
Expand Down
Loading
Loading