Skip to content

fix: preserve current AQE logical-stage links on Comet operators - #5483

Merged
sunchao merged 5 commits into
apache:mainfrom
sunchao:fix/aqe-logical-stage-links-20260826
Sep 25, 2026
Merged

sunchao merged 5 commits into
apache:mainfrom
sunchao:fix/aqe-logical-stage-links-20260826

Conversation

@sunchao

@sunchao sunchao commented Aug 26, 2026 •

Copy link
Copy Markdown
Member

Which issue does this PR close?

Closes #5482.

Rationale for this change

Spark's Adaptive Query Execution (AQE) can use runtime statistics to replace a join with a broadcast hash join when an input is small enough. Comet can disrupt this adaptation for native aggregate inputs, leaving queries with correct results but without the expected native broadcast joins.

For example, the added regression joins several aggregates built from generated data. One of its join inputs is:

SELECT id % 64 AS k, SUM(id + 1) AS v
FROM range(0, 3072, 1, 4)
GROUP BY id % 64

This reduces 3,072 input rows to just 64 groups. The full regression uses two 64-row grouped inputs on the right sides of successive left joins. It disables automatic broadcast hash join selection during initial planning but allows AQE to select those joins later, with a 10 MiB adaptive broadcast threshold. The native broadcast hash joins must therefore come from adaptive replanning.

The local comparison shows why both the result and the final plan matter:

Full regression query Without the fix With the fix
Native broadcast hash joins in the final plan 0 2
Query result 48738816 48738816

AQE relies on links between physical operators and logical plan nodes to update a query's plan as stages finish. When Spark represents a native final aggregate above a shuffle stage as a LogicalQueryStage, it reuses the physical aggregate and links it directly to that current logical-stage object. Comet's repair step then overwrites this fresh link with the older originalPlan.logicalLink, or clears it if the original has no link.

The older logical aggregate is now hidden inside the logical-stage leaf, no longer a node in the active logical tree. If a later exchange inherits that stale link, Spark cannot insert its query stage into the active tree: replacement looks for the exact referenced object, which is no longer there. AQE can then lose the new stage during subsequent replanning.

What changes are included in this PR?

This PR lets Spark's current AQE stage assignment take precedence over Comet's saved original link. When CometExecRule revisits a native operator that already has a direct LogicalQueryStage link, it leaves that link intact. The reused operator can therefore stay associated with the stage Spark is currently planning, even across repeated replanning and even when its original plan has an older or missing link.

The exception is limited to direct stage links. An inherited link can come from an ancestor and does not establish that the operator represents that stage, so ordinary and inherited links still follow the existing repair and clearing rules. Shuffle and broadcast exchange handling also stays unchanged, including the empty-link behavior required by #323. The fix changes planning metadata; it requires no changes to Spark, native aggregate execution, or configuration. The broader audit of logical-link repair is tracked separately in #6034.

The rebase also corrects one test inherited from main's #6095: the positional round-robin retry regression now invalidates exchange.shuffleDependency.shuffleId. The Spark 3.4 exchange shim reports shuffleId as zero, which could make the test invalidate a nonexistent or unrelated shuffle after earlier tests had run. Using the dependency's actual registered ID preserves the intended retry assertion across Spark versions.

How are these changes tested?

The planner regressions exercise Spark's actual reuse path through two successive replanning passes. They check that the exact current stage link survives whether the original plan has a direct, inherited, or absent link, and that ordinary and inherited links still receive the existing repair or clearing behavior.

A DPP execution regression observes Spark's preparation rules before and after Comet repairs a reused broadcast root. It verifies that the actual root reaches the direct-stage-link case, the physical plan and standard cost remain unchanged, and the temporary logical identity survives Comet's repair. An isolated copy with conflicting temporary and direct links also exercises Spark's stage-link assignment to verify that the temporary link takes precedence. The SQL result and native broadcast/DPP operators are checked, and the test skips Spark 3.4, where native AQE DPP is unsupported.

