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

   ## Which issue does this PR close?
   
   Part of #2545. The spilling itself is being built in DataFusion 
(apache/datafusion#24768, with the sort-merge fallback in 
apache/datafusion#25217 as its first step). Until that lands, this keeps the 
forced hash join away from build sides that are unlikely to fit.
   
   ## Rationale for this change
   
   Comet's hash join is DataFusion's `HashJoinExec`, whose build side has to 
fit in memory. With `spark.comet.exec.forceShuffledHashJoin` on, `RewriteJoin` 
turns every eligible sort-merge join into a shuffled hash join and drops the 
input sorts. It reads the build side's size estimate only to choose the build 
side, never to ask whether that side can fit, so one large join fails the task 
while the rest of the query would have gained from the rewrite. Spark's own 
`JoinSelection` gates shuffled hash join on `canBuildLocalHashMapBySize`, and 
it does not choose it for a side it cannot size.
   
   ## What changes are included in this PR?
   
   - `RewriteJoin` rewrites only when the build side's size estimate is under a 
limit. Otherwise it keeps the `SortMergeJoinExec`, which still runs natively, 
and records why as plan info (not as a fallback, since nothing falls back to 
Spark). A join with no logical link, or one whose link is not a `Join`, is kept 
too, with a reason saying no statistics were available. An unknown size is 
`Long.MaxValue` in Spark, so it is over any limit.
   - `spark.comet.exec.forceShuffledHashJoin.maxBuildSize`, an optional byte 
size. Unset, the limit is Spark's own rule: 
`spark.sql.autoBroadcastJoinThreshold` times the initial shuffle partition 
count, which is what `canBuildLocalHashMapBySize` compares against. When the 
threshold is not positive, which is how broadcasts get disabled, the limit uses 
Spark's default threshold of 10 MB instead: turning broadcasts off says nothing 
about the hash table an executor can hold, and a non-positive limit would 
silently turn the rewrite off. A positive value is a fixed limit; a 
non-positive value means no limit, today's behavior. Like Spark, the comparison 
is strict: a build side at exactly the limit is kept.
   - The size is Spark's planning estimate for the build child. Under AQE the 
rewrite runs again when a stage is prepared, and the estimate is then the 
materialized shuffle size of that side, scaled by 
`spark.comet.shuffle.sizeInBytesMultiplier` when Comet's shuffle produced it. 
The tests cover the planning estimate; the AQE re-planning path is exercised 
but its stage size is not asserted.
   - `CometExecRule` passes its `SQLConf` to the rule.
   - The tuning guide describes the limit next to the force config, and the 
memory management page's note about the hash join mentions it. The two hash 
join benchmark configurations, `run_all_benchmarks.sh` and the TPC-H command in 
the macOS benchmarking guide pin `maxBuildSize=-1`, so they keep measuring the 
hash join on every join as before.
   
   ## How are these changes tested?
   
   Tests in `CometJoinSuite`, with the force config on and Parquet tables whose 
sizes come from their files:
   
   - a build side over an explicit limit keeps the sort-merge join and its 
sorts, with the reason naming the estimate and the limit and no fallback 
recorded, with AQE off and on; on main it is rewritten.
   - a build side under the limit becomes a hash join with the expected build 
side and no sorts, with AQE off and on.
   - no statistics (a `SortMergeJoinExec` built without a logical link) keeps 
the join with a reason, and a non-positive limit rewrites it.
   - a non-positive limit rewrites regardless of the threshold.
   - with the limit unset, a small threshold and partition count keep the join, 
and a large threshold rewrites it.
   - with the limit unset and `autoBroadcastJoinThreshold=-1`, a small build 
side is rewritten and one estimated over 10 MB is kept with the reason naming 
that limit.
   - at the boundary, a limit equal to the estimate keeps the join and one byte 
more rewrites it.
   
   `CometJoinSuite`, `CometConfSuite`, `CometExecSuite` and the TPC-DS plan 
stability suites pass on Spark 3.5; `CometJoinSuite` passes on Spark 3.4, 4.0 
and 4.1. Disabling the comparison, the default-threshold fallback, or the info 
tag each fails the tests written for it.
   


-- 
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