Repository navigation
fix: choose the forced shuffled hash join at planning time so Spark keeps its sorts - #6785
dwsmith1983 wants to merge 5 commits into
Conversation
…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.
|
@andygrove I'd suggest this for |
sunchao
left a comment
There was a problem hiding this comment.
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:
RewriteJoinreplaced sort-merge joins after Spark had removed redundant sorts. This could leavesortWithinPartitionsresults unordered and break partitioned inserts. - Design approach: Select
ShuffledHashJoinExecthrough 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
JoinSelectionHelperandSparkStrategiesacross Spark3.4.3,3.5.9,4.0.4,4.1.3, and4.2.0. The shared helpers retain version-specific build-side rules. The Spark 4 shim correctly excludesLeftSingleand 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/ExistenceJoinrestrictions 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.1ate9efd9f764ee0a59b7898ff028d6985d4a7a28e1confirms the intended ordering fix. The size limit andExistenceJoinguard 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
left a comment
There was a problem hiding this comment.
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.
| Seq(hint.leftHint, hint.rightHint).flatten.flatMap(_.strategy).exists { | ||
| case BROADCAST | SHUFFLE_MERGE | SHUFFLE_HASH | SHUFFLE_REPLICATE_NL => true | ||
| case _ => false |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
Could
hasStrategyHintuse 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.
| 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") { |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
could the config block become one class-level helper that takes
adaptiveand 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).
| test(s"forceShuffledHashJoin keeps SORT BY on the join key, AQE=$adaptive") { | ||
| withSortLossConf(adaptive) { | ||
| checkPartitionsSortedOverHashJoin(sql(s"$sortLossJoin SORT BY k")) |
There was a problem hiding this comment.
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?
There was a problem hiding this comment.
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.
|
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.
Which issue does this PR close?
Closes #6770.
Rationale for this change
With
spark.comet.exec.forceShuffledHashJoinon,RewriteJointurned sort-merge joins into hash joins insideCometExecRule, afterEnsureRequirementsandRemoveRedundantSortshad 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.sortWithinPartitionsandSORT BYon the join key returned unsorted partitions, and a partitioned insert keyed on the join key failed withFileAlreadyExistsException, 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?
RewriteJoinis replaced byCometShuffledHashJoinStrategy, injected withinjectPlannerStrategy. Injected strategies run before Spark'sJoinSelection, so the join is aShuffledHashJoinExecfrom the start andEnsureRequirementsandRemoveRedundantSortsplace 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.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'sLogicalQueryStageStrategy, 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;BROADCAST,MERGE,SHUFFLE_HASH,SHUFFLE_REPLICATE_NL), including aSHUFFLE_HASHthat AQE sets on a re-plan, in which case Spark plans its own shuffled hash join. AQE'sNO_BROADCAST_HASHandPREFER_SHUFFLE_HASHdo not count, so a side AQE only stops from broadcasting stays a hash join;hashJoinSupported), or the join is aLeftSingle, which Spark never plans as a sort-merge join;maxBuildSize, or it isLeftSemiorExistenceJoinbuilding 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.restoreRequiredOrderingis gone, sinceEnsureRequirementsnow adds those sorts.spark.comet.exec.enabledis off, in plan-only mode, or for a streaming join.Behavior changes:
forceShuffledHashJoin. Before, aMERGEhint was converted too.understanding-comet-plans.mdsays so.ShuffledHashJoin(Comet leaves streaming plans alone). The tuning guide says so.branch-1.1has the same rewrite, so 1.1.0 loses these sorts too. A backport needs this strategy rather than a cherry-pick, sincebranch-1.1has no size limit.How are these changes tested?
CometJoinSuite, with AQE off and on where it matters:sortWithinPartitionson 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 BYplans the same partition-local sort.MERGE,SHUFFLE_HASH(Spark builds the hinted side) andBROADCASThints are respected.NO_BROADCAST_HASHon 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.Invalid broadcast query stagere-plan failure.OptimizeSkewedJoin.UTF8_LCASEcollation keeps the sort-merge join and returns Spark's rows.spark.comet.exec.enabledoff, or in plan-only mode.MERGEhint they relied on, and an RDD-backed build side, whose size Spark reports asspark.sql.defaultSizeInBytes, keeps the sort-merge join.LeftSemi,ExistenceJoinor window above a forced hash join pass with the sortsEnsureRequirementsadds.The new sort and insert tests fail on main. Removing the plan-only check fails its test, replacing the Spark 4
hashJoinSupportedcheck withtruefails the collation test, dropping the broadcast stage check fails the planned broadcast test, and counting every join hint as a strategy hint fails theNO_BROADCAST_HASHtest.CometJoinSuiteandCometConfSuitepass 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.