The execution regression runs the generated-data query with native shuffle and checks its result and adaptive broadcast joins. It also verifies that both grouped aggregates feeding the broadcasts retain direct stage links, each produce 64 rows, and record positive native compute time.

Local validation of 0b4580c1a, rebased onto a86c9672a, used Linux x86_64, JDK 17, and Spark 4.0.4, with the native library rebuilt using make core in debug mode with debug symbols disabled. Root-reactor Maven passed 99 tests with no failures or cancellations: the full CometExecRuleSuite, the full CometNativePositionalRoundRobinSuite, three targeted execution regressions for native aggregate broadcast adaptation, shuffle logical links, and adaptive shuffled-join conversion, and both table-cache regressions from #6208. These verify that the preserved logical-stage links and the newly landed cache-stage conversion behavior work together. The DPP broadcast-root lifecycle regression also passed. Scalastyle, Spotless, git diff --check, and rebase range-diff review passed.

Before the final import-only conflict resolution against #6208, the preceding head passed 94 local tests on Spark 3.4.3 (three expected version-gated cancellations), 97 on Spark 4.0.4, and all 48 hosted PR checks, including every supported Spark profile and the Spark 4.1 SQL suite. The final rebase preserves those test bodies and merges the adjacent LogicalQueryStage and InMemoryTableScanExec imports.

The Spark 3.4 positional-shuffle suite reproduced the CI failure before the one-line test correction (10 passed, one failed) and passed all 11 tests afterward. Local Spark 4.1 validation remained blocked before test compilation by the configured Maven mirror for jackson-bom:2.21.2; the repository settings were preserved. Fresh hosted CI runs the supported Spark profile matrix and Spark 4.1 SQL tests.

The original controlled before/after comparison on Spark 4.1.3 made both the direct-link planner regression and the SQL broadcast assertion fail when only the production guard was removed. CI is configured to run Spark SQL tests and Comet's supported Spark profiles for additional version coverage.

@andygrove

Copy link
Copy Markdown
Member

Note on this review: this was generated by an LLM (Claude Code) at my request while I worked through a review backlog. I have not verified the individual findings myself. Please treat everything below as suggestions to evaluate rather than as authoritative review feedback, and push back on anything that is wrong or already handled.

The write-up on this one is really good. The explanation of why a stale link makes AQE lose the stage during replanning is the clearest description of that failure mode I have seen, and the unit test in CometExecRuleSuite that walks two replanning rounds across all three original-tag shapes is exactly the right way to pin the behavior down.

Some things I would like to see addressed.

Only one of the three link-repair blocks is guarded

CometExecRule.scala has three near-identical repair blocks in that same transform, one for CometExec, one for CometShuffleExchangeExec, and one for CometBroadcastExchangeExec. The new guard is only on the first one. The rationale in the description talks about a later exchange inheriting the stale link, so it is not obvious to me why the exchange cases are safe to leave alone. Can an exchange ever carry a direct LogicalQueryStage tag that we would then clobber? If the answer is no, a sentence in the comment saying so would save the next reader the same question.

Related to that, those three blocks are copies of each other and now they have diverged. Would you be willing to pull the repair into a single private helper that all three cases call, with the guard living inside it? Right now a future change to the repair logic has to be made in three places and it is easy to miss one, which is roughly what happened here.

Metric assertions in the CometExecSuite test

The new test asserts aggregate.longMetric("elapsed_compute").value > 0. That does not tell us anything about the logical-link fix and it is the kind of timing assertion that eventually flakes on a loaded CI runner. The output_rows == 64 check is useful because it confirms we found the right aggregate. Could the elapsed_compute one just be dropped?

Query shape and runtime

The regression builds four range scans totaling roughly 10k rows and joins them. That is fine for correctness, but CometExecSuite is already one of the slower suites. Did you check what this adds to the suite's wall clock? If it is meaningful, it might be worth shrinking the ranges, since the point is the number of groups on the build side and not the input row count.

Tracking the general problem

Is #5482 the right home for the broader question of whether Comet's link repair is still needed at all in its current form, or is that a separate issue? The repair exists because originalPlan is the source of truth for the link, and this PR establishes that it is not always the source of truth. It would be good to have that written down somewhere other than a code comment.

