dwsmith1983 opened a new pull request, #6785:
URL: https://github.com/apache/datafusion-comet/pull/6785

   ## 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;
     - either side has a join strategy hint (`BROADCAST`, `MERGE`, 
`SHUFFLE_HASH`, `SHUFFLE_REPLICATE_NL`); AQE's own hints do not count;
     - 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 (#2667, #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 #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` and `SORT BY` on the join key keep every partition 
sorted, for inner, right outer (building on the smaller left side) and full 
outer joins, and the rows match Spark.
   - 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 runtime broadcast under AQE, and a broadcast Spark plans from the static 
size, stay broadcasts.
   - 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 #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 #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, and replacing the Spark 4 `hashJoinSupported` check with `true` 
fails the collation test. `CometJoinSuite` and `CometConfSuite` pass on Spark 
3.4, 3.5, 4.0 and 4.1, Spark 4.2 compiles, and the strict warnings build passes.
   


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

To unsubscribe, e-mail: [email protected]

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


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

Reply via email to