@@ -61,7 +61,7 @@ use datafusion_common::{
6161use datafusion_datasource:: { PartitionedFile , TableSchema } ;
6262use datafusion_physical_expr:: expressions:: { Column , DynamicFilterTracking , Literal } ;
6363use datafusion_physical_expr:: simplifier:: PhysicalExprSimplifier ;
64- use datafusion_physical_expr:: utils:: collect_columns;
64+ use datafusion_physical_expr:: utils:: { collect_columns, split_optional } ;
6565use datafusion_physical_expr_adapter:: PhysicalExprAdapterFactory ;
6666use datafusion_physical_expr_common:: physical_expr:: PhysicalExpr ;
6767use datafusion_physical_expr_common:: sort_expr:: LexOrdering ;
@@ -554,7 +554,10 @@ impl DecoderReadPlans {
554554 // `RowFilter` machinery cannot evaluate on this file — the rejected
555555 // conjuncts returned by `RowFilterContext::try_new`).
556556 //
557- // Either way every conjunct is applied; nothing is silently dropped.
557+ // Either way every required conjunct is applied; nothing is silently
558+ // dropped. Optional conjuncts (see `split_optional`) are not needed
559+ // for correctness. They never go to the post-scan filter: they are
560+ // row filter predicates or they are used only for statistics pruning.
558561 // ---------------------------------------------------------------
559562 let ( row_filter_context, post_scan_conjuncts) =
560563 match ( prepared. pushdown_filters , prepared. predicate . as_ref ( ) ) {
@@ -572,15 +575,13 @@ impl DecoderReadPlans {
572575 prepared. file_metrics . clone ( ) ,
573576 prepared. max_predicate_cache_size ,
574577 ) ,
575- // Pushdown disabled: the whole predicate runs post-scan (in-scan
576- // equivalent of a `FilterExec`).
577- ( false , Some ( predicate) ) => (
578- None ,
579- datafusion_physical_expr:: split_conjunction ( predicate)
580- . into_iter ( )
581- . cloned ( )
582- . collect ( ) ,
583- ) ,
578+ // Pushdown disabled: the required conjuncts run post-scan
579+ // (in-scan equivalent of a `FilterExec`). Optional conjuncts
580+ // (for example hash join dynamic filters) are not needed for
581+ // correctness and are expensive to evaluate for each row,
582+ // thus they are used only for statistics pruning, as before
583+ // the scan accepted the filters.
584+ ( false , Some ( predicate) ) => ( None , split_optional ( predicate) . 0 ) ,
584585 ( _, None ) => ( None , Vec :: new ( ) ) ,
585586 } ;
586587
@@ -5384,13 +5385,101 @@ mod test {
53845385 /// post-scan filter, so only the rows with a non-null struct survive.
53855386 #[ tokio:: test]
53865387 async fn rejected_struct_conjunct_runs_post_scan_not_dropped ( ) {
5388+ let store = Arc :: new ( InMemory :: new ( ) ) as Arc < dyn ObjectStore > ;
5389+ let ( schema, file) = write_struct_file ( & store) . await ;
5390+
5391+ // `s IS NOT NULL` references a whole struct, which `PushdownChecker`
5392+ // flags as non-primitive — `FilterCandidateBuilder::build` returns
5393+ // `Ok(None)` and the conjunct lands in `rejected`.
5394+ let predicate = logical2physical ( & col ( "s" ) . is_not_null ( ) , & schema) ;
5395+
5396+ let morselizer = ParquetMorselizerBuilder :: new ( )
5397+ . with_store ( Arc :: clone ( & store) )
5398+ . with_schema ( Arc :: clone ( & schema) )
5399+ . with_predicate ( predicate)
5400+ // The RowFilter path: emulates the post-`try_pushdown_filters`
5401+ // state where the parent `FilterExec` has already been removed
5402+ // and the scan owns the conjunct.
5403+ . with_pushdown_filters ( true )
5404+ . build ( ) ;
5405+
5406+ let stream = open_file ( & morselizer, file) . await . unwrap ( ) ;
5407+ let ( _, rows) = count_batches_and_rows ( stream) . await ;
5408+
5409+ // 2 rows have a non-null struct. Before the fix this returned 3
5410+ // (the conjunct was silently dropped).
5411+ assert_eq ! (
5412+ rows, 2 ,
5413+ "expected 2 rows with non-null struct; the rejected conjunct must \
5414+ be applied post-scan, not silently dropped"
5415+ ) ;
5416+ }
5417+
5418+ /// An optional conjunct is not needed for correctness. The scan never
5419+ /// evaluates it after the decode: not when `pushdown_filters` is false,
5420+ /// and not when the `RowFilter` rejects it for the file. A required
5421+ /// conjunct in the same situations is evaluated after the decode.
5422+ #[ tokio:: test]
5423+ async fn optional_conjunct_is_never_evaluated_post_scan ( ) {
5424+ let store = Arc :: new ( InMemory :: new ( ) ) as Arc < dyn ObjectStore > ;
5425+ let ( schema, file) = write_struct_file ( & store) . await ;
5426+
5427+ // `s IS NOT NULL` is rejected by the `RowFilter` (whole struct).
5428+ // `id > 1` can be a `RowFilter` predicate.
5429+ let rejected = logical2physical ( & col ( "s" ) . is_not_null ( ) , & schema) ;
5430+ let pushable = logical2physical ( & col ( "id" ) . gt ( lit ( 1 ) ) , & schema) ;
5431+ let optional = |expr : & Arc < dyn PhysicalExpr > | -> Arc < dyn PhysicalExpr > {
5432+ Arc :: new (
5433+ datafusion_physical_expr:: expressions:: OptionalFilterPhysicalExpr :: new (
5434+ Arc :: clone ( expr) ,
5435+ ) ,
5436+ )
5437+ } ;
5438+
5439+ // (predicate, pushdown_filters, expected rows, expected post-scan rows)
5440+ let cases: Vec < ( Arc < dyn PhysicalExpr > , bool , usize , usize ) > = vec ! [
5441+ // Required conjuncts: applied post-scan (#22384).
5442+ ( Arc :: clone( & rejected) , false , 2 , 3 ) ,
5443+ ( Arc :: clone( & rejected) , true , 2 , 3 ) ,
5444+ ( Arc :: clone( & pushable) , false , 2 , 3 ) ,
5445+ // Optional conjuncts: never evaluated post-scan.
5446+ ( optional( & rejected) , false , 3 , 0 ) ,
5447+ ( optional( & rejected) , true , 3 , 0 ) ,
5448+ ( optional( & pushable) , false , 3 , 0 ) ,
5449+ // An optional conjunct that the `RowFilter` accepts is a row
5450+ // filter predicate.
5451+ ( optional( & pushable) , true , 2 , 0 ) ,
5452+ ] ;
5453+ for ( predicate, pushdown, expected_rows, expected_post_scan_rows) in cases {
5454+ let metrics = ExecutionPlanMetricsSet :: new ( ) ;
5455+ let morselizer = ParquetMorselizerBuilder :: new ( )
5456+ . with_store ( Arc :: clone ( & store) )
5457+ . with_schema ( Arc :: clone ( & schema) )
5458+ . with_predicate ( Arc :: clone ( & predicate) )
5459+ . with_pushdown_filters ( pushdown)
5460+ . with_metrics ( metrics. clone ( ) )
5461+ . build ( ) ;
5462+ let stream = open_file ( & morselizer, file. clone ( ) ) . await . unwrap ( ) ;
5463+ let ( _, rows) = count_batches_and_rows ( stream) . await ;
5464+ let post_scan_rows = counter_metric_value ( & metrics, "post_scan_rows_pruned" )
5465+ + counter_metric_value ( & metrics, "post_scan_rows_matched" ) ;
5466+ assert_eq ! (
5467+ ( rows, post_scan_rows) ,
5468+ ( expected_rows, expected_post_scan_rows) ,
5469+ "predicate {predicate}, pushdown_filters {pushdown}"
5470+ ) ;
5471+ }
5472+ }
5473+
5474+ /// Writes a file with the columns `id` (Int32: 1, 2, 3) and `s`
5475+ /// (Struct{value: Int32, label: Utf8}; row 1 is null).
5476+ async fn write_struct_file (
5477+ store : & Arc < dyn ObjectStore > ,
5478+ ) -> ( SchemaRef , PartitionedFile ) {
53875479 use arrow:: array:: { Int32Array , StringArray , StructArray } ;
53885480 use arrow:: buffer:: NullBuffer ;
53895481 use arrow:: datatypes:: Fields ;
53905482
5391- let store = Arc :: new ( InMemory :: new ( ) ) as Arc < dyn ObjectStore > ;
5392-
5393- // Schema: id (Int32), s (Struct{value: Int32, label: Utf8}).
53945483 let struct_fields: Fields = vec ! [
53955484 Arc :: new( Field :: new( "value" , DataType :: Int32 , true ) ) ,
53965485 Arc :: new( Field :: new( "label" , DataType :: Utf8 , true ) ) ,
@@ -5420,40 +5509,15 @@ mod test {
54205509 . unwrap ( ) ;
54215510
54225511 let data_size = write_parquet_batches (
5423- Arc :: clone ( & store) ,
5512+ Arc :: clone ( store) ,
54245513 "rejected.parquet" ,
54255514 vec ! [ batch] ,
54265515 None ,
54275516 )
54285517 . await ;
54295518
54305519 let file = PartitionedFile :: new ( "rejected.parquet" . to_string ( ) , data_size as u64 ) ;
5431-
5432- // `s IS NOT NULL` references a whole struct, which `PushdownChecker`
5433- // flags as non-primitive — `FilterCandidateBuilder::build` returns
5434- // `Ok(None)` and the conjunct lands in `rejected`.
5435- let predicate = logical2physical ( & col ( "s" ) . is_not_null ( ) , & schema) ;
5436-
5437- let morselizer = ParquetMorselizerBuilder :: new ( )
5438- . with_store ( Arc :: clone ( & store) )
5439- . with_schema ( Arc :: clone ( & schema) )
5440- . with_predicate ( predicate)
5441- // The RowFilter path: emulates the post-`try_pushdown_filters`
5442- // state where the parent `FilterExec` has already been removed
5443- // and the scan owns the conjunct.
5444- . with_pushdown_filters ( true )
5445- . build ( ) ;
5446-
5447- let stream = open_file ( & morselizer, file) . await . unwrap ( ) ;
5448- let ( _, rows) = count_batches_and_rows ( stream) . await ;
5449-
5450- // 2 rows have a non-null struct. Before the fix this returned 3
5451- // (the conjunct was silently dropped).
5452- assert_eq ! (
5453- rows, 2 ,
5454- "expected 2 rows with non-null struct; the rejected conjunct must \
5455- be applied post-scan, not silently dropped"
5456- ) ;
5520+ ( schema, file)
54575521 }
54585522
54595523 /// Helpers for tests that exercise parquet virtual columns
0 commit comments