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]