Comment on lines +653 to +656
// AQE replanning reuses this physical root and links it to the current logical stage.
// originalPlan can still point to a subtree hidden inside that logical leaf, which
// AQE cannot replace in the current logical plan. Only preserve a direct stage link,
// not a link inherited from an ancestor.

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.

is it possible to link some spark code snippet here with version tag to show why this case is added? so readers can understand these comment with more context

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

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

Updated in 92b0baa. Added version-pinned Spark 4.1.3 links beside the guard: LogicalQueryStageStrategy returns the existing physical root, then SparkStrategies.plan assigns its direct logical link. The comment also distinguishes the ordinary exchange path, where the exchange is behind a query-stage leaf, and keeps the direct/inherited distinction explicit.

@sunchao

sunchao commented Aug 27, 2026

Copy link
Copy Markdown
Member Author

Updated in 92b0baac9.

I checked the exchange question against the actual AQE/DPP lifecycle before widening the guard. On the ordinary exchange path, Spark wraps exchanges in query-stage leaves, so this transform does not revisit the exchange inside them.

There is an exception worth distinguishing: the manually inserted DPP broadcast can carry a direct LogicalQueryStage link as the whole adaptive-plan root. I observed this in real SQL with a grouped fact input joining a grouped, filtered dimension, and in a grouped self-join. In both observed replans, repair restored the original aggregate link, but the prepared physical tree was equal to the reused tree, both costs were SimpleCost(0), and TEMP_LOGICAL_PLAN_TAG retained the original aggregate identity. With Spark's standard cost evaluator, AQE rejects that unchanged replan, and new stage creation prioritizes the retained temporary link. All three probe queries matched Spark and retained Comet DPP subqueries. This does not support an additional exchange guard for the lifecycle tested here, so I kept the existing exchange restore/clear behavior and the #323 empty-link contract. I also left helper refactoring out of this focused fix.

I retained elapsed_compute > 0: it checks a recorded native measurement, not a latency target. Comet exposes nanoseconds, and the locked DataFusion timer adds at least one nanosecond per recorded duration. A loaded CI runner does not create the proposed zero-duration failure.

For runtime, the existing Spark 4.1 CI job recorded 471 ms for this regression in an approximately 128-second suite, so I kept the input sizes. #5482 records the concrete stale-stage-link defect; removing all original-plan repair would be separate work and still needs to preserve #323.

The follow-up adds the requested version-pinned Spark references beside the guard. Local validation rebuilt the full Spark 4.1.3 JVM reactor and passed all 31 planner tests plus five execution/AQE/DPP/shuffle-link regressions (36 total). The separate DPP probe also asserts physical-plan equality, equal cost, and preservation of the temporary logical identity. These local runs reused a previously built OSS native library; no native code changed.

@sunchao
sunchao force-pushed the fix/aqe-logical-stage-links-20260826 branch from 92b0baa to f81c995 Compare August 28, 2026 23:28
@andygrove andygrove added the bug Something isn't working label Sep 6, 2026

@andygrove andygrove 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.

The follow-up commit and your replies cover the points I raised earlier. The guard question, the elapsed_compute assertion, and the runtime numbers all check out on a closer look. I checked the LogicalQueryStageStrategy and SparkStrategies.plan sequence you cited against Spark 3.4.3, 3.5.8, and 4.0.1, not only 4.1.0. The reuse-by-identity return and the direct setLogicalLink call are unchanged across all of them, so the guard in CometExecRule.scala should generalize the way you describe. I also traced the other places that copy a logicalLink onto a Comet node, in EliminateRedundantTransitions.scala and CometPlanAdaptiveDynamicPruningFilters.scala. Both build a brand new destination node in the same statement, so neither can hit the clobber this PR fixes. On the metric assertion, DataFusion's Time::add_duration floors every recorded duration at 1 nanosecond, so elapsed_compute > 0 can't read zero on a loaded runner either.

