Repository navigation
perf: keep the sort-merge join when the forced hash join's build side is too large - #6384
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.
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: Check the selected build child’s statistics before rewriting. Preserve the sort-merge join and its sorts when the estimate exceeds the limit or statistics are unavailable.
- Correctness / compatibility analysis: The strict comparison and initial partition-count lookup match the relevant Spark sources across 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. The 10 MiB fallback when broadcasts are disabled is an explicit, documented Comet choice. No introduced P1/P2 issues found within this review.
- Key design decisions:
BigIntavoids threshold multiplication overflow. Informational tags explain retained joins without forcing Spark fallback. The helpers keep the change localized, with additional work confined to planning. - Implementation sketch:
CometExecRulepassesSQLConfintoRewriteJoin. The PR adds the size configuration, regression tests, documentation and benchmark overrides. - Behavioral changes worth calling out: Compared with
branch-1.1at992c806a7e38c2e88bd018aa5774164b0850e1fa, forced conversion intentionally becomes size-limited.maxBuildSize<=0restores unrestricted conversion for eligible joins. The force option remains disabled by default. This heuristic does not guarantee that a hash table fits available memory. The separateExistenceJoinrestriction already exists in the base. - Suggested improvements: No additional P1/P2 changes requested.
Reviewed all 10 changed files in the full diff from fef94f6cd78b18151dff57b7a936798385356de5 to 432cf15cdfc25f17f45a989fea2fc708c1be867e. The PR is not a draft. Read repository guidance, the snapshot and refreshed GitHub 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 CI build/test verdict.
Validation: Compiled the current rule and configuration and passed 19 focused planner cases on each of Spark 3.5.9 and 4.1.3. A Spark 3.5.9 AQE execution probe returned 1,000 rows and retained sort merge after materialized statistics increased the build estimate from 7,803 to 280,000 bytes against a 7,804-byte limit. Shell/TOML syntax and diff whitespace checks passed.
Validation limits: The probes use utility/tag adapters and Spark execution, not Comet native execution. Full Comet/native and Spark SQL regression suites, workload benchmarks and memory-pressure tests were not run.
andygrove
left a comment
There was a problem hiding this comment.
When the guard keeps an outer SortMergeJoin, transformUp may already have rewritten an inner one under it, and if both joins are on the same key that loses rows. EnsureRequirements ran before this rule and saw that the inner join's output ordering already satisfied the outer one, so it put no sort between them. Once the inner join becomes a ShuffledHashJoinExec with its sorts removed, the kept outer join reads unsorted input. CometSortMergeJoinExec.convert doesn't check child ordering, so the native join silently misses matches. On branch-1.1 both joins would have been rewritten, so this is new.
I reproduced it on this branch in CometJoinSuite with Parquet tables big as (i % 100, i) for 10000 rows, small as (i * 10, i) for 10 rows and mid as (i % 100, i) for 3000 rows, running SELECT big._1, big._2, small._2, mid._2 FROM big JOIN small ON big._1 = small._1 JOIN mid ON big._1 = mid._1. With forceShuffledHashJoin=true, AQE off, two shuffle partitions and autoBroadcastJoinThreshold=3000, the default limit of 6000 bytes falls between the two build sides. Comet returns 6240 rows and Spark returns 30000. The plan is CometSortMergeJoin over CometProject over CometHashJoin with no CometSort on that side. It happens with AQE on too once the runtime sizes fall on either side of the limit. With mid as (i % 30000, i) for 300000 rows and maxBuildSize=1000000 I got 2080 rows against 10000.
Could the rule put back the local sorts the rewrite took away? I tried a pass after the transformUp in CometExecRule that wraps any child whose outputOrdering doesn't satisfy its parent's requiredChildOrdering in SortExec(required, global = false, child), which is the same check EnsureRequirements makes. It fixed both AQE modes and CometJoinSuite still passed. Doing it after the rewrite rather than only where the size guard keeps the join would also cover the LeftSemi and ExistenceJoin branches that keep a join. Could we also add the chained join as a test with checkSparkAnswer in both AQE modes? The AQE one needs spark.sql.adaptive.autoBroadcastJoinThreshold off as well, since CometTestBase sets it to 1g and AQE otherwise broadcasts both joins.
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 because the hash join cannot spill.
- Design approach: Check the selected build child’s statistics before rewriting. Retain sort merge when statistics are unavailable or exceed the configured limit.
- Correctness / compatibility analysis: The strict size comparison and initial partition-count lookup match Spark sources across 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. However, the existing P1 row-loss report remains substantiated and unresolved. Rewriting an inner join removes ordering required by a retained outer sort-merge join. No additional introduced P1/P2 issues found within this review.
- Key design decisions:
BigIntavoids threshold multiplication overflow. The documented 10 MiB fallback handles disabled broadcasts. Informational tags explain retained joins without forcing Spark fallback. The helpers keep the change localized, with additional work confined to planning. The guard is a heuristic, not a memory guarantee. - Implementation sketch:
CometExecRulepassesSQLConftoRewriteJoin. The PR adds the optional size limit, single-join regression tests, documentation and benchmark overrides. - Behavioral changes worth calling out: Compared with
branch-1.1at992c806a7e38c2e88bd018aa5774164b0850e1fa, forced conversion intentionally becomes size-limited.maxBuildSize<=0restores unrestricted conversion. The force option remains disabled by default. Mixed hash/sort-merge plans introduce the unintended row-loss regression described above. Retaining sorts changes the performance tradeoff, but no workload benchmark was run. - Suggested improvements: Address the existing review by restoring required child ordering after rewriting and covering chained joins with AQE disabled and enabled. Keep changes requested until that blocker is resolved. No duplicate inline finding is emitted.
Reviewed all 10 changed files in the full base-relative diff from 965c8bbe289ff850b614c8c3833ed7802134115b to 9f7a615d3592b30a37a8891f4f2dc52205b4f22f. The PR is not a draft. Read repository guidance, snapshot discussions and refreshed GitHub discussions. Routed skills: review-comet-pr and review-comet-memory-pr.
Exact-head CI: Comet CI has 23 successful checks, 14 skipped checks and the expressions job still running. No failures were reported. Native build, Rust tests, Spark 4.1 execution/shuffle/scan suites, TPC-H and TPC-DS checks passed. Spark SQL suites were skipped.
Validation: Compiled the current rule and configuration and passed 19 focused planner cases on each of Spark 3.5.9 and 4.1.3. An independent Spark 3.5.9 chained-join probe returned 300 rows instead of 30,000 with build estimates of 1,384 and 14,733 bytes against an 8,058-byte limit. Removing the limit or restoring required sorts returned all 30,000 rows. Shell/TOML syntax and diff whitespace checks passed.
Validation limits: Local probes use utility/tag adapters and Spark execution, not Comet native execution. Full Comet/native and Spark SQL suites, AQE chained-join execution and memory-pressure tests were not run locally. Project code remains unchanged.
…rewrite When RewriteJoin keeps a sort-merge join above a join it has already turned into a hash join on the same key, the kept join reads unsorted input. EnsureRequirements placed no sort between the two because the inner sort-merge join's ordering satisfied the outer one, and the rewrite then removed the inner join's sorts. The native sort-merge join does not check its input ordering, so matches were silently dropped. After the rewrite, add a local SortExec wherever a child no longer satisfies its parent's required ordering, the same check EnsureRequirements makes. This covers joins kept by the build-size limit and the LeftSemi and ExistenceJoin joins the rewrite already keeps, which lost rows the same way.
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Forced shuffled hash joins ignored build-side size despite lacking spill support. Selectively retaining sort-merge joins also exposed missing ordering above rewritten inner joins.
- Design approach: Check the chosen build side against a configurable size limit, then restore any required child ordering before Comet conversion.
- Correctness / compatibility analysis: The strict size comparison, initial partition-count lookup and ordering checks align with relevant Spark sources across 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. The earlier chained-join row-loss concern is addressed. No introduced P1/P2 issues found within this review.
- Key design decisions:
BigIntavoids multiplication overflow. The documented 10 MiB fallback handles disabled broadcasts. Informational tags explain retained joins without forcing Spark fallback. The helpers remain localized, with an additional planning traversal and sorts where ordering requirements are unmet. - Implementation sketch:
CometExecRulesuppliesSQLConf, performs the rewrite and restores ordering. Configuration, regression tests, tuning documentation and benchmark overrides accompany the change. - Behavioral changes worth calling out: Compared with
branch-1.1at992c806a7e38c2e88bd018aa5774164b0850e1fa, forced conversion intentionally becomes size-limited.maxBuildSize<=0restores unrestricted conversion for eligible joins. Ordering repair also fixes the existing chained semi-join problem. The force option remains disabled by default, and the size guard remains a heuristic rather than a memory guarantee. - Suggested improvements: No additional P1/P2 changes requested. The previous ordering blocker is resolved by this head.
Reviewed all 10 changed files in the full diff from 965c8bbe289ff850b614c8c3833ed7802134115b to 6f7b639f1e2ae3d8d853900002f041e33495e819. The PR is not a draft. Read repository guidance, snapshot discussions and refreshed GitHub discussions. Routed skills: review-comet-pr and review-comet-memory-pr.
Exact-head CI: The label check passed. Comet CI, CodeQL and the title workflow remain action_required, awaiting approval. There is no exact-head CI build/test verdict.
Validation: Compiled the exact-head rule and configuration and passed 19 focused planner cases on each of Spark 3.5.9 and 4.1.3. Six Spark 3.5.9 execution probes covered chained inner, semi and existence joins with AQE off and on. Inner joins returned all 30,000 expected rows, while semi and existence joins returned all 1,000. Controls omitting ordering repair returned only 300 and 10 rows respectively. Ordering repair was also idempotent. Shell/TOML syntax and diff whitespace checks passed.
Validation limits: Local probes use utility/tag adapters and Spark execution, not Comet native execution. Full Comet/native and Spark SQL suites, workload benchmarks and memory-pressure tests were not run. Project code remains unchanged.
Done in 6f7b639, as you suggested: after the
Added, with your tables and the |
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Forced hash-join conversion ignored build-side size despite lacking spill support. Selectively retaining sort-merge joins also exposed missing ordering above rewritten inner joins.
- Design approach: Check the selected build side against a configurable limit, then restore required child ordering before Comet conversion.
- Correctness / compatibility analysis: The strict comparison, initial partition-count lookup and ordering predicate align with relevant Spark sources across 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. The previously reported chained-join row loss is addressed. No introduced P1/P2 issues found within this review.
- Key design decisions:
BigIntavoids threshold multiplication overflow. Disabled broadcasts use the documented 10 MiB fallback. Informational tags explain retained joins. The localized helpers add a planning traversal and sorts where ordering requirements are unmet, without introducing a broader abstraction. - Implementation sketch:
CometExecRulesuppliesSQLConf, applies the rewrite and repairs ordering. Tests, configuration documentation and benchmark overrides accompany the change. - Behavioral changes worth calling out: Compared with
branch-1.1at992c806a7e38c2e88bd018aa5774164b0850e1fa, forced conversion intentionally becomes size-limited.maxBuildSize<=0restores unrestricted conversion. Ordering repair also fixes the existing chained semi-join problem. The force option remains disabled by default. The size guard is a heuristic, not a memory guarantee. - Suggested improvements: No additional P1/P2 changes requested. The existing ordering blocker is resolved at this head.
Reviewed all 10 changed files in the full diff from 965c8bbe289ff850b614c8c3833ed7802134115b to 6f7b639f1e2ae3d8d853900002f041e33495e819. Confirmed the PR is not a draft and read existing reviews, comments and threads. Routed skills: review-comet-pr and review-comet-memory-pr.
Exact-head CI: The label check passed. Comet CI, CodeQL and the title workflow remain action_required, awaiting approval. No exact-head CI build/test verdict is available.
Validation: Reused prior exact-head evidence after verifying source hashes, harnesses, invocations and logs. Nineteen planner cases passed on each of Spark 3.5.9 and 4.1.3. Six Spark 3.5.9 execution probes covered chained inner, semi and existence joins with AQE off and on. Repaired plans matched Spark's 30,000 inner-join rows and 1,000 semi/existence rows. Controls without repair returned only 300 and 10 rows respectively. Ordering repair was idempotent. Fresh shell/TOML syntax and diff whitespace checks passed.
Validation limits: Execution probes use utility/tag adapters and Spark execution, not Comet native execution. Full Comet/native and Spark SQL suites, workload benchmarks and memory-pressure tests were not run. Project code remains unchanged.
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Forced hash-join conversion ignored build-side size despite lacking spill support. Selectively retaining sort-merge joins also exposed missing ordering above rewritten joins.
- Design approach: Check the selected build side against a configurable limit, then restore required child ordering before Comet conversion.
- Correctness / compatibility analysis: The strict comparison, initial partition-count lookup and ordering predicate align with relevant Spark sources across 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. The previously reported chained-join row loss is addressed. No introduced P1/P2 issues found within this review.
- Key design decisions:
BigIntavoids threshold multiplication overflow. Disabled broadcasts use the documented 10 MiB fallback. Informational tags explain retained joins. Localized helpers add one planning traversal and sorts where ordering requirements are unmet, without introducing a broader abstraction. - Implementation sketch:
CometExecRulesuppliesSQLConf, applies the rewrite and repairs ordering. Configuration, regression tests, documentation and benchmark overrides accompany the change. - Behavioral changes worth calling out: Compared with
branch-1.1at7b7eec69282a7abc33a0e8e687ca594c72f140ec, forced conversion intentionally becomes size-limited.maxBuildSize<=0restores unrestricted conversion. Ordering repair also fixes the existing chained semi-join problem. The force option remains disabled by default. The limit is an estimate-based heuristic, not a memory guarantee. - Suggested improvements: No additional P1/P2 changes requested. The existing ordering blocker is resolved at this head.
Reviewed all 10 changed files in the full diff from 4ea367aab4af430fce5dafa84351bedd288ec074 to c469de3f6290e145312c469d4bce7ead66091fd7. Confirmed the PR is not a draft and read the snapshot plus refreshed reviews, comments and threads. 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. No exact-head CI build/test verdict is available.
Validation: Freshly compiled the current rule and configuration and passed 19 planner cases on each of Spark 3.5.9 and 4.1.3. Six Spark 3.5.9 execution probes covered chained inner, semi and existence joins with AQE off and on. Repaired plans matched Spark's 30,000 inner-join rows and 1,000 semi/existence rows. Controls without repair returned only 300 and 10 rows respectively. Ordering repair was idempotent. Shell/TOML syntax and diff whitespace checks passed.
Validation limits: Probes use utility/tag adapters and Spark execution, not Comet native execution. Full Comet/native builds and Spark SQL suites, workload benchmarks and memory-pressure tests were not run. Project code remains unchanged.
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Forced shuffled hash joins ignored build-side size despite lacking spill support. Selectively retaining sort-merge joins also exposed missing ordering above rewritten joins.
- Design approach: Check the selected build side against a configurable limit, then restore required child ordering before Comet conversion.
- Correctness / compatibility analysis: Verified the strict size comparison, initial partition-count lookup and ordering predicate against Spark sources for 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. The previously reported chained-join row loss is addressed. No introduced P1/P2 issues found within this review.
- Key design decisions:
BigIntavoids multiplication overflow. Disabled broadcasts use the documented 10 MiB fallback. Informational tags explain retained joins. Localized helpers add one planning traversal and sorts where ordering requirements are unmet, without introducing a broader abstraction. - Implementation sketch:
CometExecRulesuppliesSQLConf, applies the rewrite and repairs ordering. Configuration, regression tests, documentation and benchmark overrides accompany the change. - Behavioral changes worth calling out: Compared with
branch-1.1at7b7eec69282a7abc33a0e8e687ca594c72f140ec, forced conversion intentionally becomes size-limited.maxBuildSize<=0restores unrestricted conversion for eligible joins. Ordering repair also fixes the existing chained semi-join problem. The force option remains disabled by default. The guard uses estimates and does not guarantee that a hash table fits memory. - Suggested improvements: No additional P1/P2 changes requested. The existing ordering concern is resolved at this head.
Reviewed all 10 changed files in the full base-relative diff from 9dc8c3ca89962506078831f7064a25c75da83883 to 6b9fef7e4f1b4b574b79324a83dc9fd91f9a8dc0. Confirmed the PR is not a draft and read the snapshot plus refreshed reviews, comments and threads. 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. No exact-head CI build/test verdict is available.
Validation: Freshly compiled the current rule and configuration. All 19 planner cases passed on each of Spark 3.5.9 and 4.1.3. Six Spark 3.5.9 execution probes covered chained inner, semi and existence joins with AQE off and on. Repaired plans matched Spark’s 30,000 inner-join rows and 1,000 semi/existence rows. Controls without repair returned only 300 and 10 rows respectively. Ordering repair was idempotent. Shell/TOML syntax and diff whitespace checks passed.
Validation limits: Probes use utility/tag adapters and Spark execution, not Comet native execution. Full Comet/native builds, Comet and Spark SQL regression suites, workload benchmarks and memory-pressure tests were not run locally. Project code remains unchanged.
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Forced shuffled hash joins ignored build-side size despite lacking spill support. Selectively retaining sort-merge joins also exposed missing ordering above rewritten joins.
- Design approach: Check the selected build side against a configurable limit, then restore required child ordering before Comet conversion.
- Correctness / compatibility analysis: Verified the strict size comparison, initial partition-count lookup and ordering predicate against Spark sources for 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. The previously reported chained-join row loss is addressed. No introduced P1/P2 issues found within this review.
- Key design decisions:
BigIntavoids multiplication overflow. Disabled broadcasts use the documented 10 MiB fallback. Informational tags explain retained joins. Localized helpers add one planning traversal and sorts where ordering requirements are unmet, without introducing a broader abstraction. - Implementation sketch:
CometExecRulesuppliesSQLConf, applies the rewrite and repairs ordering. Configuration, regression tests, documentation and benchmark overrides accompany the change. - Behavioral changes worth calling out: Compared with
branch-1.1ate9efd9f764ee0a59b7898ff028d6985d4a7a28e1, forced conversion intentionally becomes size-limited.maxBuildSize<=0restores unrestricted conversion for eligible joins. Ordering repair also fixes existing chained semi/existence-join cases. The force option remains disabled by default. The guard uses estimates and does not guarantee that a hash table fits memory. - Suggested improvements: No additional P1/P2 changes requested. The existing ordering concern is resolved at this head.
Reviewed all 10 changed files in the full base-relative diff from 8c783aa88104616dcf0f7876b4f7a9ba71e2bf31 to 1d35245660e779a909bc0b25c30c32639667d33f. Confirmed the PR is not a draft and read the snapshot plus refreshed reviews, comments and threads. 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. No exact-head CI build/test verdict is available.
Validation: Freshly compiled the current rule and configuration. All 19 planner cases passed on each of Spark 3.5.9 and 4.1.3. Six Spark 3.5.9 execution probes covered chained inner, semi and existence joins with AQE off and on. Repaired plans matched Spark’s 30,000 inner-join rows and 1,000 semi/existence rows. Controls without repair returned only 300 and 10 rows respectively. Ordering repair was idempotent. Shell/TOML syntax and diff whitespace checks passed.
Validation limits: Probes use utility/tag adapters and Spark execution, not Comet native execution. Full Comet/native builds, Comet and Spark SQL regression suites, workload benchmarks and memory-pressure tests were not run locally. Project code remains unchanged.
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Forced shuffled hash joins ignored build-side size despite lacking spill support. Selectively retaining sort-merge joins also exposed lost ordering above rewritten joins.
- Design approach: Check the selected build side against a configurable size limit, then restore required child ordering before Comet conversion.
- Correctness: The strict size comparison and ordering predicate match the relevant Spark sources. Boundary, missing-statistics and build-side cases passed focused checks. Ordering repair addresses the previously reported chained-join row loss.
- Compatibility analysis: Checked Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0 sources, including initial shuffle partition selection. Substituting the default 10 MiB threshold when broadcasts are disabled is a deliberate, documented Comet choice.
- Key design decisions:
BigIntavoids threshold multiplication overflow. Informational tags explain retained joins without forcing Spark fallback. Non-positivemaxBuildSizepreserves unrestricted conversion for eligible joins. - Implementation sketch:
CometExecRulepassesSQLConfintoRewriteJoin, applies the rewrite and repairs ordering. The PR adds configuration, regression tests, tuning documentation and benchmark overrides. - Performance: Additional work consists of statistics checks and a planning traversal when forced conversion is enabled. Sorts are inserted only where ordering requirements are unmet. Retaining spillable sort-merge joins avoids forcing known-large build sides into hash joins, but workload speed and memory improvements were not measured.
- Design: The guard and ordering repair are localized and straightforward. Repairing ordering after the entire rewrite covers retained inner, semi and existence joins without separate fixes for each parent shape.
- Abstraction & complexity: Small helpers separate limit calculation, refusal reasons and ordering restoration. They reuse Spark’s ordering predicate and introduce no broader framework. No complexity issue meeting the P1/P2 bar was identified.
- Behavioral changes worth calling out: Compared with
branch-1.1ate9efd9f764ee0a59b7898ff028d6985d4a7a28e1, forced conversion intentionally becomes size-limited. The force option remains disabled by default. Ordering repair also fixes existing chained semi/existence-join cases. The size guard is an estimate-based heuristic, not a guarantee that a hash table fits memory. - Suggested improvements: No additional P1/P2 changes identified. The existing ordering concern is resolved at this head.
Reviewed all 10 changed files in the full base-relative diff from b56349697b786ff2ad1c1bcf6ecf45b809af5f30 to 9e549136ca773765cc172cf824778359b53db577. Confirmed the PR is not a draft and read existing reviews, conversation comments and the empty inline/thread sets. 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. No exact-head CI build/test verdict is available.
Validation: Freshly compiled the current rule and configuration. All 19 focused planner cases passed on each of Spark 3.5.9 and 4.1.3. Reused six Spark 3.5.9 execution probes after verifying source hashes, harnesses, invocations, logs and bytecode against the fresh compilation. With AQE off and on, repaired plans matched Spark’s 30,000 inner-join rows and 1,000 semi/existence rows. Controls without repair returned only 300 and 10 rows respectively. Ordering repair was idempotent. Shell/TOML syntax and diff whitespace checks passed.
Validation limits: The probes use utility/tag adapters and Spark execution, not Comet native execution. Full Comet/native builds, Comet and Spark SQL regression suites, workload benchmarks and memory-pressure tests were not run locally. Project code remains unchanged.
No introduced P1/P2 issues found within this review.
|
@andygrove could you take another look at the sort-order fix (6f7b639, described above)? The branch is current with |
andygrove
left a comment
There was a problem hiding this comment.
I checked out 9e54913 and ran CometJoinSuite on Spark 4.1 and the new tests on 4.2, and they pass. With restoreRequiredOrdering disabled, all six chained tests fail with the row counts from before, so they catch the bug.
The pass also fixes a crash that no test covers yet. A window partitioned by the join key over a rewritten join, such as SELECT big._1, big._2, small._2, count(*) OVER (PARTITION BY big._1) FROM big JOIN small ON big._1 = small._1 with maxBuildSize=-1 and broadcasts off, fails on main with the native assertion All partition by columns should have an ordering. That happens because EnsureRequirements skipped the window's sort in favour of the join's ordering. On this branch it matches Spark with AQE on and off. Could we add it as a test next to the chained ones, with a check that the join really became a CometHashJoinExec? So far every test has a join as the parent, and this one would show the pass covers other operators too.
This also fixes #6673 and adds its tests. Could the description say Closes #6673 so the issue closes on merge?
A window partitioned by the join key gets no sort of its own, because the sort-merge join's output ordering already satisfies it. Once RewriteJoin turns that join into a hash join, the window needs the ordering restored, or the native window fails its partition-ordering assertion. The test runs with AQE on and off and checks the join became a hash join.
Added in 60730d8, after the kept-join tests: your window query with |
|
The merge queue run failed on a flaky OOM, not on this change. |
Which issue does this PR close?
Closes #6673.
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.SortExecis added wherever a child no longer satisfies its parent's required ordering, the same checkEnsureRequirementsmakes. A kept sort-merge join can sit directly above a join on the same key that the rule rewrote to a hash join, with no sort between them because the inner join's ordering satisfied the outer one before its sorts were removed. The rewrite already keepsLeftSemiandExistenceJoinjoins this way, so this also fixes Kept sort-merge join reads unsorted input after replaceSortMergeJoin rewrites the join below it #6673, where those joins return wrong rows on main and, forLeftSemi, onbranch-1.1. That part is a backport candidate.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.ExistenceJoinandLeftSemiwith no limit, in both AQE modes. Each compares rows with Spark and checks that a sort sits between the kept join and the hash join. Without the restored sort all six return too few rows.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.