From b381100da1ac9d37831aa8147b33df31cb3f6576 Mon Sep 17 00:00:00 2001 From: Wenchen Fan Date: Wed, 16 Sep 2026 18:38:23 +0000 Subject: [PATCH 1/9] [SPARK-36082][SQL][FOLLOWUP] Floor the NAAJ broadcast threshold --- .../spark/sql/catalyst/optimizer/joins.scala | 17 ++-- .../apache/spark/sql/internal/SQLConf.scala | 7 +- .../optimizer/JoinSelectionHelperSuite.scala | 84 +++++++++++++------ 3 files changed, 74 insertions(+), 34 deletions(-) diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/joins.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/joins.scala index 4628fc32ea344..ecb8ed4f7eccb 100644 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/joins.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/joins.scala @@ -438,12 +438,19 @@ trait JoinSelectionHelper extends Logging { getBroadcastBuildSide(join, hintOnly = true, conf).orElse { if (noShufflePlannedBefore) getBroadcastBuildSide(join, hintOnly = false, conf) else None } - // `JoinSelection` always builds from the right for this shape. A negative threshold preserves - // the original unbounded NAAJ behavior, while zero disables the broadcast hash optimization. + // `JoinSelection` always builds from the right for this shape. Do not reject the hash + // optimization when regular join planning would broadcast the right side, as the fallback + // would still broadcast it with a slower nested-loop join. case j @ ExtractSingleColumnNullAwareAntiJoin(_, _) => - val threshold = conf.nullAwareAntiJoinBroadcastThreshold - val rightSize = j.right.stats.sizeInBytes - if (threshold < 0 || (threshold > 0 && rightSize >= 0 && rightSize <= threshold)) { + val dedicatedThreshold = conf.nullAwareAntiJoinBroadcastThreshold + val canBroadcast = dedicatedThreshold < 0 || { + val effectiveThreshold = math.max(dedicatedThreshold, conf.autoBroadcastJoinThreshold) + effectiveThreshold > 0 && { + val rightSize = j.right.stats.sizeInBytes + rightSize >= 0 && rightSize <= effectiveThreshold + } + } + if (canBroadcast) { Some(BuildRight) } else { None diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala index 2ec67312f0218..8e61a41f4d796 100644 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala @@ -7452,8 +7452,11 @@ object SQLConf { "single-column null-aware anti join for which Spark uses the broadcast hash join " + "optimization. This configuration takes effect only when " + "spark.sql.optimizeNullAwareAntiJoin is enabled. A negative value allows the " + - "optimization regardless of the estimated size, while zero disables it. If the " + - "estimated size exceeds a positive value, Spark falls back to regular join planning. " + + "optimization regardless of the estimated size. For a nonnegative value, the effective " + + "threshold is the larger of this value and spark.sql.autoBroadcastJoinThreshold. Thus, " + + "zero disables the optimization only when automatic broadcasting is also disabled. If " + + "the estimated size exceeds the effective threshold, Spark falls back to regular join " + + "planning. " + "The fallback may still broadcast the right side with a nested-loop representation " + "that uses more memory and runs in O(M * N) time. Join hints do not override this " + "configuration when the broadcast hash optimization is selected. This configuration " + diff --git a/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/optimizer/JoinSelectionHelperSuite.scala b/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/optimizer/JoinSelectionHelperSuite.scala index 71fa91f5c37e0..d056cc972b117 100644 --- a/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/optimizer/JoinSelectionHelperSuite.scala +++ b/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/optimizer/JoinSelectionHelperSuite.scala @@ -18,9 +18,9 @@ package org.apache.spark.sql.catalyst.optimizer import org.apache.spark.sql.catalyst.dsl.expressions._ -import org.apache.spark.sql.catalyst.expressions.{AttributeMap, EqualTo, IsNull, Or} +import org.apache.spark.sql.catalyst.expressions.{Attribute, AttributeMap, EqualTo, IsNull, Or} import org.apache.spark.sql.catalyst.plans.{Inner, LeftAnti, PlanTest} -import org.apache.spark.sql.catalyst.plans.logical.{BROADCAST, HintInfo, Join, JoinHint, NO_BROADCAST_HASH, SHUFFLE_HASH} +import org.apache.spark.sql.catalyst.plans.logical.{BROADCAST, HintInfo, Join, JoinHint, LeafNode, LogicalPlan, NO_BROADCAST_HASH, SHUFFLE_HASH} import org.apache.spark.sql.catalyst.statsEstimation.StatsTestPlan import org.apache.spark.sql.internal.SQLConf @@ -40,6 +40,11 @@ class JoinSelectionHelperSuite extends PlanTest with JoinSelectionHelper { private val join = Join(left, right, Inner, None, JoinHint(None, None)) + private def nullAwareAntiJoin(rightPlan: LogicalPlan = right): Join = { + val equality = EqualTo(left.output.head, rightPlan.output.head) + Join(left, rightPlan, LeftAnti, Some(Or(equality, IsNull(equality))), JoinHint.NONE) + } + private val hintBroadcast = Some(HintInfo(Some(BROADCAST))) private val hintNotToBroadcast = Some(HintInfo(Some(NO_BROADCAST_HASH))) private val hintShuffleHash = Some(HintInfo(Some(SHUFFLE_HASH))) @@ -195,48 +200,73 @@ class JoinSelectionHelperSuite extends PlanTest with JoinSelectionHelper { } } - test("getBroadcastHashJoinBuildSide uses the null-aware anti join broadcast threshold") { - val leftKey = left.output.head - val rightKey = right.output.head - val condition = Or(EqualTo(leftKey, rightKey), IsNull(EqualTo(leftKey, rightKey))) - val nullAwareAntiJoin = Join(left, right, LeftAnti, Some(condition), JoinHint.NONE) + test("NAAJ broadcast threshold is floored by the automatic broadcast threshold") { + val autoThresholdRight = right.copy( + rowCount = 10 * 1024 * 1024, + size = Some(10 * 1024 * 1024)) val largeRight = right.copy(rowCount = 20000000, size = Some(20000000)) - val negativeSizeRight = right.copy(size = Some(-1)) - val overLongMaxRight = right.copy( - rowCount = BigInt(Long.MaxValue) + 1, - size = Some(BigInt(Long.MaxValue) + 1)) withSQLConf( SQLConf.OPTIMIZE_NULL_AWARE_ANTI_JOIN.key -> "true", - SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> "10MB") { - assert(getBroadcastHashJoinBuildSide(nullAwareAntiJoin, SQLConf.get) === Some(BuildRight)) + SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> "10MB", + SQLConf.NULL_AWARE_ANTI_JOIN_BROADCAST_THRESHOLD.key -> "0") { + assert(getBroadcastHashJoinBuildSide(nullAwareAntiJoin(), SQLConf.get) === Some(BuildRight)) assert(getBroadcastHashJoinBuildSide( - nullAwareAntiJoin.copy(right = largeRight), SQLConf.get) === Some(BuildRight)) + nullAwareAntiJoin(autoThresholdRight), SQLConf.get) === Some(BuildRight)) assert(getBroadcastHashJoinBuildSide( - nullAwareAntiJoin.copy(right = overLongMaxRight), SQLConf.get) === Some(BuildRight)) + nullAwareAntiJoin(largeRight), SQLConf.get).isEmpty) } - withSQLConf(SQLConf.NULL_AWARE_ANTI_JOIN_BROADCAST_THRESHOLD.key -> "-2") { + withSQLConf( + SQLConf.OPTIMIZE_NULL_AWARE_ANTI_JOIN.key -> "true", + SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> "10MB", + SQLConf.NULL_AWARE_ANTI_JOIN_BROADCAST_THRESHOLD.key -> "20MB") { assert(getBroadcastHashJoinBuildSide( - nullAwareAntiJoin.copy(right = overLongMaxRight), SQLConf.get) === Some(BuildRight)) + nullAwareAntiJoin(largeRight), SQLConf.get) === Some(BuildRight)) } - withSQLConf(SQLConf.NULL_AWARE_ANTI_JOIN_BROADCAST_THRESHOLD.key -> "0") { - assert(getBroadcastHashJoinBuildSide(nullAwareAntiJoin, SQLConf.get).isEmpty) + withSQLConf( + SQLConf.OPTIMIZE_NULL_AWARE_ANTI_JOIN.key -> "true", + SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> "-1", + SQLConf.NULL_AWARE_ANTI_JOIN_BROADCAST_THRESHOLD.key -> "0") { + assert(getBroadcastHashJoinBuildSide(nullAwareAntiJoin(), SQLConf.get).isEmpty) } + } - withSQLConf( - SQLConf.OPTIMIZE_NULL_AWARE_ANTI_JOIN.key -> "false", - SQLConf.NULL_AWARE_ANTI_JOIN_BROADCAST_THRESHOLD.key -> "-1") { - assert(getBroadcastHashJoinBuildSide(nullAwareAntiJoin, SQLConf.get).isEmpty) + test("NAAJ broadcast threshold short-circuits config-only decisions") { + case class ThrowingStatsPlan() extends LeafNode { + override def output: Seq[Attribute] = right.output } + val nullAwareAntiJoinWithoutStats = nullAwareAntiJoin(ThrowingStatsPlan()) - withSQLConf(SQLConf.NULL_AWARE_ANTI_JOIN_BROADCAST_THRESHOLD.key -> "10MB") { - assert(getBroadcastHashJoinBuildSide(nullAwareAntiJoin, SQLConf.get) === Some(BuildRight)) + withSQLConf( + SQLConf.OPTIMIZE_NULL_AWARE_ANTI_JOIN.key -> "true", + SQLConf.NULL_AWARE_ANTI_JOIN_BROADCAST_THRESHOLD.key -> "-2") { assert(getBroadcastHashJoinBuildSide( - nullAwareAntiJoin.copy(right = largeRight), SQLConf.get).isEmpty) + nullAwareAntiJoinWithoutStats, SQLConf.get) === Some(BuildRight)) + } + + withSQLConf( + SQLConf.OPTIMIZE_NULL_AWARE_ANTI_JOIN.key -> "true", + SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> "-1", + SQLConf.NULL_AWARE_ANTI_JOIN_BROADCAST_THRESHOLD.key -> "0") { + assert(getBroadcastHashJoinBuildSide(nullAwareAntiJoinWithoutStats, SQLConf.get).isEmpty) + } + } + + test("NAAJ broadcast threshold rejects unknown sizes and respects the optimization flag") { + withSQLConf( + SQLConf.OPTIMIZE_NULL_AWARE_ANTI_JOIN.key -> "true", + SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> "-1", + SQLConf.NULL_AWARE_ANTI_JOIN_BROADCAST_THRESHOLD.key -> "10MB") { assert(getBroadcastHashJoinBuildSide( - nullAwareAntiJoin.copy(right = negativeSizeRight), SQLConf.get).isEmpty) + nullAwareAntiJoin(right.copy(size = Some(-1))), SQLConf.get).isEmpty) + } + + withSQLConf( + SQLConf.OPTIMIZE_NULL_AWARE_ANTI_JOIN.key -> "false", + SQLConf.NULL_AWARE_ANTI_JOIN_BROADCAST_THRESHOLD.key -> "-1") { + assert(getBroadcastHashJoinBuildSide(nullAwareAntiJoin(), SQLConf.get).isEmpty) } } } From 0f5f063d366b176cef9cbf1b4aebfe103663c099 Mon Sep 17 00:00:00 2001 From: Wenchen Fan Date: Thu, 17 Sep 2026 11:56:30 +0000 Subject: [PATCH 2/9] [SPARK-36082][SQL] Address review comments --- .../spark/sql/catalyst/optimizer/joins.scala | 17 ++++-- .../apache/spark/sql/internal/SQLConf.scala | 17 +++--- .../optimizer/JoinSelectionHelperSuite.scala | 54 ++++++++++++++++++- .../LeftSemiAntiJoinPushDownSuite.scala | 27 ++++++++++ .../org/apache/spark/sql/JoinSuite.scala | 8 +-- 5 files changed, 106 insertions(+), 17 deletions(-) diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/joins.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/joins.scala index ecb8ed4f7eccb..46f00e453374b 100644 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/joins.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/joins.scala @@ -443,11 +443,18 @@ trait JoinSelectionHelper extends Logging { // would still broadcast it with a slower nested-loop join. case j @ ExtractSingleColumnNullAwareAntiJoin(_, _) => val dedicatedThreshold = conf.nullAwareAntiJoinBroadcastThreshold - val canBroadcast = dedicatedThreshold < 0 || { - val effectiveThreshold = math.max(dedicatedThreshold, conf.autoBroadcastJoinThreshold) - effectiveThreshold > 0 && { - val rightSize = j.right.stats.sizeInBytes - rightSize >= 0 && rightSize <= effectiveThreshold + val canBroadcast = if (dedicatedThreshold < 0) { + true + } else { + val automaticBroadcastDisabled = conf.autoBroadcastJoinThreshold <= 0 && + conf.getConf(SQLConf.ADAPTIVE_AUTO_BROADCASTJOIN_THRESHOLD).forall(_ <= 0) + if (dedicatedThreshold == 0 && automaticBroadcastDisabled) { + false + } else { + canBroadcastBySize(j.right, conf) || (dedicatedThreshold > 0 && { + val rightSize = j.right.stats.sizeInBytes + rightSize >= 0 && rightSize <= dedicatedThreshold + }) } } if (canBroadcast) { diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala index 8e61a41f4d796..e73cea7611d3d 100644 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala @@ -7452,15 +7452,18 @@ object SQLConf { "single-column null-aware anti join for which Spark uses the broadcast hash join " + "optimization. This configuration takes effect only when " + "spark.sql.optimizeNullAwareAntiJoin is enabled. A negative value allows the " + - "optimization regardless of the estimated size. For a nonnegative value, the effective " + - "threshold is the larger of this value and spark.sql.autoBroadcastJoinThreshold. Thus, " + - "zero disables the optimization only when automatic broadcasting is also disabled. If " + - "the estimated size exceeds the effective threshold, Spark falls back to regular join " + - "planning. " + + "optimization regardless of the estimated size. For a nonnegative value, the " + + "optimization is also allowed when regular join planning considers the right side " + + "broadcastable. " + + "Regular planning uses spark.sql.adaptive.autoBroadcastJoinThreshold for runtime " + + "statistics when it is set, and spark.sql.autoBroadcastJoinThreshold otherwise. Thus, " + + "zero disables the optimization only when automatic broadcasting is also disabled. " + "The fallback may still broadcast the right side with a nested-loop representation " + "that uses more memory and runs in O(M * N) time. Join hints do not override this " + - "configuration when the broadcast hash optimization is selected. This configuration " + - "also controls whether a null-aware anti join can be pushed below an aggregate.") + "configuration when the broadcast hash optimization is selected. Set " + + "spark.sql.optimizeNullAwareAntiJoin to false to disable the optimization without " + + "changing automatic broadcast thresholds. The same eligibility decision controls " + + "whether a null-aware anti join can be pushed below an aggregate.") .version("4.2.1") .withBindingPolicy(ConfigBindingPolicy.NOT_APPLICABLE) .bytesConf(ByteUnit.BYTE) diff --git a/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/optimizer/JoinSelectionHelperSuite.scala b/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/optimizer/JoinSelectionHelperSuite.scala index d056cc972b117..25119de02fde0 100644 --- a/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/optimizer/JoinSelectionHelperSuite.scala +++ b/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/optimizer/JoinSelectionHelperSuite.scala @@ -20,7 +20,7 @@ package org.apache.spark.sql.catalyst.optimizer import org.apache.spark.sql.catalyst.dsl.expressions._ import org.apache.spark.sql.catalyst.expressions.{Attribute, AttributeMap, EqualTo, IsNull, Or} import org.apache.spark.sql.catalyst.plans.{Inner, LeftAnti, PlanTest} -import org.apache.spark.sql.catalyst.plans.logical.{BROADCAST, HintInfo, Join, JoinHint, LeafNode, LogicalPlan, NO_BROADCAST_HASH, SHUFFLE_HASH} +import org.apache.spark.sql.catalyst.plans.logical.{BROADCAST, HintInfo, Join, JoinHint, LeafNode, LogicalPlan, NO_BROADCAST_HASH, SHUFFLE_HASH, Statistics} import org.apache.spark.sql.catalyst.statsEstimation.StatsTestPlan import org.apache.spark.sql.internal.SQLConf @@ -204,6 +204,9 @@ class JoinSelectionHelperSuite extends PlanTest with JoinSelectionHelper { val autoThresholdRight = right.copy( rowCount = 10 * 1024 * 1024, size = Some(10 * 1024 * 1024)) + val betweenThresholdsRight = right.copy( + rowCount = 8 * 1024 * 1024, + size = Some(8 * 1024 * 1024)) val largeRight = right.copy(rowCount = 20000000, size = Some(20000000)) withSQLConf( @@ -225,6 +228,14 @@ class JoinSelectionHelperSuite extends PlanTest with JoinSelectionHelper { nullAwareAntiJoin(largeRight), SQLConf.get) === Some(BuildRight)) } + withSQLConf( + SQLConf.OPTIMIZE_NULL_AWARE_ANTI_JOIN.key -> "true", + SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> "10MB", + SQLConf.NULL_AWARE_ANTI_JOIN_BROADCAST_THRESHOLD.key -> "5MB") { + assert(getBroadcastHashJoinBuildSide( + nullAwareAntiJoin(betweenThresholdsRight), SQLConf.get) === Some(BuildRight)) + } + withSQLConf( SQLConf.OPTIMIZE_NULL_AWARE_ANTI_JOIN.key -> "true", SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> "-1", @@ -233,9 +244,49 @@ class JoinSelectionHelperSuite extends PlanTest with JoinSelectionHelper { } } + test("NAAJ broadcast threshold is unlimited by default") { + val overLongMaxRight = right.copy( + rowCount = BigInt(Long.MaxValue) + 1, + size = Some(BigInt(Long.MaxValue) + 1)) + + withSQLConf( + SQLConf.OPTIMIZE_NULL_AWARE_ANTI_JOIN.key -> "true", + SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> "-1") { + assert(getBroadcastHashJoinBuildSide( + nullAwareAntiJoin(overLongMaxRight), SQLConf.get) === Some(BuildRight)) + } + } + + test("NAAJ broadcast threshold uses the adaptive threshold for runtime statistics") { + case class RuntimeStatsPlan(size: BigInt) extends LeafNode { + override def output: Seq[Attribute] = right.output + override def computeStats(): Statistics = Statistics(sizeInBytes = size, isRuntime = true) + } + val runtimeRight = RuntimeStatsPlan(5 * 1024 * 1024) + + withSQLConf( + SQLConf.OPTIMIZE_NULL_AWARE_ANTI_JOIN.key -> "true", + SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> "10MB", + SQLConf.ADAPTIVE_AUTO_BROADCASTJOIN_THRESHOLD.key -> "1MB", + SQLConf.NULL_AWARE_ANTI_JOIN_BROADCAST_THRESHOLD.key -> "0") { + assert(getBroadcastHashJoinBuildSide(nullAwareAntiJoin(runtimeRight), SQLConf.get).isEmpty) + } + + withSQLConf( + SQLConf.OPTIMIZE_NULL_AWARE_ANTI_JOIN.key -> "true", + SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> "1MB", + SQLConf.ADAPTIVE_AUTO_BROADCASTJOIN_THRESHOLD.key -> "10MB", + SQLConf.NULL_AWARE_ANTI_JOIN_BROADCAST_THRESHOLD.key -> "0") { + assert(getBroadcastHashJoinBuildSide( + nullAwareAntiJoin(runtimeRight), SQLConf.get) === Some(BuildRight)) + } + } + test("NAAJ broadcast threshold short-circuits config-only decisions") { case class ThrowingStatsPlan() extends LeafNode { override def output: Seq[Attribute] = right.output + override def computeStats(): Statistics = + throw new IllegalStateException("statistics should not be read") } val nullAwareAntiJoinWithoutStats = nullAwareAntiJoin(ThrowingStatsPlan()) @@ -259,6 +310,7 @@ class JoinSelectionHelperSuite extends PlanTest with JoinSelectionHelper { SQLConf.OPTIMIZE_NULL_AWARE_ANTI_JOIN.key -> "true", SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> "-1", SQLConf.NULL_AWARE_ANTI_JOIN_BROADCAST_THRESHOLD.key -> "10MB") { + assert(getBroadcastHashJoinBuildSide(nullAwareAntiJoin(), SQLConf.get) === Some(BuildRight)) assert(getBroadcastHashJoinBuildSide( nullAwareAntiJoin(right.copy(size = Some(-1))), SQLConf.get).isEmpty) } diff --git a/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/optimizer/LeftSemiAntiJoinPushDownSuite.scala b/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/optimizer/LeftSemiAntiJoinPushDownSuite.scala index ba7386f8c9b5e..a5c0536fcf63d 100644 --- a/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/optimizer/LeftSemiAntiJoinPushDownSuite.scala +++ b/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/optimizer/LeftSemiAntiJoinPushDownSuite.scala @@ -142,6 +142,33 @@ class LeftSemiAntiJoinPushDownSuite extends PlanTest { comparePlans(optimized, originalQuery.analyze) } + test("Aggregate: NAAJ pushdown follows the effective broadcast threshold") { + val aggregate = testRelation.groupBy($"b")($"b") + val equality = $"b" === $"d" + val originalQuery = aggregate.join( + testRelation1, + joinType = LeftAnti, + condition = Some(equality || IsNull(equality))) + val pushedDownQuery = testRelation + .join( + testRelation1, + joinType = LeftAnti, + condition = Some(equality || IsNull(equality))) + .groupBy($"b")($"b") + + withSQLConf( + SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> "10MB", + SQLConf.NULL_AWARE_ANTI_JOIN_BROADCAST_THRESHOLD.key -> "0") { + comparePlans(Optimize.execute(originalQuery.analyze), pushedDownQuery.analyze) + } + + withSQLConf( + SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> "-1", + SQLConf.NULL_AWARE_ANTI_JOIN_BROADCAST_THRESHOLD.key -> "0") { + comparePlans(Optimize.execute(originalQuery.analyze), originalQuery.analyze) + } + } + test("Aggregate: LeftSemi join no pushdown") { val originalQuery = testRelation .groupBy($"b")($"b", sum($"c").as("sum")) diff --git a/sql/core/src/test/scala/org/apache/spark/sql/JoinSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/JoinSuite.scala index 642fcfde3420c..bfddc78bee477 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/JoinSuite.scala +++ b/sql/core/src/test/scala/org/apache/spark/sql/JoinSuite.scala @@ -1308,7 +1308,7 @@ class JoinSuite extends SharedSparkSession with AdaptiveSparkPlanHelper } } - test("SPARK-36082: left-broadcast NAAJ fallback uses nested-loop join") { + test("SPARK-36082: automatic threshold enables NAAJ hash join despite left broadcast hint") { withSQLConf( SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> "false", SQLConf.OPTIMIZE_NULL_AWARE_ANTI_JOIN.key -> "true", @@ -1332,9 +1332,9 @@ class JoinSuite extends SharedSparkSession with AdaptiveSparkPlanHelper val nullAwareHashJoins = plan.collect { case join: BroadcastHashJoinExec if join.isNullAwareAntiJoin => join } - assert(nestedLoopJoins.size === 1) - assert(nestedLoopJoins.head.buildSide === BuildLeft) - assert(nullAwareHashJoins.isEmpty) + assert(nestedLoopJoins.isEmpty) + assert(nullAwareHashJoins.size === 1) + assert(nullAwareHashJoins.head.buildSide === BuildRight) checkAnswer(result, Row(2.0d)) } } From 47fceb0e550afe43dfa7e06f915f2d5687fbc2cc Mon Sep 17 00:00:00 2001 From: Wenchen Fan Date: Fri, 18 Sep 2026 14:01:11 +0000 Subject: [PATCH 3/9] [SPARK-36082][SQL] Address follow-up review comments --- .../spark/sql/catalyst/optimizer/joins.scala | 4 +- .../apache/spark/sql/internal/SQLConf.scala | 14 +++-- .../optimizer/JoinSelectionHelperSuite.scala | 19 +++++++ .../LeftSemiAntiJoinPushDownSuite.scala | 16 ++++++ .../org/apache/spark/sql/JoinSuite.scala | 51 +++++++++++++------ 5 files changed, 82 insertions(+), 22 deletions(-) diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/joins.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/joins.scala index 46f00e453374b..189442ce8bc8c 100644 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/joins.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/joins.scala @@ -446,8 +446,8 @@ trait JoinSelectionHelper extends Logging { val canBroadcast = if (dedicatedThreshold < 0) { true } else { - val automaticBroadcastDisabled = conf.autoBroadcastJoinThreshold <= 0 && - conf.getConf(SQLConf.ADAPTIVE_AUTO_BROADCASTJOIN_THRESHOLD).forall(_ <= 0) + val automaticBroadcastDisabled = conf.autoBroadcastJoinThreshold < 0 && + conf.getConf(SQLConf.ADAPTIVE_AUTO_BROADCASTJOIN_THRESHOLD).forall(_ < 0) if (dedicatedThreshold == 0 && automaticBroadcastDisabled) { false } else { diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala index e73cea7611d3d..50bbac84be0af 100644 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala @@ -7455,15 +7455,19 @@ object SQLConf { "optimization regardless of the estimated size. For a nonnegative value, the " + "optimization is also allowed when regular join planning considers the right side " + "broadcastable. " + - "Regular planning uses spark.sql.adaptive.autoBroadcastJoinThreshold for runtime " + - "statistics when it is set, and spark.sql.autoBroadcastJoinThreshold otherwise. Thus, " + - "zero disables the optimization only when automatic broadcasting is also disabled. " + + "For join selection, regular planning uses " + + "spark.sql.adaptive.autoBroadcastJoinThreshold for runtime statistics when it is set, " + + "and spark.sql.autoBroadcastJoinThreshold otherwise. The same eligibility decision " + + "controls whether a null-aware anti join can be pushed below an aggregate; this " + + "pushdown runs before adaptive execution and therefore uses estimated statistics and " + + "spark.sql.autoBroadcastJoinThreshold. Thus, zero disables the optimization only when " + + "automatic broadcasting is also disabled. When neither threshold admits the right " + + "side, Spark falls back to regular join planning. " + "The fallback may still broadcast the right side with a nested-loop representation " + "that uses more memory and runs in O(M * N) time. Join hints do not override this " + "configuration when the broadcast hash optimization is selected. Set " + "spark.sql.optimizeNullAwareAntiJoin to false to disable the optimization without " + - "changing automatic broadcast thresholds. The same eligibility decision controls " + - "whether a null-aware anti join can be pushed below an aggregate.") + "changing automatic broadcast thresholds.") .version("4.2.1") .withBindingPolicy(ConfigBindingPolicy.NOT_APPLICABLE) .bytesConf(ByteUnit.BYTE) diff --git a/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/optimizer/JoinSelectionHelperSuite.scala b/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/optimizer/JoinSelectionHelperSuite.scala index 25119de02fde0..d3066ee6081bf 100644 --- a/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/optimizer/JoinSelectionHelperSuite.scala +++ b/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/optimizer/JoinSelectionHelperSuite.scala @@ -208,6 +208,7 @@ class JoinSelectionHelperSuite extends PlanTest with JoinSelectionHelper { rowCount = 8 * 1024 * 1024, size = Some(8 * 1024 * 1024)) val largeRight = right.copy(rowCount = 20000000, size = Some(20000000)) + val emptyRight = right.copy(rowCount = 0, size = Some(0)) withSQLConf( SQLConf.OPTIMIZE_NULL_AWARE_ANTI_JOIN.key -> "true", @@ -242,6 +243,15 @@ class JoinSelectionHelperSuite extends PlanTest with JoinSelectionHelper { SQLConf.NULL_AWARE_ANTI_JOIN_BROADCAST_THRESHOLD.key -> "0") { assert(getBroadcastHashJoinBuildSide(nullAwareAntiJoin(), SQLConf.get).isEmpty) } + + withSQLConf( + SQLConf.OPTIMIZE_NULL_AWARE_ANTI_JOIN.key -> "true", + SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> "0", + SQLConf.NULL_AWARE_ANTI_JOIN_BROADCAST_THRESHOLD.key -> "0") { + assert(getBroadcastHashJoinBuildSide( + nullAwareAntiJoin(emptyRight), SQLConf.get) === Some(BuildRight)) + assert(getBroadcastHashJoinBuildSide(nullAwareAntiJoin(), SQLConf.get).isEmpty) + } } test("NAAJ broadcast threshold is unlimited by default") { @@ -280,6 +290,15 @@ class JoinSelectionHelperSuite extends PlanTest with JoinSelectionHelper { assert(getBroadcastHashJoinBuildSide( nullAwareAntiJoin(runtimeRight), SQLConf.get) === Some(BuildRight)) } + + withSQLConf( + SQLConf.OPTIMIZE_NULL_AWARE_ANTI_JOIN.key -> "true", + SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> "-1", + SQLConf.ADAPTIVE_AUTO_BROADCASTJOIN_THRESHOLD.key -> "10MB", + SQLConf.NULL_AWARE_ANTI_JOIN_BROADCAST_THRESHOLD.key -> "0") { + assert(getBroadcastHashJoinBuildSide( + nullAwareAntiJoin(runtimeRight), SQLConf.get) === Some(BuildRight)) + } } test("NAAJ broadcast threshold short-circuits config-only decisions") { diff --git a/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/optimizer/LeftSemiAntiJoinPushDownSuite.scala b/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/optimizer/LeftSemiAntiJoinPushDownSuite.scala index a5c0536fcf63d..fe6876d41a874 100644 --- a/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/optimizer/LeftSemiAntiJoinPushDownSuite.scala +++ b/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/optimizer/LeftSemiAntiJoinPushDownSuite.scala @@ -24,6 +24,7 @@ import org.apache.spark.sql.catalyst.expressions._ import org.apache.spark.sql.catalyst.plans._ import org.apache.spark.sql.catalyst.plans.logical._ import org.apache.spark.sql.catalyst.rules._ +import org.apache.spark.sql.catalyst.statsEstimation.StatsTestPlan import org.apache.spark.sql.internal.SQLConf import org.apache.spark.sql.types.IntegerType @@ -155,6 +156,15 @@ class LeftSemiAntiJoinPushDownSuite extends PlanTest { joinType = LeftAnti, condition = Some(equality || IsNull(equality))) .groupBy($"b")($"b") + val largeRight = StatsTestPlan( + outputList = testRelation1.output, + rowCount = 20 * 1024 * 1024, + attributeStats = AttributeMap.empty, + size = Some(20 * 1024 * 1024)) + val largeRightQuery = aggregate.join( + largeRight, + joinType = LeftAnti, + condition = Some(equality || IsNull(equality))) withSQLConf( SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> "10MB", @@ -167,6 +177,12 @@ class LeftSemiAntiJoinPushDownSuite extends PlanTest { SQLConf.NULL_AWARE_ANTI_JOIN_BROADCAST_THRESHOLD.key -> "0") { comparePlans(Optimize.execute(originalQuery.analyze), originalQuery.analyze) } + + withSQLConf( + SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> "10MB", + SQLConf.NULL_AWARE_ANTI_JOIN_BROADCAST_THRESHOLD.key -> "0") { + comparePlans(Optimize.execute(largeRightQuery.analyze), largeRightQuery.analyze) + } } test("Aggregate: LeftSemi join no pushdown") { diff --git a/sql/core/src/test/scala/org/apache/spark/sql/JoinSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/JoinSuite.scala index bfddc78bee477..c5a9ddc927d9c 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/JoinSuite.scala +++ b/sql/core/src/test/scala/org/apache/spark/sql/JoinSuite.scala @@ -1308,34 +1308,55 @@ class JoinSuite extends SharedSparkSession with AdaptiveSparkPlanHelper } } - test("SPARK-36082: automatic threshold enables NAAJ hash join despite left broadcast hint") { + test("SPARK-36082: NAAJ hash eligibility takes precedence over a left broadcast hint") { withSQLConf( SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> "false", - SQLConf.OPTIMIZE_NULL_AWARE_ANTI_JOIN.key -> "true", - SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> Long.MaxValue.toString, - SQLConf.NULL_AWARE_ANTI_JOIN_BROADCAST_THRESHOLD.key -> "0") { + SQLConf.OPTIMIZE_NULL_AWARE_ANTI_JOIN.key -> "true") { withTempView("naajHintedLeft", "naajHintedRight") { Seq[java.lang.Double](-0.0d, 2.0d, null).toDF("key") .createOrReplaceTempView("naajHintedLeft") Seq[java.lang.Double](0.0d, 1.0d).toDF("key") .createOrReplaceTempView("naajHintedRight") - val result = sql( + val query = "select /*+ BROADCAST(naajHintedLeft) */ naajHintedLeft.* " + "from naajHintedLeft left anti join naajHintedRight on " + "naajHintedLeft.key = naajHintedRight.key or " + - "isnull(naajHintedLeft.key = naajHintedRight.key)") - val plan = result.queryExecution.sparkPlan - val nestedLoopJoins = plan.collect { - case join: BroadcastNestedLoopJoinExec => join + "isnull(naajHintedLeft.key = naajHintedRight.key)" + + withSQLConf( + SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> Long.MaxValue.toString, + SQLConf.NULL_AWARE_ANTI_JOIN_BROADCAST_THRESHOLD.key -> "0") { + val result = sql(query) + val plan = result.queryExecution.sparkPlan + val nestedLoopJoins = plan.collect { + case join: BroadcastNestedLoopJoinExec => join + } + val nullAwareHashJoins = plan.collect { + case join: BroadcastHashJoinExec if join.isNullAwareAntiJoin => join + } + assert(nestedLoopJoins.isEmpty) + assert(nullAwareHashJoins.size === 1) + assert(nullAwareHashJoins.head.buildSide === BuildRight) + checkAnswer(result, Row(2.0d)) } - val nullAwareHashJoins = plan.collect { - case join: BroadcastHashJoinExec if join.isNullAwareAntiJoin => join + + withSQLConf( + SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> "-1", + SQLConf.NULL_AWARE_ANTI_JOIN_BROADCAST_THRESHOLD.key -> "0") { + val result = sql(query) + val plan = result.queryExecution.sparkPlan + val nestedLoopJoins = plan.collect { + case join: BroadcastNestedLoopJoinExec => join + } + val nullAwareHashJoins = plan.collect { + case join: BroadcastHashJoinExec if join.isNullAwareAntiJoin => join + } + assert(nestedLoopJoins.size === 1) + assert(nestedLoopJoins.head.buildSide === BuildLeft) + assert(nullAwareHashJoins.isEmpty) + checkAnswer(result, Row(2.0d)) } - assert(nestedLoopJoins.isEmpty) - assert(nullAwareHashJoins.size === 1) - assert(nullAwareHashJoins.head.buildSide === BuildRight) - checkAnswer(result, Row(2.0d)) } } } From 2bac4c5f4ca64cb5e149ac68495d16fce4171af0 Mon Sep 17 00:00:00 2001 From: Wenchen Fan Date: Sun, 20 Sep 2026 12:44:53 +0000 Subject: [PATCH 4/9] [SPARK-59673][SQL] Respect left broadcast hints when flooring NAAJ threshold --- .../spark/sql/catalyst/optimizer/joins.scala | 21 ++++++--- .../apache/spark/sql/internal/SQLConf.scala | 16 ++++--- .../optimizer/JoinSelectionHelperSuite.scala | 21 ++++++++- .../LeftSemiAntiJoinPushDownSuite.scala | 9 +++- .../org/apache/spark/sql/JoinSuite.scala | 47 +++++++++---------- 5 files changed, 71 insertions(+), 43 deletions(-) diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/joins.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/joins.scala index 189442ce8bc8c..e932b552269a3 100644 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/joins.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/joins.scala @@ -438,23 +438,30 @@ trait JoinSelectionHelper extends Logging { getBroadcastBuildSide(join, hintOnly = true, conf).orElse { if (noShufflePlannedBefore) getBroadcastBuildSide(join, hintOnly = false, conf) else None } - // `JoinSelection` always builds from the right for this shape. Do not reject the hash - // optimization when regular join planning would broadcast the right side, as the fallback - // would still broadcast it with a slower nested-loop join. + // `JoinSelection` always builds from the right for this shape. A dedicated threshold that + // admits the right side takes precedence over join hints. The automatic threshold is only + // used as a floor when regular planning would also broadcast the right side; a left-only + // broadcast hint makes the fallback broadcast the left side instead. case j @ ExtractSingleColumnNullAwareAntiJoin(_, _) => val dedicatedThreshold = conf.nullAwareAntiJoinBroadcastThreshold val canBroadcast = if (dedicatedThreshold < 0) { true } else { + val fallbackBuildsRight = + !hintToBroadcastLeft(j.hint) || hintToBroadcastRight(j.hint) val automaticBroadcastDisabled = conf.autoBroadcastJoinThreshold < 0 && conf.getConf(SQLConf.ADAPTIVE_AUTO_BROADCASTJOIN_THRESHOLD).forall(_ < 0) if (dedicatedThreshold == 0 && automaticBroadcastDisabled) { + // Avoid potentially expensive statistics computation when the configurations alone + // determine the result. Both automatic thresholds must be disabled because selecting + // between them requires reading `stats.isRuntime`. false } else { - canBroadcastBySize(j.right, conf) || (dedicatedThreshold > 0 && { - val rightSize = j.right.stats.sizeInBytes - rightSize >= 0 && rightSize <= dedicatedThreshold - }) + (fallbackBuildsRight && canBroadcastBySize(j.right, conf)) || + (dedicatedThreshold > 0 && { + val rightSize = j.right.stats.sizeInBytes + rightSize >= 0 && rightSize <= dedicatedThreshold + }) } } if (canBroadcast) { diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala index 50bbac84be0af..d0e919669e66a 100644 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala @@ -7453,19 +7453,21 @@ object SQLConf { "optimization. This configuration takes effect only when " + "spark.sql.optimizeNullAwareAntiJoin is enabled. A negative value allows the " + "optimization regardless of the estimated size. For a nonnegative value, the " + - "optimization is also allowed when regular join planning considers the right side " + - "broadcastable. " + + "optimization is also allowed when regular join planning would broadcast the right " + + "side. A broadcast hint only on the left prevents this automatic-threshold floor from " + + "applying. " + "For join selection, regular planning uses " + "spark.sql.adaptive.autoBroadcastJoinThreshold for runtime statistics when it is set, " + "and spark.sql.autoBroadcastJoinThreshold otherwise. The same eligibility decision " + "controls whether a null-aware anti join can be pushed below an aggregate; this " + "pushdown runs before adaptive execution and therefore uses estimated statistics and " + - "spark.sql.autoBroadcastJoinThreshold. Thus, zero disables the optimization only when " + - "automatic broadcasting is also disabled. When neither threshold admits the right " + - "side, Spark falls back to regular join planning. " + + "spark.sql.autoBroadcastJoinThreshold. Moving the join below the aggregate may increase " + + "the number of left-side rows it evaluates. Thus, zero does not disable the optimization " + + "by itself: the right side must also be ineligible for automatic broadcasting. When " + + "neither threshold admits the right side, Spark falls back to regular join planning. " + "The fallback may still broadcast the right side with a nested-loop representation " + - "that uses more memory and runs in O(M * N) time. Join hints do not override this " + - "configuration when the broadcast hash optimization is selected. Set " + + "that runs in O(M * N) time. Join hints do not override a negative dedicated threshold " + + "or a positive dedicated threshold that admits the right side. Set " + "spark.sql.optimizeNullAwareAntiJoin to false to disable the optimization without " + "changing automatic broadcast thresholds.") .version("4.2.1") diff --git a/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/optimizer/JoinSelectionHelperSuite.scala b/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/optimizer/JoinSelectionHelperSuite.scala index d3066ee6081bf..0edf91142a45c 100644 --- a/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/optimizer/JoinSelectionHelperSuite.scala +++ b/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/optimizer/JoinSelectionHelperSuite.scala @@ -207,7 +207,9 @@ class JoinSelectionHelperSuite extends PlanTest with JoinSelectionHelper { val betweenThresholdsRight = right.copy( rowCount = 8 * 1024 * 1024, size = Some(8 * 1024 * 1024)) - val largeRight = right.copy(rowCount = 20000000, size = Some(20000000)) + val largeRight = right.copy( + rowCount = 20 * 1024 * 1024, + size = Some(20 * 1024 * 1024)) val emptyRight = right.copy(rowCount = 0, size = Some(0)) withSQLConf( @@ -217,6 +219,9 @@ class JoinSelectionHelperSuite extends PlanTest with JoinSelectionHelper { assert(getBroadcastHashJoinBuildSide(nullAwareAntiJoin(), SQLConf.get) === Some(BuildRight)) assert(getBroadcastHashJoinBuildSide( nullAwareAntiJoin(autoThresholdRight), SQLConf.get) === Some(BuildRight)) + assert(getBroadcastHashJoinBuildSide( + nullAwareAntiJoin(autoThresholdRight).copy( + hint = JoinHint(hintBroadcast, None)), SQLConf.get).isEmpty) assert(getBroadcastHashJoinBuildSide( nullAwareAntiJoin(largeRight), SQLConf.get).isEmpty) } @@ -227,6 +232,9 @@ class JoinSelectionHelperSuite extends PlanTest with JoinSelectionHelper { SQLConf.NULL_AWARE_ANTI_JOIN_BROADCAST_THRESHOLD.key -> "20MB") { assert(getBroadcastHashJoinBuildSide( nullAwareAntiJoin(largeRight), SQLConf.get) === Some(BuildRight)) + assert(getBroadcastHashJoinBuildSide( + nullAwareAntiJoin(largeRight).copy( + hint = JoinHint(hintBroadcast, None)), SQLConf.get) === Some(BuildRight)) } withSQLConf( @@ -301,7 +309,7 @@ class JoinSelectionHelperSuite extends PlanTest with JoinSelectionHelper { } } - test("NAAJ broadcast threshold short-circuits config-only decisions") { + test("NAAJ broadcast threshold does not read stats when eligibility is already known") { case class ThrowingStatsPlan() extends LeafNode { override def output: Seq[Attribute] = right.output override def computeStats(): Statistics = @@ -322,6 +330,15 @@ class JoinSelectionHelperSuite extends PlanTest with JoinSelectionHelper { SQLConf.NULL_AWARE_ANTI_JOIN_BROADCAST_THRESHOLD.key -> "0") { assert(getBroadcastHashJoinBuildSide(nullAwareAntiJoinWithoutStats, SQLConf.get).isEmpty) } + + withSQLConf( + SQLConf.OPTIMIZE_NULL_AWARE_ANTI_JOIN.key -> "true", + SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> "10MB", + SQLConf.NULL_AWARE_ANTI_JOIN_BROADCAST_THRESHOLD.key -> "0") { + assert(getBroadcastHashJoinBuildSide( + nullAwareAntiJoinWithoutStats.copy( + hint = JoinHint(hintBroadcast, None)), SQLConf.get).isEmpty) + } } test("NAAJ broadcast threshold rejects unknown sizes and respects the optimization flag") { diff --git a/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/optimizer/LeftSemiAntiJoinPushDownSuite.scala b/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/optimizer/LeftSemiAntiJoinPushDownSuite.scala index fe6876d41a874..cb018d2fc302c 100644 --- a/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/optimizer/LeftSemiAntiJoinPushDownSuite.scala +++ b/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/optimizer/LeftSemiAntiJoinPushDownSuite.scala @@ -146,13 +146,18 @@ class LeftSemiAntiJoinPushDownSuite extends PlanTest { test("Aggregate: NAAJ pushdown follows the effective broadcast threshold") { val aggregate = testRelation.groupBy($"b")($"b") val equality = $"b" === $"d" + val smallRight = StatsTestPlan( + outputList = testRelation1.output, + rowCount = 5 * 1024 * 1024, + attributeStats = AttributeMap.empty, + size = Some(5 * 1024 * 1024)) val originalQuery = aggregate.join( - testRelation1, + smallRight, joinType = LeftAnti, condition = Some(equality || IsNull(equality))) val pushedDownQuery = testRelation .join( - testRelation1, + smallRight, joinType = LeftAnti, condition = Some(equality || IsNull(equality))) .groupBy($"b")($"b") diff --git a/sql/core/src/test/scala/org/apache/spark/sql/JoinSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/JoinSuite.scala index c5a9ddc927d9c..251cad3f4993e 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/JoinSuite.scala +++ b/sql/core/src/test/scala/org/apache/spark/sql/JoinSuite.scala @@ -1308,7 +1308,7 @@ class JoinSuite extends SharedSparkSession with AdaptiveSparkPlanHelper } } - test("SPARK-36082: NAAJ hash eligibility takes precedence over a left broadcast hint") { + test("SPARK-59673: NAAJ automatic threshold floor respects a left broadcast hint") { withSQLConf( SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> "false", SQLConf.OPTIMIZE_NULL_AWARE_ANTI_JOIN.key -> "true") { @@ -1318,44 +1318,41 @@ class JoinSuite extends SharedSparkSession with AdaptiveSparkPlanHelper Seq[java.lang.Double](0.0d, 1.0d).toDF("key") .createOrReplaceTempView("naajHintedRight") - val query = - "select /*+ BROADCAST(naajHintedLeft) */ naajHintedLeft.* " + - "from naajHintedLeft left anti join naajHintedRight on " + + val querySuffix = + "naajHintedLeft.* from naajHintedLeft left anti join naajHintedRight on " + "naajHintedLeft.key = naajHintedRight.key or " + "isnull(naajHintedLeft.key = naajHintedRight.key)" + val hintedQuery = "select /*+ BROADCAST(naajHintedLeft) */ " + querySuffix + val unhintedQuery = "select " + querySuffix withSQLConf( SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> Long.MaxValue.toString, SQLConf.NULL_AWARE_ANTI_JOIN_BROADCAST_THRESHOLD.key -> "0") { - val result = sql(query) - val plan = result.queryExecution.sparkPlan - val nestedLoopJoins = plan.collect { + val hintedResult = sql(hintedQuery) + val hintedPlan = hintedResult.queryExecution.sparkPlan + val hintedNestedLoopJoins = hintedPlan.collect { case join: BroadcastNestedLoopJoinExec => join } - val nullAwareHashJoins = plan.collect { + val hintedNullAwareHashJoins = hintedPlan.collect { case join: BroadcastHashJoinExec if join.isNullAwareAntiJoin => join } - assert(nestedLoopJoins.isEmpty) - assert(nullAwareHashJoins.size === 1) - assert(nullAwareHashJoins.head.buildSide === BuildRight) - checkAnswer(result, Row(2.0d)) - } - - withSQLConf( - SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> "-1", - SQLConf.NULL_AWARE_ANTI_JOIN_BROADCAST_THRESHOLD.key -> "0") { - val result = sql(query) - val plan = result.queryExecution.sparkPlan - val nestedLoopJoins = plan.collect { + assert(hintedNestedLoopJoins.size === 1) + assert(hintedNestedLoopJoins.head.buildSide === BuildLeft) + assert(hintedNullAwareHashJoins.isEmpty) + checkAnswer(hintedResult, Row(2.0d)) + + val unhintedResult = sql(unhintedQuery) + val unhintedPlan = unhintedResult.queryExecution.sparkPlan + val unhintedNestedLoopJoins = unhintedPlan.collect { case join: BroadcastNestedLoopJoinExec => join } - val nullAwareHashJoins = plan.collect { + val unhintedNullAwareHashJoins = unhintedPlan.collect { case join: BroadcastHashJoinExec if join.isNullAwareAntiJoin => join } - assert(nestedLoopJoins.size === 1) - assert(nestedLoopJoins.head.buildSide === BuildLeft) - assert(nullAwareHashJoins.isEmpty) - checkAnswer(result, Row(2.0d)) + assert(unhintedNestedLoopJoins.isEmpty) + assert(unhintedNullAwareHashJoins.size === 1) + assert(unhintedNullAwareHashJoins.head.buildSide === BuildRight) + checkAnswer(unhintedResult, Row(2.0d)) } } } From 64b6d905ec47a7d49448872d93d071f3c0353ef2 Mon Sep 17 00:00:00 2001 From: Wenchen Fan Date: Tue, 22 Sep 2026 02:46:35 +0000 Subject: [PATCH 5/9] [SPARK-59673][SQL] Address remaining review comments --- .../spark/sql/catalyst/optimizer/joins.scala | 18 +++++---- .../apache/spark/sql/internal/SQLConf.scala | 19 +++++----- .../optimizer/JoinSelectionHelperSuite.scala | 27 ++++++++++++- .../LeftSemiAntiJoinPushDownSuite.scala | 7 ++++ .../org/apache/spark/sql/JoinSuite.scala | 38 ++++++++++++++----- 5 files changed, 82 insertions(+), 27 deletions(-) diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/joins.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/joins.scala index e932b552269a3..f655896dfa3c5 100644 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/joins.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/joins.scala @@ -439,25 +439,29 @@ trait JoinSelectionHelper extends Logging { if (noShufflePlannedBefore) getBroadcastBuildSide(join, hintOnly = false, conf) else None } // `JoinSelection` always builds from the right for this shape. A dedicated threshold that - // admits the right side takes precedence over join hints. The automatic threshold is only - // used as a floor when regular planning would also broadcast the right side; a left-only - // broadcast hint makes the fallback broadcast the left side instead. + // admits the right side takes precedence over join hints. Otherwise, use the automatic + // threshold as a floor only when regular planning selects a right-side broadcast by hint or + // size. A left-only broadcast hint or a right-only no-broadcast-and-replication hint makes the + // fallback build the left side instead. This decision also controls aggregate pushdown. case j @ ExtractSingleColumnNullAwareAntiJoin(_, _) => val dedicatedThreshold = conf.nullAwareAntiJoinBroadcastThreshold val canBroadcast = if (dedicatedThreshold < 0) { true } else { - val fallbackBuildsRight = - !hintToBroadcastLeft(j.hint) || hintToBroadcastRight(j.hint) + val rightBroadcastHint = hintToBroadcastRight(j.hint) val automaticBroadcastDisabled = conf.autoBroadcastJoinThreshold < 0 && conf.getConf(SQLConf.ADAPTIVE_AUTO_BROADCASTJOIN_THRESHOLD).forall(_ < 0) - if (dedicatedThreshold == 0 && automaticBroadcastDisabled) { + if (dedicatedThreshold == 0 && automaticBroadcastDisabled && !rightBroadcastHint) { // Avoid potentially expensive statistics computation when the configurations alone // determine the result. Both automatic thresholds must be disabled because selecting // between them requires reading `stats.isRuntime`. false } else { - (fallbackBuildsRight && canBroadcastBySize(j.right, conf)) || + val rightBroadcastSelectedByHintOrSize = + !hintToNotBroadcastAndReplicateRight(j.hint) && + (rightBroadcastHint || + (!hintToBroadcastLeft(j.hint) && canBroadcastBySize(j.right, conf))) + rightBroadcastSelectedByHintOrSize || (dedicatedThreshold > 0 && { val rightSize = j.right.stats.sizeInBytes rightSize >= 0 && rightSize <= dedicatedThreshold diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala index d0e919669e66a..54fd42a28b178 100644 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala @@ -7453,21 +7453,22 @@ object SQLConf { "optimization. This configuration takes effect only when " + "spark.sql.optimizeNullAwareAntiJoin is enabled. A negative value allows the " + "optimization regardless of the estimated size. For a nonnegative value, the " + - "optimization is also allowed when regular join planning would broadcast the right " + - "side. A broadcast hint only on the left prevents this automatic-threshold floor from " + - "applying. " + - "For join selection, regular planning uses " + + "optimization is also allowed when regular join planning selects the right side for " + + "broadcast by an explicit broadcast hint or the applicable automatic threshold. A " + + "broadcast hint only on the left, or a hint that prevents broadcasting and replicating " + + "the right side, prevents this floor from applying. For join selection, regular planning " + + "uses " + "spark.sql.adaptive.autoBroadcastJoinThreshold for runtime statistics when it is set, " + "and spark.sql.autoBroadcastJoinThreshold otherwise. The same eligibility decision " + "controls whether a null-aware anti join can be pushed below an aggregate; this " + "pushdown runs before adaptive execution and therefore uses estimated statistics and " + "spark.sql.autoBroadcastJoinThreshold. Moving the join below the aggregate may increase " + "the number of left-side rows it evaluates. Thus, zero does not disable the optimization " + - "by itself: the right side must also be ineligible for automatic broadcasting. When " + - "neither threshold admits the right side, Spark falls back to regular join planning. " + - "The fallback may still broadcast the right side with a nested-loop representation " + - "that runs in O(M * N) time. Join hints do not override a negative dedicated threshold " + - "or a positive dedicated threshold that admits the right side. Set " + + "by itself: regular planning must also not select a right-side broadcast by hint or " + + "size. Otherwise, Spark falls back to regular join planning. The fallback may still be " + + "forced to broadcast the right side with a nested-loop representation that runs in " + + "O(M * N) time. Join hints do not override a negative dedicated threshold or a positive " + + "dedicated threshold that admits the right side. Set " + "spark.sql.optimizeNullAwareAntiJoin to false to disable the optimization without " + "changing automatic broadcast thresholds.") .version("4.2.1") diff --git a/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/optimizer/JoinSelectionHelperSuite.scala b/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/optimizer/JoinSelectionHelperSuite.scala index 0edf91142a45c..ab72b1d702041 100644 --- a/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/optimizer/JoinSelectionHelperSuite.scala +++ b/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/optimizer/JoinSelectionHelperSuite.scala @@ -20,7 +20,7 @@ package org.apache.spark.sql.catalyst.optimizer import org.apache.spark.sql.catalyst.dsl.expressions._ import org.apache.spark.sql.catalyst.expressions.{Attribute, AttributeMap, EqualTo, IsNull, Or} import org.apache.spark.sql.catalyst.plans.{Inner, LeftAnti, PlanTest} -import org.apache.spark.sql.catalyst.plans.logical.{BROADCAST, HintInfo, Join, JoinHint, LeafNode, LogicalPlan, NO_BROADCAST_HASH, SHUFFLE_HASH, Statistics} +import org.apache.spark.sql.catalyst.plans.logical.{BROADCAST, HintInfo, Join, JoinHint, LeafNode, LogicalPlan, NO_BROADCAST_AND_REPLICATION, NO_BROADCAST_HASH, SHUFFLE_HASH, Statistics} import org.apache.spark.sql.catalyst.statsEstimation.StatsTestPlan import org.apache.spark.sql.internal.SQLConf @@ -46,6 +46,8 @@ class JoinSelectionHelperSuite extends PlanTest with JoinSelectionHelper { } private val hintBroadcast = Some(HintInfo(Some(BROADCAST))) + private val hintNotToBroadcastAndReplicate = + Some(HintInfo(Some(NO_BROADCAST_AND_REPLICATION))) private val hintNotToBroadcast = Some(HintInfo(Some(NO_BROADCAST_HASH))) private val hintShuffleHash = Some(HintInfo(Some(SHUFFLE_HASH))) @@ -210,6 +212,9 @@ class JoinSelectionHelperSuite extends PlanTest with JoinSelectionHelper { val largeRight = right.copy( rowCount = 20 * 1024 * 1024, size = Some(20 * 1024 * 1024)) + val overDedicatedThresholdRight = right.copy( + rowCount = 21 * 1024 * 1024, + size = Some(21 * 1024 * 1024)) val emptyRight = right.copy(rowCount = 0, size = Some(0)) withSQLConf( @@ -222,6 +227,15 @@ class JoinSelectionHelperSuite extends PlanTest with JoinSelectionHelper { assert(getBroadcastHashJoinBuildSide( nullAwareAntiJoin(autoThresholdRight).copy( hint = JoinHint(hintBroadcast, None)), SQLConf.get).isEmpty) + assert(getBroadcastHashJoinBuildSide( + nullAwareAntiJoin(autoThresholdRight).copy( + hint = JoinHint(hintBroadcast, hintBroadcast)), SQLConf.get) === Some(BuildRight)) + assert(getBroadcastHashJoinBuildSide( + nullAwareAntiJoin(largeRight).copy( + hint = JoinHint(None, hintBroadcast)), SQLConf.get) === Some(BuildRight)) + assert(getBroadcastHashJoinBuildSide( + nullAwareAntiJoin().copy( + hint = JoinHint(None, hintNotToBroadcastAndReplicate)), SQLConf.get).isEmpty) assert(getBroadcastHashJoinBuildSide( nullAwareAntiJoin(largeRight), SQLConf.get).isEmpty) } @@ -235,6 +249,8 @@ class JoinSelectionHelperSuite extends PlanTest with JoinSelectionHelper { assert(getBroadcastHashJoinBuildSide( nullAwareAntiJoin(largeRight).copy( hint = JoinHint(hintBroadcast, None)), SQLConf.get) === Some(BuildRight)) + assert(getBroadcastHashJoinBuildSide( + nullAwareAntiJoin(overDedicatedThresholdRight), SQLConf.get).isEmpty) } withSQLConf( @@ -339,6 +355,15 @@ class JoinSelectionHelperSuite extends PlanTest with JoinSelectionHelper { nullAwareAntiJoinWithoutStats.copy( hint = JoinHint(hintBroadcast, None)), SQLConf.get).isEmpty) } + + withSQLConf( + SQLConf.OPTIMIZE_NULL_AWARE_ANTI_JOIN.key -> "true", + SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> "-1", + SQLConf.NULL_AWARE_ANTI_JOIN_BROADCAST_THRESHOLD.key -> "0") { + assert(getBroadcastHashJoinBuildSide( + nullAwareAntiJoinWithoutStats.copy( + hint = JoinHint(None, hintBroadcast)), SQLConf.get) === Some(BuildRight)) + } } test("NAAJ broadcast threshold rejects unknown sizes and respects the optimization flag") { diff --git a/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/optimizer/LeftSemiAntiJoinPushDownSuite.scala b/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/optimizer/LeftSemiAntiJoinPushDownSuite.scala index cb018d2fc302c..aa34e7acd6f3d 100644 --- a/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/optimizer/LeftSemiAntiJoinPushDownSuite.scala +++ b/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/optimizer/LeftSemiAntiJoinPushDownSuite.scala @@ -155,6 +155,12 @@ class LeftSemiAntiJoinPushDownSuite extends PlanTest { smallRight, joinType = LeftAnti, condition = Some(equality || IsNull(equality))) + val leftHintedQuery = Join( + aggregate, + smallRight, + LeftAnti, + Some(equality || IsNull(equality)), + JoinHint(Some(HintInfo(Some(BROADCAST))), None)) val pushedDownQuery = testRelation .join( smallRight, @@ -175,6 +181,7 @@ class LeftSemiAntiJoinPushDownSuite extends PlanTest { SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> "10MB", SQLConf.NULL_AWARE_ANTI_JOIN_BROADCAST_THRESHOLD.key -> "0") { comparePlans(Optimize.execute(originalQuery.analyze), pushedDownQuery.analyze) + comparePlans(Optimize.execute(leftHintedQuery.analyze), leftHintedQuery.analyze) } withSQLConf( diff --git a/sql/core/src/test/scala/org/apache/spark/sql/JoinSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/JoinSuite.scala index 251cad3f4993e..4ddf96c3cef6b 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/JoinSuite.scala +++ b/sql/core/src/test/scala/org/apache/spark/sql/JoinSuite.scala @@ -1308,7 +1308,7 @@ class JoinSuite extends SharedSparkSession with AdaptiveSparkPlanHelper } } - test("SPARK-59673: NAAJ automatic threshold floor respects a left broadcast hint") { + test("SPARK-59673: NAAJ automatic threshold floor respects broadcast hints") { withSQLConf( SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> "false", SQLConf.OPTIMIZE_NULL_AWARE_ANTI_JOIN.key -> "true") { @@ -1322,24 +1322,25 @@ class JoinSuite extends SharedSparkSession with AdaptiveSparkPlanHelper "naajHintedLeft.* from naajHintedLeft left anti join naajHintedRight on " + "naajHintedLeft.key = naajHintedRight.key or " + "isnull(naajHintedLeft.key = naajHintedRight.key)" - val hintedQuery = "select /*+ BROADCAST(naajHintedLeft) */ " + querySuffix + val leftHintedQuery = "select /*+ BROADCAST(naajHintedLeft) */ " + querySuffix + val rightHintedQuery = "select /*+ BROADCAST(naajHintedRight) */ " + querySuffix val unhintedQuery = "select " + querySuffix withSQLConf( SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> Long.MaxValue.toString, SQLConf.NULL_AWARE_ANTI_JOIN_BROADCAST_THRESHOLD.key -> "0") { - val hintedResult = sql(hintedQuery) - val hintedPlan = hintedResult.queryExecution.sparkPlan - val hintedNestedLoopJoins = hintedPlan.collect { + val leftHintedResult = sql(leftHintedQuery) + val leftHintedPlan = leftHintedResult.queryExecution.sparkPlan + val leftHintedNestedLoopJoins = leftHintedPlan.collect { case join: BroadcastNestedLoopJoinExec => join } - val hintedNullAwareHashJoins = hintedPlan.collect { + val leftHintedNullAwareHashJoins = leftHintedPlan.collect { case join: BroadcastHashJoinExec if join.isNullAwareAntiJoin => join } - assert(hintedNestedLoopJoins.size === 1) - assert(hintedNestedLoopJoins.head.buildSide === BuildLeft) - assert(hintedNullAwareHashJoins.isEmpty) - checkAnswer(hintedResult, Row(2.0d)) + assert(leftHintedNestedLoopJoins.size === 1) + assert(leftHintedNestedLoopJoins.head.buildSide === BuildLeft) + assert(leftHintedNullAwareHashJoins.isEmpty) + checkAnswer(leftHintedResult, Row(2.0d)) val unhintedResult = sql(unhintedQuery) val unhintedPlan = unhintedResult.queryExecution.sparkPlan @@ -1354,6 +1355,23 @@ class JoinSuite extends SharedSparkSession with AdaptiveSparkPlanHelper assert(unhintedNullAwareHashJoins.head.buildSide === BuildRight) checkAnswer(unhintedResult, Row(2.0d)) } + + withSQLConf( + SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> "0", + SQLConf.NULL_AWARE_ANTI_JOIN_BROADCAST_THRESHOLD.key -> "0") { + val rightHintedResult = sql(rightHintedQuery) + val rightHintedPlan = rightHintedResult.queryExecution.sparkPlan + val rightHintedNestedLoopJoins = rightHintedPlan.collect { + case join: BroadcastNestedLoopJoinExec => join + } + val rightHintedNullAwareHashJoins = rightHintedPlan.collect { + case join: BroadcastHashJoinExec if join.isNullAwareAntiJoin => join + } + assert(rightHintedNestedLoopJoins.isEmpty) + assert(rightHintedNullAwareHashJoins.size === 1) + assert(rightHintedNullAwareHashJoins.head.buildSide === BuildRight) + checkAnswer(rightHintedResult, Row(2.0d)) + } } } } From 3b2405b74f0c97d17bd5e5435b378e4725bb5e52 Mon Sep 17 00:00:00 2001 From: Wenchen Fan Date: Wed, 23 Sep 2026 16:07:22 +0000 Subject: [PATCH 6/9] [SPARK-59673][SQL] Simplify NAAJ threshold floor --- .../spark/sql/catalyst/optimizer/joins.scala | 50 ++----- .../apache/spark/sql/internal/SQLConf.scala | 36 ++--- .../optimizer/JoinSelectionHelperSuite.scala | 133 ++++-------------- .../LeftSemiAntiJoinPushDownSuite.scala | 16 +-- .../org/apache/spark/sql/JoinSuite.scala | 61 +++----- 5 files changed, 81 insertions(+), 215 deletions(-) diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/joins.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/joins.scala index f655896dfa3c5..a94759ec2fe5d 100644 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/joins.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/joins.scala @@ -417,14 +417,13 @@ trait JoinSelectionHelper extends Logging { /** * The build side a broadcast hash join would use, or `None` when one is ruled out by the join - * shape or by a hint. + * shape, a hint, or its size. * - * `Some` does not promise the planner picks a broadcast hash join: a `SHUFFLE_MERGE` or - * `SHUFFLE_REPLICATE_NL` hint is tried before the sizes are consulted, join keys no hash join - * supports send it to a sort merge join, and AQE re-estimates the sizes at runtime. Within the - * broadcast decision itself this does follow the planner's precedence: a hinted broadcast first, - * a hinted shuffle hash join as a veto, then the sizes. Callers that only need to know whether a - * broadcast hash join is possible should use `canPlanAsBroadcastHashJoin`. + * For equi-joins, `Some` does not promise the planner picks a broadcast hash join: other hints, + * unsupported join keys, or AQE can select another strategy. For a single-column null-aware anti + * join, the dedicated and automatic broadcast thresholds determine eligibility before hints. + * Callers that only need to know whether a broadcast hash join is possible should use + * `canPlanAsBroadcastHashJoin`. */ def getBroadcastHashJoinBuildSide(join: Join, conf: SQLConf): Option[BuildSide] = join match { case ExtractEquiJoinKeys(_, leftKeys, rightKeys, _, _, _, _, _) => @@ -438,36 +437,17 @@ trait JoinSelectionHelper extends Logging { getBroadcastBuildSide(join, hintOnly = true, conf).orElse { if (noShufflePlannedBefore) getBroadcastBuildSide(join, hintOnly = false, conf) else None } - // `JoinSelection` always builds from the right for this shape. A dedicated threshold that - // admits the right side takes precedence over join hints. Otherwise, use the automatic - // threshold as a floor only when regular planning selects a right-side broadcast by hint or - // size. A left-only broadcast hint or a right-only no-broadcast-and-replication hint makes the - // fallback build the left side instead. This decision also controls aggregate pushdown. + // `JoinSelection` always builds from the right for this shape. The applicable automatic + // broadcast threshold floors a nonnegative dedicated threshold. As before, threshold + // eligibility takes precedence over join hints. This same decision intentionally controls + // aggregate pushdown. case j @ ExtractSingleColumnNullAwareAntiJoin(_, _) => val dedicatedThreshold = conf.nullAwareAntiJoinBroadcastThreshold - val canBroadcast = if (dedicatedThreshold < 0) { - true - } else { - val rightBroadcastHint = hintToBroadcastRight(j.hint) - val automaticBroadcastDisabled = conf.autoBroadcastJoinThreshold < 0 && - conf.getConf(SQLConf.ADAPTIVE_AUTO_BROADCASTJOIN_THRESHOLD).forall(_ < 0) - if (dedicatedThreshold == 0 && automaticBroadcastDisabled && !rightBroadcastHint) { - // Avoid potentially expensive statistics computation when the configurations alone - // determine the result. Both automatic thresholds must be disabled because selecting - // between them requires reading `stats.isRuntime`. - false - } else { - val rightBroadcastSelectedByHintOrSize = - !hintToNotBroadcastAndReplicateRight(j.hint) && - (rightBroadcastHint || - (!hintToBroadcastLeft(j.hint) && canBroadcastBySize(j.right, conf))) - rightBroadcastSelectedByHintOrSize || - (dedicatedThreshold > 0 && { - val rightSize = j.right.stats.sizeInBytes - rightSize >= 0 && rightSize <= dedicatedThreshold - }) - } - } + val canBroadcast = dedicatedThreshold < 0 || + (dedicatedThreshold > 0 && { + val rightSize = j.right.stats.sizeInBytes + rightSize >= 0 && rightSize <= dedicatedThreshold + }) || canBroadcastBySize(j.right, conf) if (canBroadcast) { Some(BuildRight) } else { diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala index 54fd42a28b178..128584fb1148e 100644 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala @@ -7448,29 +7448,19 @@ object SQLConf { val NULL_AWARE_ANTI_JOIN_BROADCAST_THRESHOLD = buildConf("spark.sql.optimizeNullAwareAntiJoin.broadcastThreshold") .internal() - .doc("Configures the maximum estimated size in bytes of the right side of a " + - "single-column null-aware anti join for which Spark uses the broadcast hash join " + - "optimization. This configuration takes effect only when " + - "spark.sql.optimizeNullAwareAntiJoin is enabled. A negative value allows the " + - "optimization regardless of the estimated size. For a nonnegative value, the " + - "optimization is also allowed when regular join planning selects the right side for " + - "broadcast by an explicit broadcast hint or the applicable automatic threshold. A " + - "broadcast hint only on the left, or a hint that prevents broadcasting and replicating " + - "the right side, prevents this floor from applying. For join selection, regular planning " + - "uses " + - "spark.sql.adaptive.autoBroadcastJoinThreshold for runtime statistics when it is set, " + - "and spark.sql.autoBroadcastJoinThreshold otherwise. The same eligibility decision " + - "controls whether a null-aware anti join can be pushed below an aggregate; this " + - "pushdown runs before adaptive execution and therefore uses estimated statistics and " + - "spark.sql.autoBroadcastJoinThreshold. Moving the join below the aggregate may increase " + - "the number of left-side rows it evaluates. Thus, zero does not disable the optimization " + - "by itself: regular planning must also not select a right-side broadcast by hint or " + - "size. Otherwise, Spark falls back to regular join planning. The fallback may still be " + - "forced to broadcast the right side with a nested-loop representation that runs in " + - "O(M * N) time. Join hints do not override a negative dedicated threshold or a positive " + - "dedicated threshold that admits the right side. Set " + - "spark.sql.optimizeNullAwareAntiJoin to false to disable the optimization without " + - "changing automatic broadcast thresholds.") + .doc(s"Configures a dedicated broadcast threshold for the right side of a single-column " + + "null-aware anti join. This configuration takes effect only when " + + s"${OPTIMIZE_NULL_AWARE_ANTI_JOIN.key} is enabled. A negative value allows the " + + "broadcast hash join optimization regardless of the estimated size. For a nonnegative " + + "value, the applicable automatic broadcast threshold acts as a floor: " + + s"${ADAPTIVE_AUTO_BROADCASTJOIN_THRESHOLD.key} is used for runtime statistics when set, " + + s"and ${AUTO_BROADCASTJOIN_THRESHOLD.key} is used otherwise. Once either threshold " + + "admits the right side, the optimization takes precedence over join hints. The same " + + "eligibility decision controls aggregate pushdown, which runs before adaptive execution " + + "and uses estimated statistics; join selection may reevaluate it with runtime " + + s"statistics. Thus, zero alone does not disable the optimization. Set " + + s"${OPTIMIZE_NULL_AWARE_ANTI_JOIN.key} to false to disable it without changing automatic " + + "broadcast thresholds.") .version("4.2.1") .withBindingPolicy(ConfigBindingPolicy.NOT_APPLICABLE) .bytesConf(ByteUnit.BYTE) diff --git a/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/optimizer/JoinSelectionHelperSuite.scala b/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/optimizer/JoinSelectionHelperSuite.scala index ab72b1d702041..b3f74687329e8 100644 --- a/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/optimizer/JoinSelectionHelperSuite.scala +++ b/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/optimizer/JoinSelectionHelperSuite.scala @@ -20,7 +20,7 @@ package org.apache.spark.sql.catalyst.optimizer import org.apache.spark.sql.catalyst.dsl.expressions._ import org.apache.spark.sql.catalyst.expressions.{Attribute, AttributeMap, EqualTo, IsNull, Or} import org.apache.spark.sql.catalyst.plans.{Inner, LeftAnti, PlanTest} -import org.apache.spark.sql.catalyst.plans.logical.{BROADCAST, HintInfo, Join, JoinHint, LeafNode, LogicalPlan, NO_BROADCAST_AND_REPLICATION, NO_BROADCAST_HASH, SHUFFLE_HASH, Statistics} +import org.apache.spark.sql.catalyst.plans.logical.{BROADCAST, HintInfo, Join, JoinHint, LeafNode, LogicalPlan, NO_BROADCAST_HASH, SHUFFLE_HASH, Statistics} import org.apache.spark.sql.catalyst.statsEstimation.StatsTestPlan import org.apache.spark.sql.internal.SQLConf @@ -46,8 +46,6 @@ class JoinSelectionHelperSuite extends PlanTest with JoinSelectionHelper { } private val hintBroadcast = Some(HintInfo(Some(BROADCAST))) - private val hintNotToBroadcastAndReplicate = - Some(HintInfo(Some(NO_BROADCAST_AND_REPLICATION))) private val hintNotToBroadcast = Some(HintInfo(Some(NO_BROADCAST_HASH))) private val hintShuffleHash = Some(HintInfo(Some(SHUFFLE_HASH))) @@ -203,9 +201,6 @@ class JoinSelectionHelperSuite extends PlanTest with JoinSelectionHelper { } test("NAAJ broadcast threshold is floored by the automatic broadcast threshold") { - val autoThresholdRight = right.copy( - rowCount = 10 * 1024 * 1024, - size = Some(10 * 1024 * 1024)) val betweenThresholdsRight = right.copy( rowCount = 8 * 1024 * 1024, size = Some(8 * 1024 * 1024)) @@ -215,7 +210,6 @@ class JoinSelectionHelperSuite extends PlanTest with JoinSelectionHelper { val overDedicatedThresholdRight = right.copy( rowCount = 21 * 1024 * 1024, size = Some(21 * 1024 * 1024)) - val emptyRight = right.copy(rowCount = 0, size = Some(0)) withSQLConf( SQLConf.OPTIMIZE_NULL_AWARE_ANTI_JOIN.key -> "true", @@ -223,42 +217,29 @@ class JoinSelectionHelperSuite extends PlanTest with JoinSelectionHelper { SQLConf.NULL_AWARE_ANTI_JOIN_BROADCAST_THRESHOLD.key -> "0") { assert(getBroadcastHashJoinBuildSide(nullAwareAntiJoin(), SQLConf.get) === Some(BuildRight)) assert(getBroadcastHashJoinBuildSide( - nullAwareAntiJoin(autoThresholdRight), SQLConf.get) === Some(BuildRight)) - assert(getBroadcastHashJoinBuildSide( - nullAwareAntiJoin(autoThresholdRight).copy( - hint = JoinHint(hintBroadcast, None)), SQLConf.get).isEmpty) - assert(getBroadcastHashJoinBuildSide( - nullAwareAntiJoin(autoThresholdRight).copy( - hint = JoinHint(hintBroadcast, hintBroadcast)), SQLConf.get) === Some(BuildRight)) + nullAwareAntiJoin().copy(hint = JoinHint(hintBroadcast, None)), + SQLConf.get) === Some(BuildRight)) assert(getBroadcastHashJoinBuildSide( nullAwareAntiJoin(largeRight).copy( - hint = JoinHint(None, hintBroadcast)), SQLConf.get) === Some(BuildRight)) - assert(getBroadcastHashJoinBuildSide( - nullAwareAntiJoin().copy( - hint = JoinHint(None, hintNotToBroadcastAndReplicate)), SQLConf.get).isEmpty) - assert(getBroadcastHashJoinBuildSide( - nullAwareAntiJoin(largeRight), SQLConf.get).isEmpty) + hint = JoinHint(None, hintBroadcast)), SQLConf.get).isEmpty) } withSQLConf( SQLConf.OPTIMIZE_NULL_AWARE_ANTI_JOIN.key -> "true", SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> "10MB", - SQLConf.NULL_AWARE_ANTI_JOIN_BROADCAST_THRESHOLD.key -> "20MB") { - assert(getBroadcastHashJoinBuildSide( - nullAwareAntiJoin(largeRight), SQLConf.get) === Some(BuildRight)) - assert(getBroadcastHashJoinBuildSide( - nullAwareAntiJoin(largeRight).copy( - hint = JoinHint(hintBroadcast, None)), SQLConf.get) === Some(BuildRight)) + SQLConf.NULL_AWARE_ANTI_JOIN_BROADCAST_THRESHOLD.key -> "5MB") { assert(getBroadcastHashJoinBuildSide( - nullAwareAntiJoin(overDedicatedThresholdRight), SQLConf.get).isEmpty) + nullAwareAntiJoin(betweenThresholdsRight), SQLConf.get) === Some(BuildRight)) } withSQLConf( SQLConf.OPTIMIZE_NULL_AWARE_ANTI_JOIN.key -> "true", SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> "10MB", - SQLConf.NULL_AWARE_ANTI_JOIN_BROADCAST_THRESHOLD.key -> "5MB") { + SQLConf.NULL_AWARE_ANTI_JOIN_BROADCAST_THRESHOLD.key -> "20MB") { assert(getBroadcastHashJoinBuildSide( - nullAwareAntiJoin(betweenThresholdsRight), SQLConf.get) === Some(BuildRight)) + nullAwareAntiJoin(largeRight), SQLConf.get) === Some(BuildRight)) + assert(getBroadcastHashJoinBuildSide( + nullAwareAntiJoin(overDedicatedThresholdRight), SQLConf.get).isEmpty) } withSQLConf( @@ -267,18 +248,9 @@ class JoinSelectionHelperSuite extends PlanTest with JoinSelectionHelper { SQLConf.NULL_AWARE_ANTI_JOIN_BROADCAST_THRESHOLD.key -> "0") { assert(getBroadcastHashJoinBuildSide(nullAwareAntiJoin(), SQLConf.get).isEmpty) } - - withSQLConf( - SQLConf.OPTIMIZE_NULL_AWARE_ANTI_JOIN.key -> "true", - SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> "0", - SQLConf.NULL_AWARE_ANTI_JOIN_BROADCAST_THRESHOLD.key -> "0") { - assert(getBroadcastHashJoinBuildSide( - nullAwareAntiJoin(emptyRight), SQLConf.get) === Some(BuildRight)) - assert(getBroadcastHashJoinBuildSide(nullAwareAntiJoin(), SQLConf.get).isEmpty) - } } - test("NAAJ broadcast threshold is unlimited by default") { + test("NAAJ broadcast threshold preserves the existing controls") { val overLongMaxRight = right.copy( rowCount = BigInt(Long.MaxValue) + 1, size = Some(BigInt(Long.MaxValue) + 1)) @@ -289,6 +261,23 @@ class JoinSelectionHelperSuite extends PlanTest with JoinSelectionHelper { assert(getBroadcastHashJoinBuildSide( nullAwareAntiJoin(overLongMaxRight), SQLConf.get) === Some(BuildRight)) } + + withSQLConf( + SQLConf.OPTIMIZE_NULL_AWARE_ANTI_JOIN.key -> "true", + SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> "-1", + SQLConf.NULL_AWARE_ANTI_JOIN_BROADCAST_THRESHOLD.key -> "10MB") { + assert(getBroadcastHashJoinBuildSide(nullAwareAntiJoin(), SQLConf.get) === Some(BuildRight)) + assert(getBroadcastHashJoinBuildSide( + nullAwareAntiJoin(right.copy(size = Some(-1))), SQLConf.get).isEmpty) + assert(getBroadcastHashJoinBuildSide( + nullAwareAntiJoin(overLongMaxRight), SQLConf.get).isEmpty) + } + + withSQLConf( + SQLConf.OPTIMIZE_NULL_AWARE_ANTI_JOIN.key -> "false", + SQLConf.NULL_AWARE_ANTI_JOIN_BROADCAST_THRESHOLD.key -> "-1") { + assert(getBroadcastHashJoinBuildSide(nullAwareAntiJoin(), SQLConf.get).isEmpty) + } } test("NAAJ broadcast threshold uses the adaptive threshold for runtime statistics") { @@ -315,71 +304,5 @@ class JoinSelectionHelperSuite extends PlanTest with JoinSelectionHelper { nullAwareAntiJoin(runtimeRight), SQLConf.get) === Some(BuildRight)) } - withSQLConf( - SQLConf.OPTIMIZE_NULL_AWARE_ANTI_JOIN.key -> "true", - SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> "-1", - SQLConf.ADAPTIVE_AUTO_BROADCASTJOIN_THRESHOLD.key -> "10MB", - SQLConf.NULL_AWARE_ANTI_JOIN_BROADCAST_THRESHOLD.key -> "0") { - assert(getBroadcastHashJoinBuildSide( - nullAwareAntiJoin(runtimeRight), SQLConf.get) === Some(BuildRight)) - } - } - - test("NAAJ broadcast threshold does not read stats when eligibility is already known") { - case class ThrowingStatsPlan() extends LeafNode { - override def output: Seq[Attribute] = right.output - override def computeStats(): Statistics = - throw new IllegalStateException("statistics should not be read") - } - val nullAwareAntiJoinWithoutStats = nullAwareAntiJoin(ThrowingStatsPlan()) - - withSQLConf( - SQLConf.OPTIMIZE_NULL_AWARE_ANTI_JOIN.key -> "true", - SQLConf.NULL_AWARE_ANTI_JOIN_BROADCAST_THRESHOLD.key -> "-2") { - assert(getBroadcastHashJoinBuildSide( - nullAwareAntiJoinWithoutStats, SQLConf.get) === Some(BuildRight)) - } - - withSQLConf( - SQLConf.OPTIMIZE_NULL_AWARE_ANTI_JOIN.key -> "true", - SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> "-1", - SQLConf.NULL_AWARE_ANTI_JOIN_BROADCAST_THRESHOLD.key -> "0") { - assert(getBroadcastHashJoinBuildSide(nullAwareAntiJoinWithoutStats, SQLConf.get).isEmpty) - } - - withSQLConf( - SQLConf.OPTIMIZE_NULL_AWARE_ANTI_JOIN.key -> "true", - SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> "10MB", - SQLConf.NULL_AWARE_ANTI_JOIN_BROADCAST_THRESHOLD.key -> "0") { - assert(getBroadcastHashJoinBuildSide( - nullAwareAntiJoinWithoutStats.copy( - hint = JoinHint(hintBroadcast, None)), SQLConf.get).isEmpty) - } - - withSQLConf( - SQLConf.OPTIMIZE_NULL_AWARE_ANTI_JOIN.key -> "true", - SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> "-1", - SQLConf.NULL_AWARE_ANTI_JOIN_BROADCAST_THRESHOLD.key -> "0") { - assert(getBroadcastHashJoinBuildSide( - nullAwareAntiJoinWithoutStats.copy( - hint = JoinHint(None, hintBroadcast)), SQLConf.get) === Some(BuildRight)) - } - } - - test("NAAJ broadcast threshold rejects unknown sizes and respects the optimization flag") { - withSQLConf( - SQLConf.OPTIMIZE_NULL_AWARE_ANTI_JOIN.key -> "true", - SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> "-1", - SQLConf.NULL_AWARE_ANTI_JOIN_BROADCAST_THRESHOLD.key -> "10MB") { - assert(getBroadcastHashJoinBuildSide(nullAwareAntiJoin(), SQLConf.get) === Some(BuildRight)) - assert(getBroadcastHashJoinBuildSide( - nullAwareAntiJoin(right.copy(size = Some(-1))), SQLConf.get).isEmpty) - } - - withSQLConf( - SQLConf.OPTIMIZE_NULL_AWARE_ANTI_JOIN.key -> "false", - SQLConf.NULL_AWARE_ANTI_JOIN_BROADCAST_THRESHOLD.key -> "-1") { - assert(getBroadcastHashJoinBuildSide(nullAwareAntiJoin(), SQLConf.get).isEmpty) - } } } diff --git a/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/optimizer/LeftSemiAntiJoinPushDownSuite.scala b/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/optimizer/LeftSemiAntiJoinPushDownSuite.scala index aa34e7acd6f3d..acc03addbb051 100644 --- a/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/optimizer/LeftSemiAntiJoinPushDownSuite.scala +++ b/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/optimizer/LeftSemiAntiJoinPushDownSuite.scala @@ -155,12 +155,6 @@ class LeftSemiAntiJoinPushDownSuite extends PlanTest { smallRight, joinType = LeftAnti, condition = Some(equality || IsNull(equality))) - val leftHintedQuery = Join( - aggregate, - smallRight, - LeftAnti, - Some(equality || IsNull(equality)), - JoinHint(Some(HintInfo(Some(BROADCAST))), None)) val pushedDownQuery = testRelation .join( smallRight, @@ -178,23 +172,19 @@ class LeftSemiAntiJoinPushDownSuite extends PlanTest { condition = Some(equality || IsNull(equality))) withSQLConf( + SQLConf.OPTIMIZE_NULL_AWARE_ANTI_JOIN.key -> "true", SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> "10MB", SQLConf.NULL_AWARE_ANTI_JOIN_BROADCAST_THRESHOLD.key -> "0") { comparePlans(Optimize.execute(originalQuery.analyze), pushedDownQuery.analyze) - comparePlans(Optimize.execute(leftHintedQuery.analyze), leftHintedQuery.analyze) + comparePlans(Optimize.execute(largeRightQuery.analyze), largeRightQuery.analyze) } withSQLConf( + SQLConf.OPTIMIZE_NULL_AWARE_ANTI_JOIN.key -> "true", SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> "-1", SQLConf.NULL_AWARE_ANTI_JOIN_BROADCAST_THRESHOLD.key -> "0") { comparePlans(Optimize.execute(originalQuery.analyze), originalQuery.analyze) } - - withSQLConf( - SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> "10MB", - SQLConf.NULL_AWARE_ANTI_JOIN_BROADCAST_THRESHOLD.key -> "0") { - comparePlans(Optimize.execute(largeRightQuery.analyze), largeRightQuery.analyze) - } } test("Aggregate: LeftSemi join no pushdown") { diff --git a/sql/core/src/test/scala/org/apache/spark/sql/JoinSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/JoinSuite.scala index 4ddf96c3cef6b..a822d364e8782 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/JoinSuite.scala +++ b/sql/core/src/test/scala/org/apache/spark/sql/JoinSuite.scala @@ -27,7 +27,7 @@ import org.apache.spark.internal.config.SHUFFLE_SPILL_NUM_ELEMENTS_FORCE_SPILL_T import org.apache.spark.sql.catalyst.TableIdentifier import org.apache.spark.sql.catalyst.analysis.UnresolvedRelation import org.apache.spark.sql.catalyst.expressions.{Ascending, GenericRow, SortOrder} -import org.apache.spark.sql.catalyst.optimizer.{BuildLeft, BuildRight, JoinSelectionHelper} +import org.apache.spark.sql.catalyst.optimizer.{BuildLeft, BuildRight, BuildSide, JoinSelectionHelper} import org.apache.spark.sql.catalyst.plans.logical.{Filter, HintInfo, Join, JoinHint, NO_BROADCAST_AND_REPLICATION} import org.apache.spark.sql.execution.{BinaryExecNode, FilterExec, ProjectExec, SortExec, SparkPlan, WholeStageCodegenExec} import org.apache.spark.sql.execution.adaptive.AdaptiveSparkPlanHelper @@ -1308,7 +1308,7 @@ class JoinSuite extends SharedSparkSession with AdaptiveSparkPlanHelper } } - test("SPARK-59673: NAAJ automatic threshold floor respects broadcast hints") { + test("SPARK-36082, SPARK-59673: NAAJ threshold eligibility precedes broadcast hints") { withSQLConf( SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> "false", SQLConf.OPTIMIZE_NULL_AWARE_ANTI_JOIN.key -> "true") { @@ -1326,51 +1326,34 @@ class JoinSuite extends SharedSparkSession with AdaptiveSparkPlanHelper val rightHintedQuery = "select /*+ BROADCAST(naajHintedRight) */ " + querySuffix val unhintedQuery = "select " + querySuffix + def checkNAAJ( + query: String, + expectedClass: Class[_ <: BinaryExecNode], + expectedBuildSide: BuildSide): Unit = { + assertJoin((query, expectedClass)) match { + case join: BroadcastHashJoinExec => + assert(join.isNullAwareAntiJoin) + assert(join.buildSide === expectedBuildSide) + case join: BroadcastNestedLoopJoinExec => + assert(join.buildSide === expectedBuildSide) + case join => + fail(s"Unexpected join: $join") + } + checkAnswer(sql(query), Row(2.0d)) + } + withSQLConf( SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> Long.MaxValue.toString, SQLConf.NULL_AWARE_ANTI_JOIN_BROADCAST_THRESHOLD.key -> "0") { - val leftHintedResult = sql(leftHintedQuery) - val leftHintedPlan = leftHintedResult.queryExecution.sparkPlan - val leftHintedNestedLoopJoins = leftHintedPlan.collect { - case join: BroadcastNestedLoopJoinExec => join - } - val leftHintedNullAwareHashJoins = leftHintedPlan.collect { - case join: BroadcastHashJoinExec if join.isNullAwareAntiJoin => join - } - assert(leftHintedNestedLoopJoins.size === 1) - assert(leftHintedNestedLoopJoins.head.buildSide === BuildLeft) - assert(leftHintedNullAwareHashJoins.isEmpty) - checkAnswer(leftHintedResult, Row(2.0d)) - - val unhintedResult = sql(unhintedQuery) - val unhintedPlan = unhintedResult.queryExecution.sparkPlan - val unhintedNestedLoopJoins = unhintedPlan.collect { - case join: BroadcastNestedLoopJoinExec => join - } - val unhintedNullAwareHashJoins = unhintedPlan.collect { - case join: BroadcastHashJoinExec if join.isNullAwareAntiJoin => join - } - assert(unhintedNestedLoopJoins.isEmpty) - assert(unhintedNullAwareHashJoins.size === 1) - assert(unhintedNullAwareHashJoins.head.buildSide === BuildRight) - checkAnswer(unhintedResult, Row(2.0d)) + checkNAAJ(leftHintedQuery, classOf[BroadcastHashJoinExec], BuildRight) + checkNAAJ(unhintedQuery, classOf[BroadcastHashJoinExec], BuildRight) } withSQLConf( SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> "0", SQLConf.NULL_AWARE_ANTI_JOIN_BROADCAST_THRESHOLD.key -> "0") { - val rightHintedResult = sql(rightHintedQuery) - val rightHintedPlan = rightHintedResult.queryExecution.sparkPlan - val rightHintedNestedLoopJoins = rightHintedPlan.collect { - case join: BroadcastNestedLoopJoinExec => join - } - val rightHintedNullAwareHashJoins = rightHintedPlan.collect { - case join: BroadcastHashJoinExec if join.isNullAwareAntiJoin => join - } - assert(rightHintedNestedLoopJoins.isEmpty) - assert(rightHintedNullAwareHashJoins.size === 1) - assert(rightHintedNullAwareHashJoins.head.buildSide === BuildRight) - checkAnswer(rightHintedResult, Row(2.0d)) + checkNAAJ(leftHintedQuery, classOf[BroadcastNestedLoopJoinExec], BuildLeft) + checkNAAJ(rightHintedQuery, classOf[BroadcastNestedLoopJoinExec], BuildRight) } } } From b3a7fcbd2329b24c8849d0212837c09dcc6f3b4c Mon Sep 17 00:00:00 2001 From: Wenchen Fan Date: Thu, 24 Sep 2026 02:45:04 +0000 Subject: [PATCH 7/9] [SPARK-59673][SQL] Avoid extra projection in NAAJ join test --- sql/core/src/test/scala/org/apache/spark/sql/JoinSuite.scala | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/sql/core/src/test/scala/org/apache/spark/sql/JoinSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/JoinSuite.scala index a822d364e8782..a9d085e889f74 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/JoinSuite.scala +++ b/sql/core/src/test/scala/org/apache/spark/sql/JoinSuite.scala @@ -1319,7 +1319,7 @@ class JoinSuite extends SharedSparkSession with AdaptiveSparkPlanHelper .createOrReplaceTempView("naajHintedRight") val querySuffix = - "naajHintedLeft.* from naajHintedLeft left anti join naajHintedRight on " + + "* from naajHintedLeft left anti join naajHintedRight on " + "naajHintedLeft.key = naajHintedRight.key or " + "isnull(naajHintedLeft.key = naajHintedRight.key)" val leftHintedQuery = "select /*+ BROADCAST(naajHintedLeft) */ " + querySuffix From 7c4829dbb423cf0f3734d62effdd3315cccb3961 Mon Sep 17 00:00:00 2001 From: Wenchen Fan Date: Thu, 24 Sep 2026 03:44:55 +0000 Subject: [PATCH 8/9] [SPARK-59673][SQL] Clarify NAAJ threshold planning tradeoffs --- .../scala/org/apache/spark/sql/catalyst/optimizer/joins.scala | 4 +++- .../main/scala/org/apache/spark/sql/internal/SQLConf.scala | 3 ++- 2 files changed, 5 insertions(+), 2 deletions(-) diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/joins.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/joins.scala index a94759ec2fe5d..b1cb873509fcd 100644 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/joins.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/joins.scala @@ -440,7 +440,9 @@ trait JoinSelectionHelper extends Logging { // `JoinSelection` always builds from the right for this shape. The applicable automatic // broadcast threshold floors a nonnegative dedicated threshold. As before, threshold // eligibility takes precedence over join hints. This same decision intentionally controls - // aggregate pushdown. + // aggregate pushdown. If neither threshold admits the hash join, regular planning may still + // broadcast the right side for a nested-loop join. The thresholds limit hash relation + // construction, not all broadcasts. case j @ ExtractSingleColumnNullAwareAntiJoin(_, _) => val dedicatedThreshold = conf.nullAwareAntiJoinBroadcastThreshold val canBroadcast = dedicatedThreshold < 0 || diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala index 128584fb1148e..d9b283cd464e9 100644 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala @@ -7458,7 +7458,8 @@ object SQLConf { "admits the right side, the optimization takes precedence over join hints. The same " + "eligibility decision controls aggregate pushdown, which runs before adaptive execution " + "and uses estimated statistics; join selection may reevaluate it with runtime " + - s"statistics. Thus, zero alone does not disable the optimization. Set " + + "statistics. A lower adaptive threshold can leave a pushed-down join using a " + + s"nested-loop plan. Thus, zero alone does not disable the optimization. Set " + s"${OPTIMIZE_NULL_AWARE_ANTI_JOIN.key} to false to disable it without changing automatic " + "broadcast thresholds.") .version("4.2.1") From e964feaf011445a11b1ac7fd37a9b8ef1fd1f68d Mon Sep 17 00:00:00 2001 From: Wenchen Fan Date: Thu, 24 Sep 2026 11:49:27 +0000 Subject: [PATCH 9/9] [SPARK-59673][SQL] Handle projected plans in join test helper --- .../src/test/scala/org/apache/spark/sql/JoinSuite.scala | 7 +++++-- 1 file changed, 5 insertions(+), 2 deletions(-) diff --git a/sql/core/src/test/scala/org/apache/spark/sql/JoinSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/JoinSuite.scala index a9d085e889f74..c7f9a9cb1b0cc 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/JoinSuite.scala +++ b/sql/core/src/test/scala/org/apache/spark/sql/JoinSuite.scala @@ -66,6 +66,9 @@ class JoinSuite extends SharedSparkSession with AdaptiveSparkPlanHelper val c = pair._2 val df = sql(sqlString) val optimized = df.queryExecution.optimizedPlan + val optimizedJoin = optimized.collectFirst { case join: Join => join }.getOrElse { + fail(s"No join in optimized plan:\n$optimized") + } val physical = df.queryExecution.sparkPlan val operators = physical.collect { case j: BroadcastHashJoinExec => j @@ -80,13 +83,13 @@ class JoinSuite extends SharedSparkSession with AdaptiveSparkPlanHelper fail(s"$sqlString expected operator: $c, but got ${operators.head}\n physical: \n$physical") } assert( - canPlanAsBroadcastHashJoin(optimized.asInstanceOf[Join], conf) === + canPlanAsBroadcastHashJoin(optimizedJoin, conf) === operators.head.isInstanceOf[BroadcastHashJoinExec], "canPlanAsBroadcastHashJoin not in sync with join selection codepath!") operators.head match { case bhj: BroadcastHashJoinExec => assert( - getBroadcastHashJoinBuildSide(optimized.asInstanceOf[Join], conf) + getBroadcastHashJoinBuildSide(optimizedJoin, conf) .contains(bhj.buildSide), "getBroadcastHashJoinBuildSide not in sync with join selection codepath!") case _ =>