Your reply about the exchange cases mentions a separate local probe for the DPP broadcast scenario, where a direct LogicalQueryStage link can land on a CometBroadcastExchangeExec at the root of a DPP subquery's own adaptive plan. Leaving that case unguarded relies on Spark's cost evaluator rejecting a no-op replan, and on TEMP_LOGICAL_PLAN_TAG taking priority inside setLogicalLinkForNewQueryStage. That's real Spark behavior you found by testing, not something Comet's own code guarantees. The probe itself isn't part of this diff. Would it be worth committing it as a regression test next to the two already in CometExecRuleSuite? If either of those Spark internals shifts in a later release, this path could regress the same way the original bug did, with correct results and no native broadcast join.

Is #5482 still the right place to track whether the general link-repair mechanism should exist in its current form, or would a new issue be better? I don't see that broader question filed anywhere yet, and I'd rather it not end up living only in this comment thread.

@sunchao sunchao added the run-spark-4.1-tests Run the Spark 4.1 SQL tests on this pull request instead of waiting for the merge queue label Sep 19, 2026
@sunchao
sunchao force-pushed the fix/aqe-logical-stage-links-20260826 branch from f81c995 to 0d82334 Compare September 19, 2026 04:18
@sunchao sunchao added the run-all-spark-profiles Run the Comet test suites against every Spark profile on this pull request, ahead of the merge queue label Sep 19, 2026
@sunchao

sunchao commented Sep 19, 2026

Copy link
Copy Markdown
Member Author

Addressed the two remaining follow-ups from the September 9 review.

The new CometExecRuleSuite regression executes a grouped fact/dimension join and requires both native broadcast joins and Comet DPP, as well as the expected rows. Observers around Comet's AQE preparation rules require the direct LogicalQueryStage broadcast-root case to occur, then check that the physical tree and standard cost are unchanged and the original temporary logical identity survives Comet's repair. A captured copy with conflicting direct and temporary links calls Spark's actual stage-link assignment helper and verifies that the temporary link wins. The observer state is scoped to the test, and Spark 3.4 skips this native AQE DPP case.

The broader question of which originalPlan repairs remain necessary is now tracked in #6034. It covers link ownership, repeated AQE replanning, DPP broadcast roots, and the #323 clearing contract; #5482 stays focused on the reproduced stale-stage-link defect.

Committed in 69b2ca86a. Local validation on Spark 4.1.3 / JDK 17 passed all 35 planner tests and three targeted execution regressions (38 total), using the native library built from the rebased source. Spotless, Scalastyle, and git diff --check passed. The Spark SQL and all-Spark-profile CI labels are enabled; those checks are pending.

@sunchao
sunchao force-pushed the fix/aqe-logical-stage-links-20260826 branch from 1bc753b to c3bd613 Compare September 22, 2026 23:45

@andygrove andygrove 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.

LGTM. Thanks @sunchao

@sunchao
sunchao force-pushed the fix/aqe-logical-stage-links-20260826 branch 2 times, most recently from ae9ee74 to b2e625a Compare September 25, 2026 03:09
@sunchao
sunchao enabled auto-merge September 25, 2026 03:12
@sunchao
sunchao added this pull request to the merge queue Sep 25, 2026
@github-merge-queue
github-merge-queue Bot removed this pull request from the merge queue due to a conflict with the base branch Sep 25, 2026
@sunchao
sunchao force-pushed the fix/aqe-logical-stage-links-20260826 branch from 54d55d4 to 0b4580c Compare September 25, 2026 07:46
@sunchao
sunchao enabled auto-merge September 25, 2026 07:47
@sunchao
sunchao added this pull request to the merge queue Sep 25, 2026
Merged via the queue into apache:main with commit 88a1f48 Sep 25, 2026
62 checks passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

bug Something isn't working run-all-spark-profiles Run the Comet test suites against every Spark profile on this pull request, ahead of the merge queue run-spark-4.1-tests Run the Spark 4.1 SQL tests on this pull request instead of waiting for the merge queue

Projects

None yet

Development

Successfully merging this pull request may close these issues.

CometExecRule overwrites direct AQE LogicalQueryStage links during replanning

3 participants