Skip to content

fix: choose the forced shuffled hash join at planning time so Spark keeps its sorts - #6785

Open
dwsmith1983 wants to merge 5 commits into
apache:mainfrom
dwsmith1983:fix/6770-force-shj-planner-strategy
Open

dwsmith1983 wants to merge 5 commits into
apache:mainfrom
dwsmith1983:fix/6770-force-shj-planner-strategy

Conversation

@dwsmith1983

@dwsmith1983 dwsmith1983 commented Oct 8, 2026 •

Copy link
Copy Markdown
Contributor

Which issue does this PR close?

Closes #6770.

Rationale for this change

With spark.comet.exec.forceShuffledHashJoin on, RewriteJoin turned sort-merge joins into hash joins inside CometExecRule, after EnsureRequirements and RemoveRedundantSorts had run. When a sort-merge join's output ordering already satisfied a sort above it, Spark had removed that sort as redundant, and once the join became a hash join nothing put it back. sortWithinPartitions and SORT BY on the join key returned unsorted partitions, and a partitioned insert keyed on the join key failed with FileAlreadyExistsException, because the planned write's partition sort was the one removed. #6384 restores orderings an operator above still requires, which covers a kept sort-merge join or a window, but here nothing left in the plan requires the ordering.

What changes are included in this PR?

  • RewriteJoin is replaced by CometShuffledHashJoinStrategy, injected with injectPlannerStrategy. Injected strategies run before Spark's JoinSelection, so the join is a ShuffledHashJoinExec from the start and EnsureRequirements and RemoveRedundantSorts place and keep the sorts themselves. The strategy also runs on every AQE re-plan, where the join's children carry the materialized shuffle sizes, as the size check did before.
  • It plans a hash join only where Spark would otherwise plan a sort-merge join, and leaves the join to Spark when:
    • Spark would broadcast it, either by size (canPlanAsBroadcastHashJoin, with AQE's threshold for runtime sizes) or because AQE already planned a broadcast stage for one side. The second check is needed because Spark's LogicalQueryStageStrategy, which plans a join over a broadcast stage as a broadcast join, runs after injected strategies, and AQE throws away a re-plan that hash joins a broadcast stage together with everything else that re-plan changed;
    • either side has a join strategy hint (BROADCAST, MERGE, SHUFFLE_HASH, SHUFFLE_REPLICATE_NL), including a SHUFFLE_HASH that AQE sets on a re-plan, in which case Spark plans its own shuffled hash join. AQE's NO_BROADCAST_HASH and PREFER_SHUFFLE_HASH do not count, so a side AQE only stops from broadcasting stays a hash join;
    • the keys cannot be sorted, so there was no sort-merge join to replace;
    • on Spark 4, the keys are not binary stable under their collation (hashJoinSupported), or the join is a LeftSingle, which Spark never plans as a sort-merge join;
    • the build side is over maxBuildSize, or it is LeftSemi or ExistenceJoin building on the right (correctness: TPCDS Q69 has incorrect output between Comet and Spark #2667, feat: Investigate potential issues with HashJoin and LeftSemi #2697). Spark then plans the sort-merge join, and the strategy records the same explain info or fallback reason as before on it.
  • The build side choice and the size limit from perf: keep the sort-merge join when the forced hash join's build side is too large #6384 move into the strategy unchanged. restoreRequiredOrdering is gone, since EnsureRequirements now adds those sorts.
  • The strategy does nothing when Comet or spark.comet.exec.enabled is off, in plan-only mode, or for a streaming join.

Behavior changes:

  • A join strategy hint now wins over forceShuffledHashJoin. Before, a MERGE hint was converted too.
  • On Spark 4, joins on collated keys that are not binary stable are no longer forced into Spark's hash join.
  • The plan-only report shows the sort-merge joins Spark plans, not the hash joins a real run would use; understanding-comet-plans.md says so.
  • Joins of static tables inside a streaming query are planned by the session's strategies too, so with the force config on they run as Spark's own ShuffledHashJoin (Comet leaves streaming plans alone). The tuning guide says so.

branch-1.1 has the same rewrite, so 1.1.0 loses these sorts too. A backport needs this strategy rather than a cherry-pick, since branch-1.1 has no size limit.

How are these changes tested?

CometJoinSuite, with AQE off and on where it matters:

  • sortWithinPartitions on the join key keeps every partition sorted, for inner, right outer (building on the smaller left side) and full outer joins, and the rows match Spark. SORT BY plans the same partition-local sort.
  • A partitioned insert keyed on the join key runs a native hash join, writes all 1000 rows, and writes one file per partition.
  • A streamed side already partitioned and sorted on the key keeps its order through the hash join when no exchange is added.
  • MERGE, SHUFFLE_HASH (Spark builds the hinted side) and BROADCAST hints are respected.
  • A join whose side AQE marks NO_BROADCAST_HASH on a re-plan (1 of its 10 shuffle partitions is non-empty) stays a native hash join with no sort-merge join in the final plan. The test records the join hints of the Comet run to show AQE set the hint.
  • A runtime broadcast under AQE stays a broadcast. A join Spark broadcasts from the static estimate also stays one after AQE re-plans with the adaptive threshold off; its 10 KB threshold fits only the small side, so only the broadcast stage check keeps the strategy off the join, and the test asserts AQE logs no Invalid broadcast query stage re-plan failure.
  • A skewed join is split by OptimizeSkewedJoin.
  • On Spark 4, a join on a computed key with a UTF8_LCASE collation keeps the sort-merge join and returns Spark's rows.
  • The strategy plans no hash join with Comet off, with spark.comet.exec.enabled off, or in plan-only mode.
  • The perf: keep the sort-merge join when the forced hash join's build side is too large #6384 size limit tests now run without the MERGE hint they relied on, and an RDD-backed build side, whose size Spark reports as spark.sql.defaultSizeInBytes, keeps the sort-merge join.
  • The Kept sort-merge join reads unsorted input after replaceSortMergeJoin rewrites the join below it #6673 tests for a kept sort-merge join, LeftSemi, ExistenceJoin or window above a forced hash join pass with the sorts EnsureRequirements adds.

The new sort and insert tests fail on main. Removing the plan-only check fails its test, replacing the Spark 4 hashJoinSupported check with true fails the collation test, dropping the broadcast stage check fails the planned broadcast test, and counting every join hint as a strategy hint fails the NO_BROADCAST_HASH test. CometJoinSuite and CometConfSuite pass on Spark 3.4, 3.5, 4.0 and 4.1 (the collation test runs on Spark 4 only), Spark 4.2 compiles, and the strict warnings build passes.

…eeps its sorts

With spark.comet.exec.forceShuffledHashJoin on, RewriteJoin turned sort-merge
joins into hash joins inside CometExecRule, after EnsureRequirements and
RemoveRedundantSorts had run. A sort Spark had dropped because the sort-merge
join already produced that order was then gone for good, so
sortWithinPartitions and SORT BY returned unsorted partitions and partitioned
inserts failed with FileAlreadyExistsException.

The choice now happens in a planner strategy injected ahead of Spark's
JoinSelection, so Spark places and keeps every sort the plan needs. It also
runs on each AQE re-plan with the materialized shuffle sizes. The build side
size limit and the LeftSemi and ExistenceJoin exclusions carry over. Joins Spark
would broadcast, joins with a strategy hint, keys that cannot be sorted, and on
Spark 4 collated keys and LeftSingle joins are left to Spark.

Closes apache#6770.
@github-actions github-actions Bot added bug Something isn't working area:joins Join operators and dynamic filter pushdown labels Oct 8, 2026
@dwsmith1983

Copy link
Copy Markdown
Contributor Author

@andygrove I'd suggest this for backport-1.1. branch-1.1 has the same RewriteJoin, so with spark.comet.exec.forceShuffledHashJoin on, 1.1.0 drops the same sorts: sortWithinPartitions and SORT BY on the join key return unsorted partitions, and a partitioned insert keyed on the join key fails with FileAlreadyExistsException. The branch has no #6384 size limit, so the backport would bring the strategy over without the size check rather than cherry-pick this as is.

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Reviewed the full 10-file diff from 412468e2817207eff0c1907d1c4aac8c3f06680d to 4ffdeba6c9089e0ccf09443f028766eb63aa543c. The PR remains non-draft. No introduced P1/P2 issues found within this review. The existing backport discussion contains no unresolved substantive blocker.

Summary

  • Prior state and problem: RewriteJoin replaced sort-merge joins after Spark had removed redundant sorts. This could leave sortWithinPartitions results unordered and break partitioned inserts.
  • Design approach: Select ShuffledHashJoinExec through an injected planner strategy, allowing Spark's existing requirement and sort-removal rules to operate on the correct join type.
  • Correctness: Checked join extraction, residual conditions, build-side eligibility, size boundaries, broadcast deferral, hints, and ordering propagation against Spark sources. Spark's preparation and AQE replanning sequences support the approach. DataFusion's probe-order contract also matches the relevant Spark hash-join ordering claims. No reproducible P1/P2 correctness issue was identified.
  • Compatibility analysis: Compared JoinSelectionHelper and SparkStrategies across Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3, and 4.2.0. The shared helpers retain version-specific build-side rules. The Spark 4 shim correctly excludes LeftSingle and keys lacking binary-stable equality.
  • Key design decisions: Broadcasts and explicit strategy hints defer to Spark. The existing build-size limit, smaller-side selection, and LeftSemi/ExistenceJoin restrictions remain. Declined conversions retain explain information.
  • Implementation sketch: Register CometShuffledHashJoinStrategy, remove the late rewrite and ordering-repair pass, add two version-family shims, and expand join regression coverage for sorting, writes, hints, AQE, collations, and disabled modes.
  • Performance: The native join algorithm is unchanged. The change removes a whole-plan rewrite and ordering-repair traversal, while adding planning-time checks per eligible join. Required output sorts restore correctness. No measured speedup or reproducible performance regression was established.
  • Design: Choosing the join before Spark establishes physical requirements addresses the cause of lost sorts. It is simpler to reason about than reconstructing ordering after Spark has erased explicit sorts.
  • Abstraction & complexity: The strategy keeps selection policy together, and the small shims isolate genuine Spark-version differences. Reusing Spark's selection helpers avoids duplicating its broadcast and build-side eligibility rules. No complexity concern reached the P1/P2 bar.
  • Behavioral changes worth calling out: Strategy hints now win over forcing, non-binary-stable collated keys stay with Spark, plan-only reports omit forced conversion, and eligible static joins inside streaming queries can use Spark's hash join. These are documented. Comparison with branch-1.1 at e9efd9f764ee0a59b7898ff028d6985d4a7a28e1 confirms the intended ordering fix. The size limit and ExistenceJoin guard already existed in this PR's base.
  • Suggested improvements: No additional code change meets the P1/P2 reporting bar. Obtain the pending CI and Spark SQL/profile verdicts before queueing this planner change.

Routed skills: Read AGENTS.md and applied .ai/skills/review-comet-pr/SKILL.md. No sibling subsystem skill applies to the changed files. Read the supplied discussion snapshot and checked live reviews and inline comments. No Copilot feedback was used.

Exact-head CI: Label pull requests passed. Comet CI and CodeQL both report action_required. There is no exact-head CI build or test verdict.

Validation: The exact-head Spark 4.1.3 JVM reactor test-compile passed with formatting/style checks skipped. CometConfSuite passed all 21 tests. A disposable planner harness passed 10 checks against the freshly compiled classes, covering build sides, hints, broadcast deferral, the size boundary, disabled Comet, and retention of a local sort through Spark preparation. git diff --check passed.

Validation limits: The planner harness directly exercised planJoin and Spark preparation. It did not execute native joins. Native CometJoinSuite, AQE runtime behavior, partitioned writes, streaming execution, Spark SQL suites, and other-profile builds were not independently run. No matching native build was available, and shared disk exhaustion constrained further preparation. The author's broader test claims were not treated as independently verified results. Tracked project files remain unchanged. Nothing was published.

@comphead comphead left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks @dwsmith1983. Choosing the join at planning time addresses the cause of the lost sorts, which looks better to me than repairing them afterwards. The inline comments cover reusing Spark's own strategy hint set, and trimming repeated setup and an equivalent case in the new tests.

Comment on lines +112 to +114
Seq(hint.leftHint, hint.rightHint).flatten.flatMap(_.strategy).exists {
case BROADCAST | SHUFFLE_MERGE | SHUFFLE_HASH | SHUFFLE_REPLICATE_NL => true
case _ => false

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Spark already keeps this exact list as JoinStrategyHint.strategies (BROADCAST, SHUFFLE_MERGE, SHUFFLE_HASH, SHUFFLE_REPLICATE_NL, without the AQE-only hints). It is public in hints.scala on 3.4.3, 3.5.8, 4.0.1 and 4.1.3, and ResolveHints uses it to decide which hint names a user can write. Could hasStrategyHint use it, for example Seq(hint.leftHint, hint.rightHint).exists(_.exists(_.strategy.exists(JoinStrategyHint.strategies.contains)))? That drops the four constant imports and the match, and walks the hints once instead of building two intermediate collections. If Spark adds another user-facing hint later, the join would also keep honoring it without another edit here. I have not compiled the suggested form.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Could hasStrategyHint use it

Done. It now checks JoinStrategyHint.strategies, so the four constant imports and the match are gone. I also added a test that a side AQE marks NO_BROADCAST_HASH keeps the hash join, since that hint is outside the set and nothing covered it.

Comment on lines +1256 to +1263
private def withSortLossConf(adaptive: Boolean)(f: => Unit): Unit = withSQLConf(
SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> adaptive.toString,
SQLConf.ADAPTIVE_AUTO_BROADCASTJOIN_THRESHOLD.key -> "-1",
SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> "-1",
SQLConf.SHUFFLE_PARTITIONS.key -> "2",
CometConf.COMET_FORCE_SHJ.key -> "true") {
withParquetTable((0 until 10000).map(i => (i % 100, i)), "big") {
withParquetTable((0 until 10).map(i => (i * 10, i)), "small") {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This repeats setup the suite already has. big and small are the same data that withChainedJoinTables builds, and the five configs here are the ones withKeptSemiJoinConf and the window test already set (they also set maxBuildSize=-1). Could this helper call withChainedJoinTables(midRows = 3000, midKeys = 100) for the tables, and could the config block become one class-level helper that takes adaptive and serves the semi join, window and sort tests? I have not run it. The unused mid table and maxBuildSize=-1 should not change these plans, since the build side is always small.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

could the config block become one class-level helper that takes adaptive and serves the semi join, window and sort tests?

Done. The semi join, window and sort tests now share one helper that takes adaptive, sets the five configs plus maxBuildSize=-1, and builds the tables with withChainedJoinTables(midRows = 3000, midKeys = 100).

Comment on lines +1301 to +1303
test(s"forceShuffledHashJoin keeps SORT BY on the join key, AQE=$adaptive") {
withSortLossConf(adaptive) {
checkPartitionsSortedOverHashJoin(sql(s"$sortLossJoin SORT BY k"))

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

sql(sortLossJoin).sortWithinPartitions("k") and sql(s"$sortLossJoin SORT BY k") should analyze to the same Sort with global = false over the same projection and join. In Spark 3.5.8 AstBuilder builds SORT BY as Sort(..., global = false, query) and Dataset.sortInternal builds the same node. From reading the code I expect identical plans and identical assertions here, but I have not compared the plans. The issue lists both forms, so I see why both are here. Would one of them be enough, given that the second adds a run per AQE setting without a different plan?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Would one of them be enough, given that the second adds a run per AQE setting without a different plan?

Yes. Both forms analyze to the same Sort with global = false, so I dropped the SORT BY test and kept the sortWithinPartitions one.

@andygrove andygrove added the backport-1.1 Candidate for backporting to 1.1 release branch label Oct 9, 2026
@andygrove

Copy link
Copy Markdown
Member

Thanks @dwsmith1983. I added the backport label. I won't review yet since there are existing review comments to be addressed.

…broadcast stage check

hasStrategyHint now uses Spark's own set of strategy hints. The forced
hash join tests share one helper, and the SORT BY case is dropped since
it plans the same local sort as sortWithinPartitions.

New coverage: a side AQE marks NO_BROADCAST_HASH keeps the hash join,
and the planned broadcast test now fails without the broadcast stage
check, which keeps AQE from discarding a re-plan that hash joins the
stage.
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:joins Join operators and dynamic filter pushdown backport-1.1 Candidate for backporting to 1.1 release branch bug Something isn't working

Projects

None yet

Development

Successfully merging this pull request may close these issues.

forceShuffledHashJoin leaves sortWithinPartitions output unsorted and fails partitioned inserts

4 participants