From e3519868c9ead026d8c116796bc64930ca528433 Mon Sep 17 00:00:00 2001 From: yashrb24 Date: Fri, 18 Sep 2026 13:19:08 +0530 Subject: [PATCH 1/3] feat: derive per-column pruning guarantees from tuple IN lists --- .../datasource-parquet/src/bloom_filter.rs | 108 ++++++++++ datafusion/functions/src/core/struct.rs | 18 +- .../physical-expr/src/utils/guarantee.rs | 189 +++++++++++++++++- .../test_files/push_down_filter_parquet.slt | 6 +- 4 files changed, 312 insertions(+), 9 deletions(-) diff --git a/datafusion/datasource-parquet/src/bloom_filter.rs b/datafusion/datasource-parquet/src/bloom_filter.rs index 2384f95df9885..71c5a0f194412 100644 --- a/datafusion/datasource-parquet/src/bloom_filter.rs +++ b/datafusion/datasource-parquet/src/bloom_filter.rs @@ -270,6 +270,114 @@ mod tests { .unwrap() } + #[tokio::test] + async fn test_tuple_in_bloom_pruning_preserves_correlation() -> Result<()> { + use arrow::array::{Int32Array, RecordBatch, StructArray, as_boolean_array}; + use datafusion_physical_expr::PhysicalExpr; + use datafusion_physical_expr::expressions::{ + DynamicFilterPhysicalExpr, InListExpr, + }; + use parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder; + + let schema = Arc::new(Schema::new(vec![ + Field::new("a", DataType::Int32, false), + Field::new("b", DataType::Int32, false), + ])); + // Row groups contain exact matches, impossible values, and crossed pairs. + let batch = RecordBatch::try_new( + Arc::clone(&schema), + vec![ + Arc::new(Int32Array::from(vec![1, 2, 99, 98, 1, 2])), + Arc::new(Int32Array::from(vec![10, 20, 99, 98, 20, 10])), + ], + )?; + let tuple = logical2physical( + &datafusion_functions::core::r#struct().call(vec![col("a"), col("b")]), + &schema, + ); + let columns = tuple.children().into_iter().cloned().collect(); + let DataType::Struct(fields) = tuple.data_type(&schema)? else { + unreachable!() + }; + let values = Arc::new(StructArray::new( + fields, + vec![ + Arc::new(Int32Array::from(vec![1, 2])), + Arc::new(Int32Array::from(vec![10, 20])), + ], + None, + )); + let exact: Arc = Arc::new(InListExpr::try_new_from_array( + tuple, values, false, &schema, + )?); + let live = Arc::new(DynamicFilterPhysicalExpr::new( + columns, + datafusion_physical_expr::expressions::lit(true), + )); + assert!( + build_test_pruning_predicate(live.clone(), schema.as_ref().clone()) + .literal_guarantees() + .is_empty() + ); + live.update(Arc::clone(&exact))?; + + for bloom_enabled in [true, false] { + let props = WriterProperties::builder() + .set_max_row_group_row_count(Some(2)) + .set_bloom_filter_fpp(0.000001) + .set_bloom_filter_enabled(bloom_enabled) + .build(); + let mut writer = + ArrowWriter::try_new(Vec::new(), Arc::clone(&schema), Some(props))?; + writer.write(&batch)?; + let data = bytes::Bytes::from(writer.into_inner()?); + let predicate = + build_test_pruning_predicate(live.clone(), schema.as_ref().clone()); + let pruned = test_row_group_bloom_filter_pruning_predicate( + "tuple-in.parquet", + data.clone(), + &predicate, + ) + .await?; + let selected = pruned.access_plan().row_group_indexes(); + assert_eq!( + selected, + if bloom_enabled { + vec![0, 2] + } else { + vec![0, 1, 2] + } + ); + let reader = ParquetRecordBatchReaderBuilder::try_new(data)? + .with_row_groups(selected) + .build()?; + let mut actual = Vec::new(); + let mut decoded = 0; + for batch in reader { + let batch = batch?; + decoded += batch.num_rows(); + let mask = exact.evaluate(&batch)?.into_array(batch.num_rows())?; + actual.push(arrow::compute::filter_record_batch( + &batch, + as_boolean_array(&mask), + )?); + } + assert_eq!(decoded, if bloom_enabled { 4 } else { 6 }); + datafusion_common::assert_batches_eq!( + [ + "+---+----+", + "| a | b |", + "+---+----+", + "| 1 | 10 |", + "| 2 | 20 |", + "+---+----+" + ], + &actual + ); + } + Ok(()) + } + #[tokio::test] async fn test_row_group_bloom_filter_pruning_predicate_simple_expr() { BloomFilterTest::new_data_index_bloom_encoding_stats() diff --git a/datafusion/functions/src/core/struct.rs b/datafusion/functions/src/core/struct.rs index 164b9d2032f4c..d689f3aac493a 100644 --- a/datafusion/functions/src/core/struct.rs +++ b/datafusion/functions/src/core/struct.rs @@ -15,11 +15,13 @@ // specific language governing permissions and limitations // under the License. +use super::getfield::GetFieldFunc; use arrow::array::StructArray; use arrow::datatypes::{DataType, Field, FieldRef}; -use datafusion_common::{Result, exec_err, internal_err}; +use datafusion_common::{Result, ScalarValue, exec_err, internal_err}; use datafusion_expr::{ - ColumnarValue, Documentation, ReturnFieldArgs, ScalarFunctionArgs, + ColumnarValue, Documentation, ReturnFieldArgs, ScalarFunctionArgs, ScalarUDF, + StructFieldMapping, }; use datafusion_expr::{ScalarUDFImpl, Signature, Volatility}; use datafusion_macros::user_doc; @@ -150,4 +152,16 @@ impl ScalarUDFImpl for StructFunc { fn documentation(&self) -> Option<&Documentation> { self.doc() } + + fn struct_field_mapping( + &self, + literal_args: &[Option], + ) -> Option { + Some(StructFieldMapping { + field_accessor: Arc::new(ScalarUDF::from(GetFieldFunc::new())), + fields: (0..literal_args.len()) + .map(|i| (vec![ScalarValue::Utf8(Some(format!("c{i}")))], i)) + .collect(), + }) + } } diff --git a/datafusion/physical-expr/src/utils/guarantee.rs b/datafusion/physical-expr/src/utils/guarantee.rs index 8b870c573d393..7ec3f6042c653 100644 --- a/datafusion/physical-expr/src/utils/guarantee.rs +++ b/datafusion/physical-expr/src/utils/guarantee.rs @@ -20,8 +20,9 @@ use crate::utils::split_disjunction; use crate::{PhysicalExpr, split_conjunction}; +use arrow::array::{Array, RecordBatch}; use datafusion_common::{Column, HashMap, ScalarValue}; -use datafusion_expr::Operator; +use datafusion_expr::{Operator, Volatility}; use std::collections::HashSet; use std::fmt::{self, Display, Formatter}; use std::sync::Arc; @@ -135,6 +136,16 @@ impl LiteralGuarantee { inlist.guarantee, inlist.list.iter().map(|lit| lit.value()), ) + } else if let Some(projected) = project_struct_in_list(inlist) { + projected + .into_iter() + .fold(builder, |builder, (col, values)| { + builder.aggregate_multi_conjunct( + col, + Guarantee::In, + &values, + ) + }) } else { builder } @@ -310,11 +321,11 @@ impl<'a> GuaranteeBuilder<'a> { /// * `AND (a != 1 OR a != 2 OR a != 3)`: a is not in (1, 2, or 3) /// * `AND (a NOT IN (1,2,3))`: a is not in (1, 2, or 3) #[allow(clippy::allow_attributes, clippy::mutable_key_type)] // ScalarValue has interior mutability but is intentionally used as hash key - fn aggregate_multi_conjunct( + fn aggregate_multi_conjunct<'b>( mut self, col: &'a crate::expressions::Column, guarantee: Guarantee, - new_values: impl IntoIterator, + new_values: impl IntoIterator, ) -> Self { let key = (col, guarantee); if let Some(index) = self.map.get(&key) { @@ -377,6 +388,85 @@ impl<'a> GuaranteeBuilder<'a> { } } +/// Project necessary per-column guarantees; the original predicate retains tuple correlation. +fn project_struct_in_list( + inlist: &crate::expressions::InListExpr, +) -> Option)>> { + if inlist.negated() || inlist.is_empty() { + return None; + } + let expr = inlist.expr().downcast_ref::()?; + let literal_args = expr + .args() + .iter() + .map(|arg| { + arg.downcast_ref::() + .map(|lit| lit.value().clone()) + }) + .collect::>(); + let mapping = expr.fun().struct_field_mapping(&literal_args)?; + if mapping.field_accessor.signature().volatility != Volatility::Immutable { + return None; + } + let tuples = inlist + .list() + .iter() + .map(|value| { + let literal = value.downcast_ref::()?; + let ScalarValue::Struct(array) = literal.value() else { + return None; + }; + (array.len() == 1).then_some(literal.value()) + }) + .collect::>>()?; + // Null tuples cannot make a positive IN predicate true. + let tuples = ScalarValue::iter_to_array( + tuples.into_iter().filter(|tuple| !tuple.is_null()).cloned(), + ) + .ok()?; + let batch = RecordBatch::try_from_iter([("tuple", tuples)]).ok()?; + + let mut projected = Vec::new(); + for (accessor_args, source_index) in mapping.fields { + let column = expr + .args() + .get(source_index)? + .downcast_ref::()?; + let mut args: Vec> = + vec![Arc::new(crate::expressions::Column::new("tuple", 0))]; + args.extend(accessor_args.into_iter().map(crate::expressions::lit)); + let accessor = crate::ScalarFunctionExpr::try_new( + Arc::clone(&mapping.field_accessor), + args, + batch.schema_ref(), + Arc::new(expr.config_options().clone()), + ) + .ok()?; + let array = accessor + .evaluate(&batch) + .ok()? + .into_array_of_size(batch.num_rows()) + .ok()?; + let mut values = Vec::new(); + for index in 0..array.len() { + let value = ScalarValue::try_from_array(array.as_ref(), index).ok()?; + let value = match value { + ScalarValue::Dictionary(_, value) => *value, + value => value, + }; + if value.is_null() { + values.clear(); + break; + } + values.push(value); + } + if !values.is_empty() { + projected.push((column, values)); + } + } + Some(projected) +} + /// Represents a single `col [not]in literal` expression struct ColOpLit<'a> { col: &'a crate::expressions::Column, @@ -442,7 +532,6 @@ impl<'a> ColInList<'a> { /// /// Returns None otherwise fn try_new(inlist: &'a crate::expressions::InListExpr) -> Option { - // Only support single-column inlist currently, multi-column inlist is not supported let col = inlist.expr().downcast_ref::()?; let literals = inlist @@ -842,6 +931,98 @@ mod test { ); } + #[test] + fn test_struct_inlist_guarantees() { + use crate::expressions::InListExpr; + use arrow::array::{ArrayRef, Int32Array, StringArray, StructArray}; + use arrow::buffer::NullBuffer; + + let make_expr = |strings: ArrayRef, nulls: Option, negated| { + let schema = Schema::new(vec![ + Field::new("a", strings.data_type().clone(), true), + Field::new("b", DataType::Int32, true), + ]); + let expr = logical2physical( + &datafusion_functions::core::r#struct().call(vec![col("a"), col("b")]), + &schema, + ); + let DataType::Struct(fields) = expr.data_type(&schema).unwrap() else { + unreachable!() + }; + let values = Arc::new(StructArray::new( + fields, + vec![strings, Arc::new(Int32Array::from(vec![1, 2, 3]))], + nulls, + )); + Arc::new( + InListExpr::try_new_from_array(expr, values, negated, &schema).unwrap(), + ) as Arc + }; + let strings: ArrayRef = Arc::new(StringArray::from(vec!["foo", "foo", "bar"])); + assert_eq!( + LiteralGuarantee::analyze(&make_expr(Arc::clone(&strings), None, false)), + vec![ + in_guarantee("a", ["foo", "bar"]), + in_guarantee("b", [1, 2, 3]) + ] + ); + assert!( + LiteralGuarantee::analyze(&make_expr(Arc::clone(&strings), None, true)) + .is_empty() + ); + assert_eq!( + LiteralGuarantee::analyze(&make_expr( + strings, + Some(vec![false, true, true].into()), + false + )), + vec![in_guarantee("a", ["foo", "bar"]), in_guarantee("b", [2, 3])] + ); + + let strings = StringArray::from(vec![Some("foo"), None, Some("bar")]); + let dictionary = + DataType::Dictionary(Box::new(DataType::Int32), Box::new(DataType::Utf8)); + for data_type in [DataType::Utf8, dictionary] { + let strings = arrow::compute::cast(&strings, &data_type).unwrap(); + assert_eq!( + LiteralGuarantee::analyze(&make_expr(strings, None, false)), + vec![in_guarantee("b", [1, 2, 3])] + ); + } + let strings = arrow::compute::cast( + &StringArray::from(vec!["foo", "foo", "bar"]), + &DataType::Dictionary(Box::new(DataType::Int32), Box::new(DataType::Utf8)), + ) + .unwrap(); + assert_eq!( + LiteralGuarantee::analyze(&make_expr(strings, None, false)), + vec![ + in_guarantee("a", ["foo", "bar"]), + in_guarantee("b", [1, 2, 3]) + ] + ); + + // Named fields map to nonconsecutive arguments, in a different column order. + let tuple = RecordBatch::try_from_iter([ + ("right", Arc::new(Int32Array::from(vec![1])) as ArrayRef), + ("left", Arc::new(StringArray::from(vec!["foo"])) as ArrayRef), + ]) + .unwrap(); + let expr = datafusion_functions::core::named_struct().call(vec![ + lit("right"), + col("b"), + lit("left"), + col("a"), + ]); + test_analyze( + expr.in_list( + vec![lit(ScalarValue::Struct(Arc::new(StructArray::from(tuple))))], + false, + ), + vec![in_guarantee("a", ["foo"]), in_guarantee("b", [1])], + ); + } + #[test] fn test_inlist_conjunction() { // b IN (1, 2, 3) AND b IN (2, 3, 4) diff --git a/datafusion/sqllogictest/test_files/push_down_filter_parquet.slt b/datafusion/sqllogictest/test_files/push_down_filter_parquet.slt index 6774d4f3a01db..7a96850f9630c 100644 --- a/datafusion/sqllogictest/test_files/push_down_filter_parquet.slt +++ b/datafusion/sqllogictest/test_files/push_down_filter_parquet.slt @@ -389,7 +389,7 @@ FROM join_probe p INNER JOIN join_build AS build Plan with Metrics 01)HashJoinExec: mode=CollectLeft, join_type=Inner, on=[(a@0, a@0), (b@1, b@1)], projection=[a@3, b@4, c@2, e@5], metrics=[output_rows=2, output_batches=1, array_map_created_count=0, build_input_batches=1, build_input_rows=2, input_batches=1, input_rows=2, avg_fanout=100% (2/2), probe_hit_rate=100% (2/2)] 02)--DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/join_build.parquet]]}, projection=[a, b, c], file_type=parquet, metrics=[output_rows=2, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=0, pushdown_rows_pruned=0, predicate_cache_inner_records=0, predicate_cache_records=0, scan_efficiency_ratio=19.58% (196/1.00 K)] -03)--DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/join_probe.parquet]]}, projection=[a, b, e], file_type=parquet, predicate=DynamicFilter [ a@0 >= aa AND a@0 <= ab AND b@1 >= ba AND b@1 <= bb AND struct(a@0, b@1) IN (SET) ([{c0:aa,c1:ba}, {c0:ab,c1:bb}]) ], dynamic_rg_pruning=eligible, pruning_predicate=a_null_count@1 != row_count@2 AND a_max@0 >= aa AND a_null_count@1 != row_count@2 AND a_min@3 <= ab AND b_null_count@5 != row_count@2 AND b_max@4 >= ba AND b_null_count@5 != row_count@2 AND b_min@6 <= bb, required_guarantees=[], metrics=[output_rows=2, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=2, pushdown_rows_pruned=2, predicate_cache_inner_records=8, predicate_cache_records=4, scan_efficiency_ratio=22.05% (228/1.03 K)] +03)--DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/join_probe.parquet]]}, projection=[a, b, e], file_type=parquet, predicate=DynamicFilter [ a@0 >= aa AND a@0 <= ab AND b@1 >= ba AND b@1 <= bb AND struct(a@0, b@1) IN (SET) ([{c0:aa,c1:ba}, {c0:ab,c1:bb}]) ], dynamic_rg_pruning=eligible, pruning_predicate=a_null_count@1 != row_count@2 AND a_max@0 >= aa AND a_null_count@1 != row_count@2 AND a_min@3 <= ab AND b_null_count@5 != row_count@2 AND b_max@4 >= ba AND b_null_count@5 != row_count@2 AND b_min@6 <= bb, required_guarantees=[a in (aa, ab), b in (ba, bb)], metrics=[output_rows=2, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=2, pushdown_rows_pruned=2, predicate_cache_inner_records=8, predicate_cache_records=4, scan_efficiency_ratio=22.05% (228/1.03 K)] statement ok reset datafusion.explain.analyze_categories; @@ -874,7 +874,7 @@ ON lj_build.a = lj_probe.a AND lj_build.b = lj_probe.b; Plan with Metrics 01)HashJoinExec: mode=CollectLeft, join_type=Left, on=[(a@0, a@0), (b@1, b@1)], metrics=[output_rows=2, output_batches=1, array_map_created_count=0, build_input_batches=1, build_input_rows=2, input_batches=2, input_rows=2, avg_fanout=100% (2/2), probe_hit_rate=100% (2/2)] 02)--DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/lj_build.parquet]]}, projection=[a, b, c], file_type=parquet, metrics=[output_rows=2, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=0, pushdown_rows_pruned=0, predicate_cache_inner_records=0, predicate_cache_records=0, scan_efficiency_ratio=19.58% (196/1.00 K)] -03)--DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/lj_probe.parquet]]}, projection=[a, b, e], file_type=parquet, predicate=DynamicFilter [ a@0 >= aa AND a@0 <= ab AND b@1 >= ba AND b@1 <= bb AND struct(a@0, b@1) IN (SET) ([{c0:aa,c1:ba}, {c0:ab,c1:bb}]) ], dynamic_rg_pruning=eligible, pruning_predicate=a_null_count@1 != row_count@2 AND a_max@0 >= aa AND a_null_count@1 != row_count@2 AND a_min@3 <= ab AND b_null_count@5 != row_count@2 AND b_max@4 >= ba AND b_null_count@5 != row_count@2 AND b_min@6 <= bb, required_guarantees=[], metrics=[output_rows=2, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=2, pushdown_rows_pruned=2, predicate_cache_inner_records=8, predicate_cache_records=4, scan_efficiency_ratio=22.05% (228/1.03 K)] +03)--DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/lj_probe.parquet]]}, projection=[a, b, e], file_type=parquet, predicate=DynamicFilter [ a@0 >= aa AND a@0 <= ab AND b@1 >= ba AND b@1 <= bb AND struct(a@0, b@1) IN (SET) ([{c0:aa,c1:ba}, {c0:ab,c1:bb}]) ], dynamic_rg_pruning=eligible, pruning_predicate=a_null_count@1 != row_count@2 AND a_max@0 >= aa AND a_null_count@1 != row_count@2 AND a_min@3 <= ab AND b_null_count@5 != row_count@2 AND b_max@4 >= ba AND b_null_count@5 != row_count@2 AND b_min@6 <= bb, required_guarantees=[a in (aa, ab), b in (ba, bb)], metrics=[output_rows=2, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=2, pushdown_rows_pruned=2, predicate_cache_inner_records=8, predicate_cache_records=4, scan_efficiency_ratio=22.05% (228/1.03 K)] # LEFT SEMI JOIN: only matching build rows are returned; probe scan still # receives the dynamic filter. @@ -890,7 +890,7 @@ WHERE EXISTS ( Plan with Metrics 01)HashJoinExec: mode=CollectLeft, join_type=LeftSemi, on=[(a@0, a@0), (b@1, b@1)], metrics=[output_rows=2, output_batches=1, array_map_created_count=0, build_input_batches=1, build_input_rows=2, input_batches=2, input_rows=4, avg_fanout=100% (2/2), probe_hit_rate=100% (2/2)] 02)--DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/lj_build.parquet]]}, projection=[a, b, c], file_type=parquet, metrics=[output_rows=2, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=0, pushdown_rows_pruned=0, predicate_cache_inner_records=0, predicate_cache_records=0, scan_efficiency_ratio=19.58% (196/1.00 K)] -03)--DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/lj_probe.parquet]]}, projection=[a, b], file_type=parquet, predicate=DynamicFilter [ a@0 >= aa AND a@0 <= ab AND b@1 >= ba AND b@1 <= bb AND struct(a@0, b@1) IN (SET) ([{c0:aa,c1:ba}, {c0:ab,c1:bb}]) ], dynamic_rg_pruning=eligible, pruning_predicate=a_null_count@1 != row_count@2 AND a_max@0 >= aa AND a_null_count@1 != row_count@2 AND a_min@3 <= ab AND b_null_count@5 != row_count@2 AND b_max@4 >= ba AND b_null_count@5 != row_count@2 AND b_min@6 <= bb, required_guarantees=[], metrics=[output_rows=2, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=2, pushdown_rows_pruned=2, predicate_cache_inner_records=8, predicate_cache_records=4, scan_efficiency_ratio=14.89% (154/1.03 K)] +03)--DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/lj_probe.parquet]]}, projection=[a, b], file_type=parquet, predicate=DynamicFilter [ a@0 >= aa AND a@0 <= ab AND b@1 >= ba AND b@1 <= bb AND struct(a@0, b@1) IN (SET) ([{c0:aa,c1:ba}, {c0:ab,c1:bb}]) ], dynamic_rg_pruning=eligible, pruning_predicate=a_null_count@1 != row_count@2 AND a_max@0 >= aa AND a_null_count@1 != row_count@2 AND a_min@3 <= ab AND b_null_count@5 != row_count@2 AND b_max@4 >= ba AND b_null_count@5 != row_count@2 AND b_min@6 <= bb, required_guarantees=[a in (aa, ab), b in (ba, bb)], metrics=[output_rows=2, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=2, pushdown_rows_pruned=2, predicate_cache_inner_records=8, predicate_cache_records=4, scan_efficiency_ratio=14.89% (154/1.03 K)] statement ok reset datafusion.explain.analyze_categories; From cde900735ec3a879ca619c61b8de94a1aef34c8a Mon Sep 17 00:00:00 2001 From: yashrb24 Date: Fri, 18 Sep 2026 17:46:34 +0530 Subject: [PATCH 2/3] test: simplify tuple IN Bloom pruning regression --- .../datasource-parquet/src/bloom_filter.rs | 94 ++++++++----------- 1 file changed, 41 insertions(+), 53 deletions(-) diff --git a/datafusion/datasource-parquet/src/bloom_filter.rs b/datafusion/datasource-parquet/src/bloom_filter.rs index 71c5a0f194412..96dab54835d28 100644 --- a/datafusion/datasource-parquet/src/bloom_filter.rs +++ b/datafusion/datasource-parquet/src/bloom_filter.rs @@ -321,60 +321,48 @@ mod tests { ); live.update(Arc::clone(&exact))?; - for bloom_enabled in [true, false] { - let props = WriterProperties::builder() - .set_max_row_group_row_count(Some(2)) - .set_bloom_filter_fpp(0.000001) - .set_bloom_filter_enabled(bloom_enabled) - .build(); - let mut writer = - ArrowWriter::try_new(Vec::new(), Arc::clone(&schema), Some(props))?; - writer.write(&batch)?; - let data = bytes::Bytes::from(writer.into_inner()?); - let predicate = - build_test_pruning_predicate(live.clone(), schema.as_ref().clone()); - let pruned = test_row_group_bloom_filter_pruning_predicate( - "tuple-in.parquet", - data.clone(), - &predicate, - ) - .await?; - let selected = pruned.access_plan().row_group_indexes(); - assert_eq!( - selected, - if bloom_enabled { - vec![0, 2] - } else { - vec![0, 1, 2] - } - ); - let reader = ParquetRecordBatchReaderBuilder::try_new(data)? - .with_row_groups(selected) - .build()?; - let mut actual = Vec::new(); - let mut decoded = 0; - for batch in reader { - let batch = batch?; - decoded += batch.num_rows(); - let mask = exact.evaluate(&batch)?.into_array(batch.num_rows())?; - actual.push(arrow::compute::filter_record_batch( - &batch, - as_boolean_array(&mask), - )?); - } - assert_eq!(decoded, if bloom_enabled { 4 } else { 6 }); - datafusion_common::assert_batches_eq!( - [ - "+---+----+", - "| a | b |", - "+---+----+", - "| 1 | 10 |", - "| 2 | 20 |", - "+---+----+" - ], - &actual - ); + let props = WriterProperties::builder() + .set_max_row_group_row_count(Some(2)) + .set_bloom_filter_fpp(0.000001) + .set_bloom_filter_enabled(true) + .build(); + let mut writer = + ArrowWriter::try_new(Vec::new(), Arc::clone(&schema), Some(props))?; + writer.write(&batch)?; + let data = bytes::Bytes::from(writer.into_inner()?); + let predicate = + build_test_pruning_predicate(live.clone(), schema.as_ref().clone()); + let pruned = test_row_group_bloom_filter_pruning_predicate( + "tuple-in.parquet", + data.clone(), + &predicate, + ) + .await?; + let selected = pruned.access_plan().row_group_indexes(); + assert_eq!(selected, vec![0, 2]); + let reader = ParquetRecordBatchReaderBuilder::try_new(data)? + .with_row_groups(selected) + .build()?; + let mut actual = Vec::new(); + for batch in reader { + let batch = batch?; + let mask = exact.evaluate(&batch)?.into_array(batch.num_rows())?; + actual.push(arrow::compute::filter_record_batch( + &batch, + as_boolean_array(&mask), + )?); } + datafusion_common::assert_batches_eq!( + [ + "+---+----+", + "| a | b |", + "+---+----+", + "| 1 | 10 |", + "| 2 | 20 |", + "+---+----+" + ], + &actual + ); Ok(()) } From 879fc8779254d67e1105cd4a56bf133dc4eb6864 Mon Sep 17 00:00:00 2001 From: yashrb24 Date: Mon, 21 Sep 2026 13:11:53 +0530 Subject: [PATCH 3/3] test: cover tuple Bloom pruning through a Parquet scan --- .../datasource-parquet/src/bloom_filter.rs | 72 ------------------- .../test_files/push_down_filter_parquet.slt | 21 ++++-- 2 files changed, 15 insertions(+), 78 deletions(-) diff --git a/datafusion/datasource-parquet/src/bloom_filter.rs b/datafusion/datasource-parquet/src/bloom_filter.rs index e59f40a1904d5..2384f95df9885 100644 --- a/datafusion/datasource-parquet/src/bloom_filter.rs +++ b/datafusion/datasource-parquet/src/bloom_filter.rs @@ -270,78 +270,6 @@ mod tests { .unwrap() } - #[tokio::test] - async fn test_tuple_in_bloom_pruning() -> Result<()> { - use arrow::array::{Int32Array, RecordBatch, StructArray}; - use datafusion_physical_expr::PhysicalExpr; - use datafusion_physical_expr::expressions::{ - DynamicFilterPhysicalExpr, InListExpr, - }; - - let schema = Arc::new(Schema::new(vec![ - Field::new("a", DataType::Int32, false), - Field::new("b", DataType::Int32, false), - ])); - // Row groups contain exact matches, impossible values, and crossed pairs. - let batch = RecordBatch::try_new( - Arc::clone(&schema), - vec![ - Arc::new(Int32Array::from(vec![1, 2, 99, 98, 1, 2])), - Arc::new(Int32Array::from(vec![10, 20, 99, 98, 20, 10])), - ], - )?; - let tuple = logical2physical( - &datafusion_functions::core::r#struct().call(vec![col("a"), col("b")]), - &schema, - ); - let columns = tuple.children().into_iter().cloned().collect(); - let DataType::Struct(fields) = tuple.data_type(&schema)? else { - unreachable!() - }; - let values = Arc::new(StructArray::new( - fields, - vec![ - Arc::new(Int32Array::from(vec![1, 2])), - Arc::new(Int32Array::from(vec![10, 20])), - ], - None, - )); - let exact: Arc = Arc::new(InListExpr::try_new_from_array( - tuple, values, false, &schema, - )?); - let live = Arc::new(DynamicFilterPhysicalExpr::new( - columns, - datafusion_physical_expr::expressions::lit(true), - )); - assert!( - build_test_pruning_predicate(live.clone(), schema.as_ref().clone()) - .literal_guarantees() - .is_empty() - ); - live.update(exact)?; - - let props = WriterProperties::builder() - .set_max_row_group_row_count(Some(2)) - .set_bloom_filter_fpp(0.000001) - .set_bloom_filter_enabled(true) - .build(); - let mut writer = - ArrowWriter::try_new(Vec::new(), Arc::clone(&schema), Some(props))?; - writer.write(&batch)?; - let data = bytes::Bytes::from(writer.into_inner()?); - let predicate = build_test_pruning_predicate(live, schema.as_ref().clone()); - let pruned = test_row_group_bloom_filter_pruning_predicate( - "tuple-in.parquet", - data, - &predicate, - ) - .await?; - let selected = pruned.access_plan().row_group_indexes(); - // Per-column Bloom checks cannot exclude the crossed pairs in group 2. - assert_eq!(selected, vec![0, 2]); - Ok(()) - } - #[tokio::test] async fn test_row_group_bloom_filter_pruning_predicate_simple_expr() { BloomFilterTest::new_data_index_bloom_encoding_stats() diff --git a/datafusion/sqllogictest/test_files/push_down_filter_parquet.slt b/datafusion/sqllogictest/test_files/push_down_filter_parquet.slt index 6fc4db203453c..caaa2856519b6 100644 --- a/datafusion/sqllogictest/test_files/push_down_filter_parquet.slt +++ b/datafusion/sqllogictest/test_files/push_down_filter_parquet.slt @@ -346,24 +346,33 @@ COPY ( ) TO 'test_files/scratch/push_down_filter_parquet/join_build.parquet' STORED AS PARQUET; +# Three row groups: exact matches, Bloom-absent values within the min/max ranges, +# and crossed pairs that survive per-column pruning but fail the tuple filter. statement ok COPY ( SELECT * FROM (VALUES ('aa', 'ba', 1.0), ('ab', 'bb', 2.0), - ('ac', 'bc', 3.0), - ('ad', 'bd', 4.0) + ('aa', 'baz', 3.0), + ('ab', 'baz', 4.0), + ('aa', 'bb', 5.0), + ('ab', 'ba', 6.0) ) AS v(a, b, e) ) TO 'test_files/scratch/push_down_filter_parquet/join_probe.parquet' -STORED AS PARQUET; +STORED AS PARQUET +OPTIONS ( + 'format.max_row_group_size' '2', + 'format.bloom_filter_on_write' 'true', + 'format.bloom_filter_fpp' '0.000001' +); statement ok -CREATE EXTERNAL TABLE join_build (a VARCHAR, b VARCHAR, c DOUBLE) +CREATE EXTERNAL TABLE join_build (a VARCHAR NOT NULL, b VARCHAR NOT NULL, c DOUBLE) STORED AS PARQUET LOCATION 'test_files/scratch/push_down_filter_parquet/join_build.parquet'; statement ok -CREATE EXTERNAL TABLE join_probe (a VARCHAR, b VARCHAR, e DOUBLE) +CREATE EXTERNAL TABLE join_probe (a VARCHAR NOT NULL, b VARCHAR NOT NULL, e DOUBLE) STORED AS PARQUET LOCATION 'test_files/scratch/push_down_filter_parquet/join_probe.parquet'; @@ -389,7 +398,7 @@ FROM join_probe p INNER JOIN join_build AS build Plan with Metrics 01)HashJoinExec: mode=CollectLeft, join_type=Inner, on=[(a@0, a@0), (b@1, b@1)], projection=[a@3, b@4, c@2, e@5], metrics=[output_rows=2, output_batches=1, array_map_created_count=0, build_input_batches=1, build_input_rows=2, input_batches=1, input_rows=2, avg_fanout=100% (2/2), probe_hit_rate=100% (2/2)] 02)--DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/join_build.parquet]]}, projection=[a, b, c], file_type=parquet, metrics=[output_rows=2, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=0, pushdown_rows_pruned=0, predicate_cache_inner_records=0, predicate_cache_records=0, scan_efficiency_ratio=19.88% (196/986)] -03)--DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/join_probe.parquet]]}, projection=[a, b, e], file_type=parquet, predicate=DynamicFilter [ a@0 >= aa AND a@0 <= ab AND b@1 >= ba AND b@1 <= bb AND struct(a@0, b@1) IN (SET) ([{c0:aa,c1:ba}, {c0:ab,c1:bb}]) ], dynamic_rg_pruning=eligible, pruning_predicate=a_null_count@1 != row_count@2 AND a_max@0 >= aa AND a_null_count@1 != row_count@2 AND a_min@3 <= ab AND b_null_count@5 != row_count@2 AND b_max@4 >= ba AND b_null_count@5 != row_count@2 AND b_min@6 <= bb, required_guarantees=[a in (aa, ab), b in (ba, bb)], metrics=[output_rows=2, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=2, pushdown_rows_pruned=2, predicate_cache_inner_records=8, predicate_cache_records=4, scan_efficiency_ratio=22.37% (228/1.02 K)] +03)--DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/join_probe.parquet]]}, projection=[a, b, e], file_type=parquet, predicate=DynamicFilter [ a@0 >= aa AND a@0 <= ab AND b@1 >= ba AND b@1 <= bb AND struct(a@0, b@1) IN (SET) ([{c0:aa,c1:ba}, {c0:ab,c1:bb}]) ], dynamic_rg_pruning=eligible, pruning_predicate=a_null_count@1 != row_count@2 AND a_max@0 >= aa AND a_null_count@1 != row_count@2 AND a_min@3 <= ab AND b_null_count@5 != row_count@2 AND b_max@4 >= ba AND b_null_count@5 != row_count@2 AND b_min@6 <= bb, required_guarantees=[a in (aa, ab), b in (ba, bb)], metrics=[output_rows=2, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=3 total → 3 matched, row_groups_pruned_bloom_filter=3 total → 2 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=2, pushdown_rows_pruned=2, predicate_cache_inner_records=4, predicate_cache_records=4, scan_efficiency_ratio=23.61% (606/2.57 K)] statement ok reset datafusion.explain.analyze_categories;