neilconway opened a new issue, #24002:
URL: https://github.com/apache/datafusion/issues/24002

   ### Describe the bug
   
   `HashJoinExec::gather_filters_for_pushdown` lets parent filters push to both 
sides of LeftAnti/RightAnti joins: `lr_is_preserved` returns `(true, true)`, 
with the non-preserved side's allowed column set restricted to join keys. Two 
behaviors combine into a wrong-results bug:
   
   1. A column-free filter (e.g. a constant false) passes the allowed-indices 
check vacuously, so it is marked supported for the non-preserved (probe) side.
   2. `HashJoinExec::handle_child_pushdown_result` combines child results with 
`FilterPushdownPropagation::if_any`, so when the probe side absorbs a filter 
but the build side does not (e.g. the build side is a source without filter 
pushdown support), the join reports the filter as pushed down and the 
`FilterExec` above deletes it.
   
   The filter then applies only to the probe side. For an anti join that is 
unsound: filtering the probe side creates output rows. With `FilterExec(false)` 
above a `LeftAnti` join, the pushed false empties the probe side and the join 
emits every build-side row instead of zero rows.
   
   Key-referencing filters routed to the non-preserved side are unsound under 
the same asymmetry: for a semi join a match implies key equality, so dropping 
the filter above is safe; for an anti join a non-match implies nothing, so it 
is not.
   
   ### To Reproduce
   
   On current `main`, with the patch below saved as `/tmp/antijoin_repro.patch`:
   
   ```shell
   git apply /tmp/antijoin_repro.patch
   cargo test -p datafusion --test core_integration -- 
repro_constant_filter_lost_above_anti_join --nocapture
   ```
   
   <details>
   <summary><code>antijoin_repro.patch</code></summary>
   
   ```diff
   diff --git a/datafusion/core/tests/physical_optimizer/filter_pushdown.rs 
b/datafusion/core/tests/physical_optimizer/filter_pushdown.rs
   index 80b3f37b7..bbc0c45ab 100644
   --- a/datafusion/core/tests/physical_optimizer/filter_pushdown.rs
   +++ b/datafusion/core/tests/physical_optimizer/filter_pushdown.rs
   @@ -3451,3 +3451,87 @@ fn post_phase_is_idempotent_on_hash_join() {
            "second invocation of FilterPushdown::new_post_optimization mutated 
the plan",
        );
    }
   +
   +/// Repro: a constant FALSE filter above a LeftAnti hash join is deleted by 
the
   +/// filter pushdown optimizer when only the probe (right) side absorbs it,
   +/// changing query results. The join reports the filter as pushed down via
   +/// `if_any`, so `FilterExec` removes it, but only the right side applies 
it --
   +/// emptying the probe side makes the anti join emit every build-side row.
   +#[tokio::test]
   +async fn repro_constant_filter_lost_above_anti_join() {
   +    use datafusion_common::JoinType;
   +    use datafusion_physical_plan::joins::{HashJoinExec, PartitionMode};
   +
   +    // The test scans are single-use, so build an identical plan for each 
run.
   +    let make_plan = || {
   +        let left_schema =
   +            Arc::new(Schema::new(vec![Field::new("a", DataType::Utf8, 
false)]));
   +        // Build (left, preserved) side: scan that does NOT absorb pushed 
filters.
   +        let left_scan = TestScanBuilder::new(Arc::clone(&left_schema))
   +            .with_support(false)
   +            .with_batches(vec![record_batch!(("a", Utf8, ["x", "y", 
"z"])).unwrap()])
   +            .build();
   +        let right_schema =
   +            Arc::new(Schema::new(vec![Field::new("a", DataType::Utf8, 
false)]));
   +        // Probe (right, non-preserved) side: scan that DOES absorb pushed 
filters.
   +        let right_scan = TestScanBuilder::new(Arc::clone(&right_schema))
   +            .with_support(true)
   +            .with_batches(vec![record_batch!(("a", Utf8, ["x"])).unwrap()])
   +            .build();
   +        let on = vec![(
   +            col("a", &left_schema).unwrap(),
   +            col("a", &right_schema).unwrap(),
   +        )];
   +        let join = Arc::new(
   +            HashJoinExec::try_new(
   +                left_scan,
   +                right_scan,
   +                on,
   +                None,
   +                &JoinType::LeftAnti,
   +                None,
   +                PartitionMode::CollectLeft,
   +                datafusion_common::NullEquality::NullEqualsNothing,
   +                false,
   +            )
   +            .unwrap(),
   +        );
   +        // Without the filter the anti join returns rows y and z; the 
constant
   +        // FALSE filter above it means the query must return 0 rows.
   +        let predicate = 
Arc::new(Literal::new(ScalarValue::Boolean(Some(false))));
   +        Arc::new(FilterExec::try_new(predicate, join).unwrap())
   +            as Arc<dyn ExecutionPlan>
   +    };
   +
   +    let session_ctx = SessionContext::new();
   +    session_ctx.register_object_store(
   +        ObjectStoreUrl::parse("test://").unwrap().as_ref(),
   +        Arc::new(InMemory::new()),
   +    );
   +    let count = |plan: Arc<dyn ExecutionPlan>| {
   +        let ctx = session_ctx.state().task_ctx();
   +        async move {
   +            collect(plan, ctx)
   +                .await
   +                .unwrap()
   +                .iter()
   +                .map(|b| b.num_rows())
   +                .sum::<usize>()
   +        }
   +    };
   +
   +    let unoptimized_rows = count(make_plan()).await;
   +
   +    let mut config = ConfigOptions::default();
   +    config.execution.parquet.pushdown_filters = true;
   +    let optimized = FilterPushdown::new().optimize(make_plan(), 
&config).unwrap();
   +    println!("optimized plan:\n{}", format_plan_for_test(&optimized));
   +
   +    let optimized_rows = count(optimized).await;
   +
   +    println!("unoptimized_rows={unoptimized_rows} 
optimized_rows={optimized_rows}");
   +    assert_eq!(
   +        unoptimized_rows, optimized_rows,
   +        "filter pushdown changed query results"
   +    );
   +}
   ```
   
   </details>
   
   ### Expected behavior
   
   _No response_
   
   ### Additional context
   
   _No response_


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to