Skip to content

Commit f427fb7

Browse files
authored
perf(parquet): skip Bloom reads for fully matched row groups (#25854)
## Which issue does this PR close? No linked issue. This is a follow-up to the fully matched Parquet row-group work in PR #23696. ## Rationale for this change When row-group statistics prove that every row satisfies a scan predicate, a Bloom filter cannot prune that group. Reading its Bloom filters still adds object-store I/O, which can be especially costly for remote files. ## What changes are included in this PR? - Skip Bloom filter reads and predicate evaluation for fully matched row groups. - Avoid creating a Bloom reader when every surviving row group is fully matched. - Keep the Bloom pruning matched metric accounting for skipped groups. ## What is the testing strategy for this PR? The new `fully_matched_row_groups_skip_bloom_filter_reads` test verifies that a partially matched group still reads Bloom filters, a fully matched group reduces `bytes_scanned`, and an all-fully-matched file reads zero Bloom bytes during open. It also checks that the returned rows are unchanged. ## Are there any user-facing changes? No API or query-result changes. Scans avoid unnecessary Bloom filter reads for fully matched row groups.
1 parent 1d9be2e commit f427fb7

2 files changed

Lines changed: 99 additions & 0 deletions

File tree

‎datafusion/datasource-parquet/src/opener/mod.rs‎

Lines changed: 94 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1482,6 +1482,10 @@ impl RowGroupsPrunedParquetOpen {
14821482
self.prepared.pruning_predicate.as_ref().map(|p| p.as_ref())
14831483
&& self.prepared.loaded.prepared.enable_bloom_filter
14841484
&& !self.row_groups.is_empty()
1485+
&& self
1486+
.row_groups
1487+
.row_group_indexes()
1488+
.any(|idx| !self.row_groups.access_plan().is_fully_matched(idx))
14851489
{
14861490
// Use the existing reader for bloom filter I/O;
14871491
// replace with a fresh reader for decoding below.
@@ -1521,6 +1525,11 @@ impl RowGroupsPrunedParquetOpen {
15211525
.collect();
15221526

15231527
for idx in self.row_groups.row_group_indexes() {
1528+
// Statistics have already proved that every row in this group
1529+
// matches the predicate, so its Bloom filters cannot prune it.
1530+
if self.row_groups.access_plan().is_fully_matched(idx) {
1531+
continue;
1532+
}
15241533
let mut row_group_filters =
15251534
BloomFilterStatistics::with_capacity(parquet_columns.len());
15261535
for (column_name, column_idx, physical_type, type_length) in
@@ -2485,6 +2494,11 @@ mod test {
24852494
self
24862495
}
24872496

2497+
fn with_enable_bloom_filter(mut self, enable: bool) -> Self {
2498+
self.enable_bloom_filter = enable;
2499+
self
2500+
}
2501+
24882502
fn with_metrics(mut self, metrics: ExecutionPlanMetricsSet) -> Self {
24892503
self.metrics = metrics;
24902504
self
@@ -4570,6 +4584,86 @@ mod test {
45704584
);
45714585
}
45724586

4587+
#[tokio::test]
4588+
async fn fully_matched_row_groups_skip_bloom_filter_reads() {
4589+
let store = Arc::new(InMemory::new()) as Arc<dyn ObjectStore>;
4590+
let batches = vec![
4591+
record_batch!(("a", Int32, vec![1, 1, 1])).unwrap(),
4592+
record_batch!(("a", Int32, vec![0, 1, 2])).unwrap(),
4593+
];
4594+
let schema = batches[0].schema();
4595+
let props = WriterProperties::builder()
4596+
.set_max_row_group_row_count(Some(3))
4597+
.set_bloom_filter_enabled(true)
4598+
.set_statistics_enabled(EnabledStatistics::Chunk)
4599+
.build();
4600+
let data_size = write_parquet_batches(
4601+
Arc::clone(&store),
4602+
"bloom.parquet",
4603+
batches,
4604+
Some(props),
4605+
)
4606+
.await;
4607+
let file = PartitionedFile::new("bloom.parquet".to_string(), data_size as u64);
4608+
let predicate = logical2physical(&col("a").eq(lit(1)), &schema);
4609+
4610+
let mut bloom_bytes = Vec::new();
4611+
for stats_pruning in [false, true] {
4612+
let metrics = ExecutionPlanMetricsSet::new();
4613+
let morselizer = ParquetMorselizerBuilder::new()
4614+
.with_store(Arc::clone(&store))
4615+
.with_schema(Arc::clone(&schema))
4616+
.with_predicate(Arc::clone(&predicate))
4617+
.with_pushdown_filters(true)
4618+
.with_row_group_stats_pruning(stats_pruning)
4619+
.with_enable_bloom_filter(true)
4620+
.with_metrics(metrics.clone())
4621+
.build();
4622+
let stream = open_file(&morselizer, file.clone()).await.unwrap();
4623+
// The decoder has not been polled yet, so only Bloom filter reads
4624+
// contribute to this metric.
4625+
bloom_bytes.push(counter_metric_value(&metrics, "bytes_scanned"));
4626+
assert_eq!(collect_int32_values(stream).await, vec![1, 1, 1, 1]);
4627+
}
4628+
assert!(bloom_bytes[1] > 0, "partial row group needs Bloom I/O");
4629+
assert!(
4630+
bloom_bytes[1] < bloom_bytes[0],
4631+
"fully matched row group should not load Bloom filters: {bloom_bytes:?}"
4632+
);
4633+
4634+
let all_matched = record_batch!(("a", Int32, vec![1, 1, 1])).unwrap();
4635+
let data_size = write_parquet_batches(
4636+
Arc::clone(&store),
4637+
"all-matched.parquet",
4638+
vec![all_matched],
4639+
Some(
4640+
WriterProperties::builder()
4641+
.set_bloom_filter_enabled(true)
4642+
.set_statistics_enabled(EnabledStatistics::Chunk)
4643+
.build(),
4644+
),
4645+
)
4646+
.await;
4647+
let metrics = ExecutionPlanMetricsSet::new();
4648+
let morselizer = ParquetMorselizerBuilder::new()
4649+
.with_store(Arc::clone(&store))
4650+
.with_schema(Arc::clone(&schema))
4651+
.with_predicate(predicate)
4652+
.with_pushdown_filters(true)
4653+
.with_row_group_stats_pruning(true)
4654+
.with_enable_bloom_filter(true)
4655+
.with_metrics(metrics.clone())
4656+
.build();
4657+
let stream = open_file(
4658+
&morselizer,
4659+
PartitionedFile::new("all-matched.parquet".to_string(), data_size as u64),
4660+
)
4661+
.await
4662+
.unwrap();
4663+
assert_eq!(counter_metric_value(&metrics, "bytes_scanned"), 0);
4664+
assert_eq!(collect_int32_values(stream).await, vec![1, 1, 1]);
4665+
}
4666+
45734667
#[test]
45744668
fn should_load_page_index_with_row_selection() {
45754669
use parquet::arrow::arrow_reader::{RowSelection, RowSelector};

‎datafusion/datasource-parquet/src/row_group_filter.rs‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -451,6 +451,11 @@ impl RowGroupAccessPlanFilter {
451451
continue;
452452
}
453453

454+
if self.access_plan.is_fully_matched(idx) {
455+
metrics.row_groups_pruned_bloom_filter.add_matched(1);
456+
continue;
457+
}
458+
454459
// A row group without any loaded bloom filter cannot be pruned by this pass: with no
455460
// bloom statistics to consult, evaluation can only conclude "may match". Skip the
456461
// evaluation in that case: it runs once per row group and can be expensive for wide

0 commit comments

Comments
 (0)