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]