Skip to content

fix: restore Spark write execs when reverting transition-heavy stages - #5957

Open
sam-1112 wants to merge 16 commits into
apache:mainfrom
sam-1112:fix-5719-write-exec-revert
Open

sam-1112 wants to merge 16 commits into
apache:mainfrom
sam-1112:fix-5719-write-exec-revert

Conversation

@sam-1112

@sam-1112 sam-1112 commented Sep 15, 2026 •

Copy link
Copy Markdown
Contributor

Which issue does this PR close?

Closes #5719.

Rationale for this change

RevertNativeForTransitionHeavyStages.revertToSpark treated every CometExec as a like-for-like swap of originalPlan. CometIcebergWriteExec and CometNativeWriteExec reported their own child as originalPlan, so reverting a transition-heavy write stage erased the write node:

  • a leaf or multi-child input dropped the write entirely
  • a unary input also duplicated the child (child.withNewChildren(Seq(child)))

After that, IcebergCommitExec no longer received iceberg_commit_message rows. On Spark 3.5 the reverted native Parquet plan is just Range: the write reports success and no output directory is created. The same aliasing is never a valid Spark restore for any operator.

This is pre-existing and independent of #5696. Part of #5649.

What changes are included in this PR?

This implements both fixes from #5719:

  • Give the write execs a real Spark original plan:
    • CometIcebergWriteExec.originalPlan is the IcebergWriteExec it replaced
    • CometNativeWriteExec.originalPlan is the DataWritingCommandExec it replaced
  • Restore through CometExec.sparkFallback instead of grafting originalPlan onto itself.
  • Native Parquet writes override sparkFallback so a reverted WriteFilesExec (when present) keeps wrapping the restored input.
  • Reject operators whose originalPlan is one of their own children before rewriting. If restore is invalid, skip reversion for the whole stage rather than emitting a broken write plan.
  • Re-apply transformStageDown after rewriting a node so stacked transitions such as SparkToColumnar(C2R(...)) unwrap fully.
  • Stop at a stage boundary when that rewrite is itself the boundary. Unwrapping a transition that sits directly on a shuffle used to keep walking into the next stage. transformStageUp and insertTransitions do not cross the exchange, so those transitions were never restored (RevertNativeForTransitionHeavyStages strips transitions in the stage below when AQE is off #6152).
  • Remove the two CometExecRule arms that unwrapped a second Spark write already sitting on a native write. They only matched the old originalPlan = child link (see below).
  • Document the restore contract. iceberg-writes.md points the write-node pitfall at sparkFallback, and adding_a_new_operator.md says when an operator must override it.

AQE re-plan of native writes

This changes how AQE re-plans a native write, not only how fallback restores it.

CometExecRule copies originalPlan.logicalLink onto every CometExec. On main the write execs used their child as originalPlan, so they picked up that child's link. When the child was a shuffle stage, AQE folded the write into the stage's LogicalQueryStage and the re-plan came back double-wrapped: a Spark write operator around the native write that was already there. Two arms in CometExecRule existed only to unwrap that shape:

  • Spark 3.x: DataWritingCommandExec(_, WriteFilesExec) whose child is CometNativeWriteExec
  • IcebergWriteExec whose child is CometIcebergWriteExec

originalPlan is now the DataWritingCommandExec or IcebergWriteExec that was replaced, so AQE re-plans the write as that operator. A shuffle directly under a native Iceberg write stays in the child stage and is not wrapped again. Those arms no longer match, the suites stay green without them, and they are removed. The comments that remain describe this contract rather than the old re-plan.

How are these changes tested?

Direct revertToSpark coverage in RevertNativeForTransitionHeavyStagesSuite:

  • Iceberg write over a leaf, over SparkToColumnar, and over a unary Comet child (no duplicated Filter)
  • stacked SparkToColumnar(C2R) under a native write
  • native Parquet write with and without WriteFilesExec
  • invalid originalPlan == child aliases leave the stage unchanged
  • unwrapping a transition that sits on a shuffle leaves the transitions in the stage below it

End-to-end, with AQE on and off:

  • Iceberg INSERT with transitionRevert.enabled=true and maxTransitions=0 restores IcebergWriteExec and still writes the rows
  • Partitioned copy-on-write DELETE ... WHERE id IN (SELECT ...) restores the write, keeps a shuffle under it, and matches the rows and partition directories of a sibling table written natively. The AQE-off case is the one that used to fail with ColumnarBatch cannot be cast to InternalRow (RevertNativeForTransitionHeavyStages strips transitions in the stage below when AQE is off #6152). An unpartitioned INSERT never puts an exchange under the write, so it does not catch that.
  • Parquet writes (including a union of row sources) restore DataWritingCommandExec → WriteFilesExec and write the expected row counts
./mvnw test -Dtest=none -Dsuites="org.apache.comet.rules.RevertNativeForTransitionHeavyStagesSuite"
./mvnw test -Dtest=none -Dsuites="org.apache.comet.CometIcebergWriteActionSuite transition-heavy fallback preserves Iceberg writes"
./mvnw test -Dtest=none -Dsuites="org.apache.comet.CometIcebergWriteActionSuite transition-heavy fallback preserves partitioned CoW"
./mvnw test -Dtest=none -Dsuites="org.apache.comet.parquet.CometParquetWriterSuite transition-heavy fallback"

@sam-1112
sam-1112 marked this pull request as draft September 15, 2026 13:43
@github-actions github-actions Bot added bug Something isn't working area:writer Native Parquet writer area:Iceberg labels Sep 15, 2026
@sam-1112
sam-1112 force-pushed the fix-5719-write-exec-revert branch from b0a9862 to 812b06c Compare September 16, 2026 02:55
@sam-1112
sam-1112 marked this pull request as ready for review September 16, 2026 02:58

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

Correctness

The previous fallback could discard a native writer because its originalPlan was the writer's input, rather than the Spark write operator. For a unary input, restoring children could also reconstruct the wrong shape. This change keeps the original DataWritingCommandExec or IcebergWriteExec, restores the Parquet WriteFilesExec wrapper when present, and removes stacked row/columnar transitions before rebuilding the stage. It also rejects missing, aliased, or arity-incompatible fallback plans and leaves an unsafe stage unreverted.

The restored Parquet shape agrees with the maintained Spark 3.5 and 4.0 sources: DataWritingCommandExec owns command execution, and FileFormatWriter uses WriteFilesExec.executeWrite for planned writes. Retaining the original command and wrapper preserves their write metadata and save-mode handling. For Iceberg, carrying the original write node and stable commit-message output preserves the relationship with the outer commit operator and its existing commit/abort path. The native writer execution methods are unchanged. I found no additional production correctness issue in the reviewed changes. Spark 3.4 and 4.1 source compatibility remains unverified because the required maintained branches were unavailable.

One [P2] test reliability issue remains, detailed inline: captureDataWritingCommand unregisters its listener before asynchronous callback completion is guaranteed, so three new regression tests can fail after a successful write. Waiting for the callback before unregistering fixes that ordering.

Validation

Reviewed all 10 changed files at 812b06c9 against 4e69a248, including the leaf/unary fallback cases, alias guard, wrapper restoration, and the added AQE-on/off write tests. A compiled Scala probe containing the exact new helper reproduced the delayed-callback failure, with passing immediate-delivery and wait-before-unregister controls. This was an isolated test-double check, not a Spark/JNI run. I did not execute the full Comet suites locally. At September 16, 03:50 UTC, CI and CodeQL were action_required with zero jobs. The successful label job provides no test coverage.

Performance

The added alias validation makes one extra linear walk of an eligible stage during driver-side planning. It runs only when transition-heavy reversion is enabled and the transition threshold is exceeded. That feature remains disabled by default. The transition-stripping recursion removes wrappers, and the changes add no per-row work or additional write execution. I found no material new hot-path overhead requiring a separate benchmark. No runtime speedup, memory improvement, or benchmark result is claimed by this review.

Design

Keeping the actual Spark writer as the fallback source is a sound way to preserve ownership of write and commit behavior. The generic bottom-up reconstruction handles ordinary Comet nodes, while the Parquet override reconstructs the extra WriteFilesExec layer that conversion previously removed. Checking aliases before rewriting children is necessary because child replacement could otherwise hide the old reference relationship. Catching only the dedicated fallback exception keeps this conservative decision local to an unsafe restoration without swallowing unrelated failures. The source changes and the plan-shape tests address the original defect directly. The callback wait is the concrete improvement needed for dependable regression coverage.

Abstraction & complexity

sparkFallback(newChildren) gives restoration a narrow operator-level extension point instead of embedding write-specific reconstruction in the stage rule. Its shared null/alias/arity checks and the single Parquet specialization earn their complexity because these operators have different restoration shapes. The dedicated exception makes the rule's refusal to revert explicit. I found no additional abstraction or simplification issue that warrants a separate finding.

spark.range(1).toDF("id").write.mode("overwrite").parquet(path)
}
} finally {
spark.listenerManager.unregister(listener)

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.

Correctness

[P2] Wait for the write callback before unregistering the listener

A successful parquet(path) call does not guarantee that QueryExecutionListener.onSuccess has run: Spark posts the SQL execution-end event to the shared asynchronous listener queue. If that queue delivers the event after this finally block, the listener has already been removed and captured stays null, so the three new Parquet reversion tests can fail with expected a captured parquet write plan even though the write succeeded. I verified this delivery path in the maintained Spark 3.5 and 4.0 branches and reproduced the failing schedule with the exact helper in isolated Scala test doubles. Immediate delivery and a wait-before-unregister control both pass. Please wait with a bounded latch/future while the listener remains registered, then unregister in finally, so completion and visibility of the captured plan are guaranteed.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Fixed. The helper now waits on a bounded CountDownLatch from onSuccess while the listener is still registered, then unregisters in finally.

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

Re-reviewed 6d7be018. The previous P2 listener race is addressed: the helper stores the plan before releasing the latch, waits up to ten seconds while the listener remains registered, and unregisters in finally. The production changes are unchanged from the previous review. I found no new or remaining P1/P2 issues.

The exact old helper still reproduces the delayed-callback failure. The exact updated helper passed immediate and delayed delivery, visibility, unrelated-event filtering, write failure, real ten-second timeout, interruption, and cleanup checks in a compiled Scala probe. These are isolated test doubles, not Spark/JNI execution. At September 16, 05:28 UTC, CI and CodeQL remain action_required with zero jobs.

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

Rechecked 76da1090bcab against 58ab5f618e1e, including the full ten-file change and the base merge since 6d7be0184a31. No new or remaining verified P1/P2 findings. Preserving the existing approval.

The authored edits are unchanged. Parquet fallback retains the Spark write command and restores its WriteFilesExec wrapper when present. Iceberg fallback preserves the original writer, output attributes and commit boundary. Null/alias/arity guards, stacked transition removal and AQE stage boundaries remain intact. The inherited plan-cache and scan changes introduce no additional fallback edits.

The listener-race fix is also unchanged: it publishes the captured plan before signaling completion and waits while the listener is registered. The earlier eleven-scenario isolated probe remains historical evidence. I did not rerun it or execute Spark/JNI/write-IO tests on this head. Source checks used the maintained Spark 3.5/4.0 branches. Versions 3.4/4.1 remain unavailable.

At 2026-09-17 04:37:19 UTC, current-head CI and CodeQL require approval with zero jobs. Only labeling passed.

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

Rechecked 08bec891 against 4bf725f6 and the last reviewed head 76da1090. One new [P1] build blocker is detailed inline: the rebase retained two private Parquet test helpers that now conflict with protected methods in CometParquetWriterTestBase. Both fail with weaker-access errors in isolated compiler checks on Scala 2.12.18 and 2.13.17. Removing the duplicates makes the control compile.

The production write-restoration changes, whole-stage guards and prior listener-race fix remain intact. I also checked the inherited empty-relation and transition integrations against the maintained Spark 3.5/4.0 sources. Spark 3.4/4.1 source coverage remains unavailable.

The earlier eleven-scenario listener probe is historical evidence and was not rerun. No full Comet, Spark/JNI or write-IO suite was executed. Current Comet CI and CodeQL require approval with zero jobs. Label success provides no product validation.

* @return
* The captured execution plan
*/
private def captureWritePlan(writeOp: String => Unit, outputPath: String): SparkPlan = {

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.

Correctness

[P1] Remove the duplicate helpers after the base-class extraction

Could you remove this retained captureWritePlan and assertNoCometNativeWriteExec at line 1123, and use the inherited implementations? The rebased suite now extends CometParquetWriterTestBase, which already defines both methods as protected. These private declarations override them with weaker access, so Spark test compilation fails before any suite can run. I reproduced both compiler errors with the exact signatures and inheritance on Scala 2.12.18 and 2.13.17. Removing the duplicate declarations makes the isolated control compile. Their bodies already match the inherited implementations.

@sam-1112
sam-1112 force-pushed the fix-5719-write-exec-revert branch from 8c56374 to 7798597 Compare September 19, 2026 01:02

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

Rechecked 86e518f3 against 5fdc9619. The P1 compilation blocker is fixed: the suite now uses the inherited protected helpers. I reran the original failure case and the updated inheritance/call-site control on Scala 2.12.18 and 2.13.17. The old declarations fail with both weaker-access errors, and the updated controls compile.

The listener-race fix and write-restoration behavior remain intact. The base update preserves the Parquet command/wrapper and Iceberg writer/commit boundaries. No new or remaining verified P1/P2 findings.

These are isolated compiler checks with trivial bodies and a SparkPlan stand-in, plus source validation against maintained Spark 3.5/4.0. The earlier eleven-scenario listener probe was not rerun. Spark 3.4/4.1 source coverage remains unavailable, and no full Comet, Spark/JNI or write-IO suite was executed. At September 19, 02:01 UTC, Comet CI and CodeQL require approval with zero jobs. Only labeling passed.

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

Rechecked fbe1d17c against b7f35b6a. No new or remaining verified P1/P2 findings. Preserving the existing approval.

All ten authored files are unchanged from 86e518f3. The eighteen changed files match the update from main. The registry and native length changes introduce no additional interaction with write restoration. Parquet fallback still restores the original command and optional WriteFilesExec wrapper. Iceberg fallback retains the original writer, output attributes and commit boundary.

The P1 helper-visibility fix and P2 listener-race fix remain intact. I verified source identity with the earlier isolated Scala compiler and listener evidence. Those probes were not rerun. Source validation used the maintained Spark 3.5/4.0 branches. Spark 3.4/4.1 source coverage remains unavailable, and no full Comet, Spark/JNI or write-IO suite was executed.

At September 19, 12:24 UTC, Comet CI and CodeQL require approval with zero jobs. Only labeling passed, which provides no product validation.

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

I ran this merged into current main on 3.4, 3.5, 4.0, 4.1 and 4.2, and the fix holds. With main's originalPlan = child restored, the new unit tests and both Iceberg end-to-end tests fail on 4.1, and all eight end-to-end tests fail on 3.5. It also fixes something worse than #5719 describes. On 3.5, main's reverted plan for a native Parquet write is just Range, so the write reports success and no output directory is ever created.

This also changes how AQE re-plans native writes, not only how they revert. CometExecRule copies originalPlan.logicalLink onto every CometExec, so on main the write execs picked up their child's link. When that child was a shuffle stage, AQE folded the write into the stage's LogicalQueryStage, and the re-plan came back double-wrapped. That is what the two unwrap arms in CometExecRule.scala:441-459 handle. With this PR both arms stop firing, and the suites stay green. That looks right to me, but could we say so in the PR description? Could we also either remove those arms or update their comments so they stop describing a re-plan that no longer happens?

docs/source/contributor-guide/iceberg-writes.md:374-376 still says the write execs' originalPlan is their child. The branch predates that page, so it needs a merge from main first. Could that pitfall point at sparkFallback, with a sentence in adding_a_new_operator.md about when an operator needs to override it? I'm adding run-iceberg-tests as well, since this touches every AQE re-plan of a native Iceberg write that has a shuffle directly below it.

Comment on lines +164 to +168
if (transformed ne plan) {
transformStageDown(transformed)(rule)
} else {
val newChildren = transformed.children.map { child =>
if (isStageBoundary(child)) child else transformStageDown(child)(rule)

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.

With AQE off this still descends into the stage below an exchange. The boundary check only runs on the children of the transformed node. So when the stripped transition sits directly on a shuffle, the recursion strips the transitions inside the next stage down, and nothing puts them back, because transformStageUp and insertTransitions stop at the exchange. This is the problem in #6152, and here's a concrete repro: a copy-on-write DELETE ... WHERE id IN (SELECT ...) on a partitioned table, with spark.sql.adaptive.enabled=false and maxTransitions=0, fails the shuffle map stage with ColumnarBatch cannot be cast to InternalRow. Since this PR already rewrites transformStageDown, could it return transformed unchanged when it is a stage boundary, with if (isStageBoundary(transformed)) transformed else transformStageDown(transformed)(rule)? With that, the same DELETE commits once with the right rows and the same partition layout as the native write, and the rest of the suite still passes.

}

for (adaptive <- Seq(false, true)) {
test(s"transition-heavy fallback preserves Iceberg writes with AQE=$adaptive") {

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 end-to-end case only covers an unpartitioned INSERT ... VALUES. There is no exchange under the write, and the restored IcebergWriteExec has no ReplaceDataDispatchInfo. Could we add a copy-on-write DELETE against a partitioned table, with AQE on and off, and compare its rows and partition directories with a sibling table written natively? That is the one shape where the restored node carries state the unit tests build by hand. With AQE off it's also the case that catches the boundary problem above.

…che#6152)

Document sparkFallback and remove the CometExecRule unwraps that only matched the old originalPlan = child link.
…-revert

# Conflicts:
#	spark/src/main/scala/org/apache/comet/rules/RevertNativeForTransitionHeavyStages.scala
#	spark/src/test/scala/org/apache/comet/rules/RevertNativeForTransitionHeavyStagesSuite.scala

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

I ran this at 4aafb0629 on 3.4, 3.5, 4.0 and 4.1, and it covers everything from my last review. With the new boundary check taken out, the new unit test fails, and the AQE-off copy-on-write DELETE fails again with ColumnarBatch cannot be cast to InternalRow, so the tests really do catch #6152. I also checked the removed CometExecRule arms. If the write execs take their child's logical link, as on main, AQE builds the double-wrapped plan again and 13 tests fail on 4.1 (11 on 3.5). With this PR's link it never shows up in either write suite. Removing the arms looks right to me.

There's one more shape where the revert still reaches into the stage below. I've left it inline. It fails the same way on main, but it's the case #6152 describes, and this PR is already changing that code.

I've also added the run-iceberg-tests label I mentioned last time. Iceberg's own extension tests pick AQE at random, and this changes how every AQE re-plan of a native write is built.

val newChildren = transformed.children.map { child =>
if (isStageBoundary(child)) child else transformStageDown(child)(rule)
if (transformed ne plan) {
if (isStageBoundary(transformed)) transformed

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.

This stops the strip at the exchange, but if the stage's root is itself the transition sitting on the exchange, stripped is the exchange. Then transformStageUp and insertTransitions are called with a boundary as the root. They only check children for boundaries, so they walk into the map stage. With AQE off, transitionRevert.enabled=true and maxTransitions=0, SELECT _1, _2 FROM tbl DISTRIBUTE BY _2 plans as CometColumnarToRow over CometExchange over CometNativeScan. The revert turns the scan into Spark's FileScan under the still-native shuffle and drops the row transition at the top, and the query fails with a ClassCastException casting Spark's OnHeapColumnVector to a Comet vector. Could revertToSpark leave the stage alone when stripped is a stage boundary? When I tried if (isStageBoundary(stripped)) throwing InvalidSparkFallbackException, the query returned the right rows, the plan stayed native, and the rest of RevertNativeForTransitionHeavyStagesSuite passed. A test with that DISTRIBUTE BY query would cover it, and then this PR could close #6152 as well.

@andygrove

Copy link
Copy Markdown
Member

This is a light fully automated review since there are so many PRs open.

The new check at RevertNativeForTransitionHeavyStages.scala:239 covers the AQE-off shape, but I think SELECT _1, _2 FROM tbl DISTRIBUTE BY _2 with maxTransitions=0 still breaks with AQE on, which is the default. AQE coalesces that small shuffle, so the result stage reaches this rule as ColumnarToRowExec(AQEShuffleReadExec(ShuffleQueryStageExec)). stripped is then the AQEShuffleReadExec, which isStageBoundary doesn't match, so the check doesn't fire. insertTransitions only adds a ColumnarToRowExec under a row-based parent, and revertStageIfNeeded at line 114 only fixes up the columnar-output direction, so the stage comes back as a bare columnar AQEShuffleReadExec. Its doExecute just casts the CometShuffledBatchRDD, so I'd expect executeCollect to fail casting a ColumnarBatch to UnsafeRow. I think any row-output stage whose reverted root stays columnar has the same gap, for example a scan-only query where CometNativeScanExec reverts to a vectorized FileSourceScanExec. Could revertStageIfNeeded wrap the result in ColumnarToRowExec when !outputColumnar && reverted.supportsColumnar, mirroring the RowToColumnarExec branch? Could the DISTRIBUTE BY test at RevertNativeForTransitionHeavyStagesSuite.scala:665 also run with AQE on?

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

Labels

area:Iceberg area:writer Native Parquet writer 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-iceberg-tests

Projects

None yet

Development

Successfully merging this pull request may close these issues.

revertToSpark erases CometIcebergWriteExec / CometNativeWriteExec because originalPlan is the node's own child

3 participants