Skip to content

perf: keep the sort-merge join when the forced hash join's build side is too large - #6384

Open
dwsmith1983 wants to merge 7 commits into
apache:mainfrom
dwsmith1983:perf/2545-hash-join-build-size-guard
Open

dwsmith1983 wants to merge 7 commits into
apache:mainfrom
dwsmith1983:perf/2545-hash-join-build-size-guard

Conversation

@dwsmith1983

Copy link
Copy Markdown
Contributor

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.

… 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.
@github-actions github-actions Bot added enhancement New feature or request performance area:Iceberg area:joins Join operators and dynamic filter pushdown labels Sep 29, 2026
@dwsmith1983

Copy link
Copy Markdown
Contributor Author

@andygrove could you approve a CI run on cc9a9a167 and take a look when you have time? This is the build size guard for #2545.

@dwsmith1983

Copy link
Copy Markdown
Contributor Author

@andygrove thanks for running CI here. Main picked up a conflict in CometConf.scala from #6474 (both added a config after forceShuffledHashJoin), so the head is now 1506a55c8 with only that resolved. Could you approve CI on it again?

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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: BigInt avoids 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: CometExecRule supplies SQLConf, and RewriteJoin checks 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 maxBuildSize restores 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.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:Iceberg area:joins Join operators and dynamic filter pushdown enhancement New feature or request performance

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants