From db584a8d16f2fa79f4d9decb58e72388f2db042f Mon Sep 17 00:00:00 2001 From: siddu Date: Thu, 8 Oct 2026 00:51:13 +0530 Subject: [PATCH 1/4] fix: skip HashJoinExec nodes carrying dynamic filters in JoinSelection (#26106) ## Which issue does this PR close? - Closes #26106. ## Rationale for this change `JoinSelection` is not safe to run on plans that already carry dynamic filters, which occurs during re-optimization passes (e.g. downstream pipelines that wrap an already-optimized plan into a writer sink and re-run physical optimization). Once `FilterPushdown` has wired up a dynamic filter between the build side and probe side of a `HashJoinExec`, swapping the inputs via `swap_inputs` panics with: `Internal error: Cannot swap HashJoinExec inputs after dynamic filter has been constructed` The join's build side is already committed at that stage, and swapping inputs would invalidate the dynamic filter expressions that reference probe-side columns. `JoinSelection` should detect this and leave the `HashJoinExec` unchanged instead of failing the query. ## What changes are included in this PR? - In `JoinSelection::statistical_join_selection_subrule`, check if `!hash_join.dynamic_expressions_produced().is_empty()` and return `None` (leaving the plan unchanged). - In `can_swap_hash_join`, guard against swapping when dynamic expressions are produced. - In `hash_join_swap_subrule`, guard against swapping unbounded left inputs when dynamic expressions are produced. - Added tests in `join_selection.rs` verifying that `JoinSelection` skips `HashJoinExec` carrying dynamic filters in both `CollectLeft` and `Partitioned` modes. ## What is the testing strategy for this PR? Added `test_join_selection_skips_hash_join_with_dynamic_filter` in `datafusion/core/tests/physical_optimizer/join_selection.rs` verifying that `JoinSelection` leaves the plan unchanged for both `CollectLeft` and `Partitioned` modes without error. ## Are there any user-facing changes? No API changes. Fixes an internal error when re-optimizing plans that contain dynamic filters. --- .../physical_optimizer/join_selection.rs | 57 ++++++++++++++++- .../physical-optimizer/src/join_selection.rs | 63 +++++++++++-------- 2 files changed, 93 insertions(+), 27 deletions(-) diff --git a/datafusion/core/tests/physical_optimizer/join_selection.rs b/datafusion/core/tests/physical_optimizer/join_selection.rs index d4561decc6331..d82f290051b1e 100644 --- a/datafusion/core/tests/physical_optimizer/join_selection.rs +++ b/datafusion/core/tests/physical_optimizer/join_selection.rs @@ -33,7 +33,9 @@ use datafusion_execution::{RecordBatchStream, SendableRecordBatchStream, TaskCon use datafusion_expr::Operator; use datafusion_physical_expr::PhysicalExprRef; use datafusion_physical_expr::expressions::col; -use datafusion_physical_expr::expressions::{BinaryExpr, Column, NegativeExpr}; +use datafusion_physical_expr::expressions::{ + BinaryExpr, Column, DynamicFilterPhysicalExpr, NegativeExpr, lit, +}; use datafusion_physical_expr::intervals::utils::check_support; use datafusion_physical_expr::{EquivalenceProperties, Partitioning, PhysicalExpr}; use datafusion_physical_expr_common::sort_expr::PhysicalSortExpr; @@ -1969,3 +1971,56 @@ fn test_join_with_maybe_swap_unbounded_case(t: TestCase) -> Result<()> { } Ok(()) } + +#[rstest] +#[case(PartitionMode::CollectLeft)] +#[case(PartitionMode::Partitioned)] +#[tokio::test] +async fn test_join_selection_skips_hash_join_with_dynamic_filter( + #[case] partition_mode: PartitionMode, +) -> Result<()> { + // Left has larger statistics than right, which would normally trigger swap + let (big, small) = create_big_and_small(); + let on = vec![( + Arc::new(Column::new_with_schema("big_col", &big.schema())?) as PhysicalExprRef, + Arc::new(Column::new_with_schema("small_col", &small.schema())?) as PhysicalExprRef, + )]; + + let dynamic_filter = Arc::new(DynamicFilterPhysicalExpr::new( + vec![Arc::clone(&on[0].1)], + lit(true), + )); + + #[allow(deprecated)] + let join = Arc::new( + HashJoinExec::try_new( + Arc::clone(&big), + Arc::clone(&small), + on, + None, + &JoinType::Inner, + None, + partition_mode, + NullEquality::NullEqualsNothing, + false, + )? + .with_dynamic_filter_expr(dynamic_filter)?, + ); + + let original_schema = join.schema(); + + // JoinSelection must not fail and must leave the join unchanged + let optimized = JoinSelection::new().optimize(join, &ConfigOptions::new())?; + let optimized_join = optimized + .downcast_ref::() + .expect("join should remain HashJoinExec without wrapping projection"); + + assert_eq!(optimized_join.partition_mode(), partition_mode); + assert_eq!(*optimized_join.join_type(), JoinType::Inner); + assert!(Arc::ptr_eq(optimized_join.left(), &big)); + assert!(Arc::ptr_eq(optimized_join.right(), &small)); + assert_eq!(optimized_join.schema(), original_schema); + assert_eq!(optimized_join.dynamic_expressions_produced().len(), 1); + + Ok(()) +} diff --git a/datafusion/physical-optimizer/src/join_selection.rs b/datafusion/physical-optimizer/src/join_selection.rs index 2b342bece7040..68cfe8abef62a 100644 --- a/datafusion/physical-optimizer/src/join_selection.rs +++ b/datafusion/physical-optimizer/src/join_selection.rs @@ -173,9 +173,11 @@ impl PhysicalOptimizerRule for JoinSelection { } /// Determines whether it is possible to swap inputs of a hash join - for null-aware joins, we can only swap an uncorrelated `LeftAnti` -/// (a single join key and no filter), because the swapped `RightAnti` has no per-row NULL handling +/// (a single join key and no filter), because the swapped `RightAnti` has no per-row NULL handling. +/// Joins carrying a dynamic filter cannot be swapped. fn can_swap_hash_join(hash_join: &HashJoinExec) -> bool { hash_join.join_type().supports_swap() + && hash_join.dynamic_expressions_produced().is_empty() && (!hash_join.null_aware || (*hash_join.join_type() == JoinType::LeftAnti && hash_join.on().len() == 1 @@ -301,32 +303,40 @@ fn statistical_join_selection_subrule( context: &dyn PhysicalOptimizerContext, ) -> Result>> { let transformed = if let Some(hash_join) = plan.downcast_ref::() { - match hash_join.partition_mode() { - PartitionMode::Auto => try_collect_left(hash_join, false, context)? - .map_or_else( - || partitioned_hash_join(hash_join, context).map(Some), - |v| Ok(Some(v)), - )?, - PartitionMode::CollectLeft => try_collect_left(hash_join, true, context)? - .map_or_else( - || partitioned_hash_join(hash_join, context).map(Some), - |v| Ok(Some(v)), - )?, - PartitionMode::Partitioned => { - let left = hash_join.left(); - let right = hash_join.right(); - if can_swap_hash_join(hash_join) - && should_swap_join_order(&**left, &**right, context)? - { - // Null-aware RightAnti only supports CollectLeft - let partition_mode = if hash_join.null_aware { - PartitionMode::CollectLeft + if !hash_join.dynamic_expressions_produced().is_empty() { + // Once a HashJoinExec carries a dynamic filter, its build side + // has already been determined and the dynamic filter has been wired + // to the probe side. Reordering inputs would invalidate the dynamic + // filter, so skip this join and leave it unchanged. + None + } else { + match hash_join.partition_mode() { + PartitionMode::Auto => try_collect_left(hash_join, false, context)? + .map_or_else( + || partitioned_hash_join(hash_join, context).map(Some), + |v| Ok(Some(v)), + )?, + PartitionMode::CollectLeft => try_collect_left(hash_join, true, context)? + .map_or_else( + || partitioned_hash_join(hash_join, context).map(Some), + |v| Ok(Some(v)), + )?, + PartitionMode::Partitioned => { + let left = hash_join.left(); + let right = hash_join.right(); + if can_swap_hash_join(hash_join) + && should_swap_join_order(&**left, &**right, context)? + { + // Null-aware RightAnti only supports CollectLeft + let partition_mode = if hash_join.null_aware { + PartitionMode::CollectLeft + } else { + PartitionMode::Partitioned + }; + hash_join.swap_inputs(partition_mode).map(Some)? } else { - PartitionMode::Partitioned - }; - hash_join.swap_inputs(partition_mode).map(Some)? - } else { - None + None + } } } } @@ -524,6 +534,7 @@ pub fn hash_join_swap_subrule( _config_options: &ConfigOptions, ) -> Result> { if let Some(hash_join) = input.downcast_ref::() + && hash_join.dynamic_expressions_produced().is_empty() && hash_join.left.boundedness().is_unbounded() && !hash_join.right.boundedness().is_unbounded() && !hash_join.null_aware // Don't swap null-aware anti joins From dd8e781db7b5f389630400ac931e5610eace06b9 Mon Sep 17 00:00:00 2001 From: siddu Date: Thu, 8 Oct 2026 16:15:32 +0530 Subject: [PATCH 2/4] fix: address cargo fmt and PartitionMode comparison in tests --- datafusion/core/tests/physical_optimizer/join_selection.rs | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/datafusion/core/tests/physical_optimizer/join_selection.rs b/datafusion/core/tests/physical_optimizer/join_selection.rs index d82f290051b1e..26513d20552a2 100644 --- a/datafusion/core/tests/physical_optimizer/join_selection.rs +++ b/datafusion/core/tests/physical_optimizer/join_selection.rs @@ -1983,7 +1983,8 @@ async fn test_join_selection_skips_hash_join_with_dynamic_filter( let (big, small) = create_big_and_small(); let on = vec![( Arc::new(Column::new_with_schema("big_col", &big.schema())?) as PhysicalExprRef, - Arc::new(Column::new_with_schema("small_col", &small.schema())?) as PhysicalExprRef, + Arc::new(Column::new_with_schema("small_col", &small.schema())?) + as PhysicalExprRef, )]; let dynamic_filter = Arc::new(DynamicFilterPhysicalExpr::new( @@ -2015,7 +2016,7 @@ async fn test_join_selection_skips_hash_join_with_dynamic_filter( .downcast_ref::() .expect("join should remain HashJoinExec without wrapping projection"); - assert_eq!(optimized_join.partition_mode(), partition_mode); + assert_eq!(*optimized_join.partition_mode(), partition_mode); assert_eq!(*optimized_join.join_type(), JoinType::Inner); assert!(Arc::ptr_eq(optimized_join.left(), &big)); assert!(Arc::ptr_eq(optimized_join.right(), &small)); From 5fb39d7ce7f47d137401f9cd6ea954e7d11a8527 Mon Sep 17 00:00:00 2001 From: siddu Date: Thu, 8 Oct 2026 23:35:43 +0530 Subject: [PATCH 3/4] Fix clippy and codecov issues in join selection tests --- .../physical_optimizer/join_selection.rs | 68 ++++++++++++++++++- .../physical-optimizer/src/join_selection.rs | 4 +- 2 files changed, 68 insertions(+), 4 deletions(-) diff --git a/datafusion/core/tests/physical_optimizer/join_selection.rs b/datafusion/core/tests/physical_optimizer/join_selection.rs index 26513d20552a2..b259bd866ae5f 100644 --- a/datafusion/core/tests/physical_optimizer/join_selection.rs +++ b/datafusion/core/tests/physical_optimizer/join_selection.rs @@ -1992,7 +1992,7 @@ async fn test_join_selection_skips_hash_join_with_dynamic_filter( lit(true), )); - #[allow(deprecated)] + #[expect(deprecated)] let join = Arc::new( HashJoinExec::try_new( Arc::clone(&big), @@ -2025,3 +2025,69 @@ async fn test_join_selection_skips_hash_join_with_dynamic_filter( Ok(()) } + +#[tokio::test] +async fn test_join_selection_skips_unbounded_hash_join_with_dynamic_filter() -> Result<()> { + let left_exec: Arc = Arc::new(UnboundedExec::new( + None, + RecordBatch::new_empty(Arc::new(Schema::new(vec![Field::new( + "a", + DataType::Int32, + false, + )]))), + 2, + )); + let right_exec: Arc = Arc::new(UnboundedExec::new( + Some(1), + RecordBatch::new_empty(Arc::new(Schema::new(vec![Field::new( + "b", + DataType::Int32, + false, + )]))), + 2, + )); + + let on = vec![( + col("a", &left_exec.schema())?, + col("b", &right_exec.schema())?, + )]; + + let dynamic_filter = Arc::new(DynamicFilterPhysicalExpr::new( + vec![Arc::clone(&on[0].1)], + lit(true), + )); + + #[expect(deprecated)] + let join = Arc::new( + HashJoinExec::try_new( + Arc::clone(&left_exec), + Arc::clone(&right_exec), + on, + None, + &JoinType::Inner, + None, + PartitionMode::Partitioned, + NullEquality::NullEqualsNothing, + false, + )? + .with_dynamic_filter_expr(dynamic_filter)?, + ); + + let original_schema = join.schema(); + + // hash_join_swap_subrule would normally swap unbounded left with bounded right, + // but must skip this join because it has a dynamic filter. + let optimized = JoinSelection::new().optimize(join, &ConfigOptions::new())?; + let optimized_join = optimized + .downcast_ref::() + .expect("join should remain HashJoinExec without wrapping projection"); + + assert_eq!(*optimized_join.partition_mode(), PartitionMode::Partitioned); + assert_eq!(*optimized_join.join_type(), JoinType::Inner); + assert!(Arc::ptr_eq(optimized_join.left(), &left_exec)); + assert!(Arc::ptr_eq(optimized_join.right(), &right_exec)); + assert_eq!(optimized_join.schema(), original_schema); + assert_eq!(optimized_join.dynamic_expressions_produced().len(), 1); + + Ok(()) +} diff --git a/datafusion/physical-optimizer/src/join_selection.rs b/datafusion/physical-optimizer/src/join_selection.rs index 68cfe8abef62a..a0eb6997a4c6b 100644 --- a/datafusion/physical-optimizer/src/join_selection.rs +++ b/datafusion/physical-optimizer/src/join_selection.rs @@ -173,11 +173,9 @@ impl PhysicalOptimizerRule for JoinSelection { } /// Determines whether it is possible to swap inputs of a hash join - for null-aware joins, we can only swap an uncorrelated `LeftAnti` -/// (a single join key and no filter), because the swapped `RightAnti` has no per-row NULL handling. -/// Joins carrying a dynamic filter cannot be swapped. +/// (a single join key and no filter), because the swapped `RightAnti` has no per-row NULL handling fn can_swap_hash_join(hash_join: &HashJoinExec) -> bool { hash_join.join_type().supports_swap() - && hash_join.dynamic_expressions_produced().is_empty() && (!hash_join.null_aware || (*hash_join.join_type() == JoinType::LeftAnti && hash_join.on().len() == 1 From 16a613590498e367de66cf5729c3882262aa6b79 Mon Sep 17 00:00:00 2001 From: siddu Date: Fri, 9 Oct 2026 23:19:48 +0530 Subject: [PATCH 4/4] fix: correct nullability of AND/OR expressions to resolve plan mismatch in CSE --- datafusion/core/tests/issue_25978.rs | 27 ++++++++++++++++++++ datafusion/expr/src/expr_schema.rs | 38 ++++++++++++++++++++++++++++ 2 files changed, 65 insertions(+) create mode 100644 datafusion/core/tests/issue_25978.rs diff --git a/datafusion/core/tests/issue_25978.rs b/datafusion/core/tests/issue_25978.rs new file mode 100644 index 0000000000000..2e5005e37dcac --- /dev/null +++ b/datafusion/core/tests/issue_25978.rs @@ -0,0 +1,27 @@ +use datafusion::prelude::*; +use datafusion::error::Result; + +#[tokio::test] +async fn test_issue_25978() -> Result<()> { + let ctx = SessionContext::new(); + + // CREATE TABLE a (k VARCHAR NOT NULL, x VARCHAR) AS VALUES ('k1', 'y'), ('k2', NULL); + // CREATE TABLE b (k VARCHAR NOT NULL) AS VALUES ('k3'); + ctx.sql("CREATE TABLE a (k VARCHAR NOT NULL, x VARCHAR) AS VALUES ('k1', 'y'), ('k2', NULL)").await?.collect().await?; + ctx.sql("CREATE TABLE b (k VARCHAR NOT NULL) AS VALUES ('k3')").await?.collect().await?; + + // SELECT k, bool_or(f) AS f FROM ( + // SELECT k, coalesce(x = 'y', FALSE) AS f FROM a + // UNION ALL + // SELECT k, TRUE AS f FROM b + // ) u GROUP BY k; + let df = ctx.sql("SELECT k, bool_or(f) AS f FROM ( + SELECT k, coalesce(x = 'y', FALSE) AS f FROM a + UNION ALL + SELECT k, TRUE AS f FROM b + ) u GROUP BY k").await?; + + df.collect().await?; + + Ok(()) +} diff --git a/datafusion/expr/src/expr_schema.rs b/datafusion/expr/src/expr_schema.rs index 8aaca2b72b2c6..31031e6ff0647 100644 --- a/datafusion/expr/src/expr_schema.rs +++ b/datafusion/expr/src/expr_schema.rs @@ -435,6 +435,24 @@ impl ExprSchemable for Expr { Expr::ScalarSubquery(subquery) => Ok(scalar_subquery_nullable(subquery)), Expr::BinaryExpr(BinaryExpr { left, right, op }) => match op { Operator::IsDistinctFrom | Operator::IsNotDistinctFrom => Ok(false), + Operator::And => { + if contains_is_not_null(left.as_ref(), right.as_ref()) { + return Ok(false); + } + if contains_is_not_null(right.as_ref(), left.as_ref()) { + return Ok(false); + } + Ok(left.nullable(input_schema)? || right.nullable(input_schema)?) + } + Operator::Or => { + if contains_is_null(left.as_ref(), right.as_ref()) { + return Ok(false); + } + if contains_is_null(right.as_ref(), left.as_ref()) { + return Ok(false); + } + Ok(left.nullable(input_schema)? || right.nullable(input_schema)?) + } _ => Ok(left.nullable(input_schema)? || right.nullable(input_schema)?), }, Expr::Like(Like { expr, pattern, .. }) @@ -1623,3 +1641,23 @@ mod tests { } } } + +fn contains_is_not_null(expr: &Expr, target: &Expr) -> bool { + if let Expr::IsNotNull(inner) = expr { + return inner.as_ref() == target; + } + if let Expr::BinaryExpr(BinaryExpr { left, right, op: Operator::And }) = expr { + return contains_is_not_null(left, target) || contains_is_not_null(right, target); + } + false +} + +fn contains_is_null(expr: &Expr, target: &Expr) -> bool { + if let Expr::IsNull(inner) = expr { + return inner.as_ref() == target; + } + if let Expr::BinaryExpr(BinaryExpr { left, right, op: Operator::Or }) = expr { + return contains_is_null(left, target) || contains_is_null(right, target); + } + false +}