perf: keep the sort-merge join when the forced hash join's build side is too large - #6384
dwsmith1983 wants to merge 7 commits into
Conversation
… is too large With spark.comet.exec.forceShuffledHashJoin on, RewriteJoin turned every eligible SortMergeJoinExec into a shuffled hash join and read the build side's size estimate only to pick the build side. DataFusion's hash join has to hold that side in memory, so one large join failed the task while the rest of the query would have gained from the rewrite. The rule now rewrites only when the build side's estimate is under a limit and otherwise keeps the SortMergeJoinExec, which still runs natively, recording the reason as plan info. The limit is Spark's own rule for choosing a shuffled hash join, spark.sql.autoBroadcastJoinThreshold times the initial shuffle partition count, with Spark's 10 MB default threshold standing in when broadcasts are disabled. spark.comet.exec.forceShuffledHashJoin.maxBuildSize sets a fixed limit, and a non-positive value removes it, which is what the hash join benchmark configurations and the macOS benchmarking guide now pin so they keep measuring every join.
|
@andygrove could you approve a CI run on |
|
@andygrove thanks for running CI here. Main picked up a conflict in |
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Forced shuffled hash joins ignored build-side size, exposing large joins to failures in the non-spilling hash join.
- Design approach: Gate the rewrite using Spark statistics and preserve the sort-merge join and its sorts when the guard refuses conversion.
- Correctness / compatibility analysis: The strict size comparison and initial partition-count selection match Spark sources across 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. Boundary, missing-statistics, build-side and configuration checks passed. No introduced P1/P2 issues found within this review.
- Key design decisions:
BigIntavoids threshold multiplication overflow. Disabled broadcasts use the documented 10 MiB fallback. Refusal uses informational tags without forcing Spark fallback. The helpers keep the change localized, and added work occurs during planning. - Implementation sketch:
CometExecRulesuppliesSQLConf, andRewriteJoinchecks the chosen build child before removing sorts. Tests, tuning documentation and benchmark configurations accompany the change. - Behavioral changes worth calling out: The forced rewrite becomes size-limited by default. A non-positive
maxBuildSizerestores unrestricted conversion. This remains an estimate-based heuristic, not a guarantee that the hash table fits memory. - Suggested improvements: No additional P1/P2 code changes identified.
Reviewed all 10 changed files in the full diff from 961fbfe768e034a339c0090d99bfc5402e6e7dad to 1506a55c8e097813064b9a6fd0307c5b13a59c00. The PR is not a draft. Read the snapshot and refreshed discussions. No unresolved substantiated P1/P2 concerns were present. Routed skills: review-comet-pr and review-comet-memory-pr.
Exact-head CI: The label check passed. Comet CI and CodeQL remain action_required, awaiting approval. There is no exact-head build/test verdict.
Validation: Compiled the current rule and configuration and passed 19 isolated planner cases on each of Spark 3.5.9 and 4.1.3. A Spark 4.1.3 AQE execution probe returned 1,000 rows and changed an initially eligible hash join to sort merge when materialized statistics exceeded the limit. Shell/TOML syntax and diff whitespace checks passed.
Validation limits: The probes used lightweight utility/tag adapters and Spark execution, not Comet native execution. Full Comet/native and Spark SQL regression suites were not run. No workload performance benchmark was run.
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. Withspark.comet.exec.forceShuffledHashJoinon,RewriteJointurns 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 ownJoinSelectiongates shuffled hash join oncanBuildLocalHashMapBySize, and it does not choose it for a side it cannot size.What changes are included in this PR?
RewriteJoinrewrites only when the build side's size estimate is under a limit. Otherwise it keeps theSortMergeJoinExec, 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 aJoin, is kept too, with a reason saying no statistics were available. An unknown size isLong.MaxValuein 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.autoBroadcastJoinThresholdtimes the initial shuffle partition count, which is whatcanBuildLocalHashMapBySizecompares 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.spark.comet.shuffle.sizeInBytesMultiplierwhen 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.CometExecRulepasses itsSQLConfto the rule.run_all_benchmarks.shand the TPC-H command in the macOS benchmarking guide pinmaxBuildSize=-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:SortMergeJoinExecbuilt without a logical link) keeps the join with a reason, and a non-positive limit rewrites it.autoBroadcastJoinThreshold=-1, a small build side is rewritten and one estimated over 10 MB is kept with the reason naming that limit.CometJoinSuite,CometConfSuite,CometExecSuiteand the TPC-DS plan stability suites pass on Spark 3.5;CometJoinSuitepasses 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.