Conversation
b0a9862 to
812b06c
Compare
sunchao
left a comment
There was a problem hiding this comment.
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) |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
Fixed. The helper now waits on a bounded CountDownLatch from onSuccess while the listener is still registered, then unregisters in finally.
sunchao
left a comment
There was a problem hiding this comment.
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
left a comment
There was a problem hiding this comment.
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.
76da109 to
08bec89
Compare
sunchao
left a comment
There was a problem hiding this comment.
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 = { |
There was a problem hiding this comment.
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.
…apache#5719) Give Iceberg and parquet native writes a real originalPlan so revertToSpark keeps the commit-message node instead of duplicating or dropping the child.
8c56374 to
7798597
Compare
sunchao
left a comment
There was a problem hiding this comment.
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
left a comment
There was a problem hiding this comment.
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
left a comment
There was a problem hiding this comment.
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.
| if (transformed ne plan) { | ||
| transformStageDown(transformed)(rule) | ||
| } else { | ||
| val newChildren = transformed.children.map { child => | ||
| if (isStageBoundary(child)) child else transformStageDown(child)(rule) |
There was a problem hiding this comment.
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") { |
There was a problem hiding this comment.
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
left a comment
There was a problem hiding this comment.
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 |
There was a problem hiding this comment.
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.
|
This is a light fully automated review since there are so many PRs open. The new check at |
Which issue does this PR close?
Closes #5719.
Rationale for this change
RevertNativeForTransitionHeavyStages.revertToSparktreated everyCometExecas a like-for-like swap oforiginalPlan.CometIcebergWriteExecandCometNativeWriteExecreported their own child asoriginalPlan, so reverting a transition-heavy write stage erased the write node:child.withNewChildren(Seq(child)))After that,
IcebergCommitExecno longer receivediceberg_commit_messagerows. On Spark 3.5 the reverted native Parquet plan is justRange: 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:
CometIcebergWriteExec.originalPlanis theIcebergWriteExecit replacedCometNativeWriteExec.originalPlanis theDataWritingCommandExecit replacedCometExec.sparkFallbackinstead of graftingoriginalPlanonto itself.sparkFallbackso a revertedWriteFilesExec(when present) keeps wrapping the restored input.originalPlanis one of their own children before rewriting. If restore is invalid, skip reversion for the whole stage rather than emitting a broken write plan.transformStageDownafter rewriting a node so stacked transitions such asSparkToColumnar(C2R(...))unwrap fully.transformStageUpandinsertTransitionsdo not cross the exchange, so those transitions were never restored (RevertNativeForTransitionHeavyStages strips transitions in the stage below when AQE is off #6152).CometExecRulearms that unwrapped a second Spark write already sitting on a native write. They only matched the oldoriginalPlan = childlink (see below).iceberg-writes.mdpoints the write-node pitfall atsparkFallback, andadding_a_new_operator.mdsays 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.
CometExecRulecopiesoriginalPlan.logicalLinkonto everyCometExec. On main the write execs used their child asoriginalPlan, so they picked up that child's link. When the child was a shuffle stage, AQE folded the write into the stage'sLogicalQueryStageand the re-plan came back double-wrapped: a Spark write operator around the native write that was already there. Two arms inCometExecRuleexisted only to unwrap that shape:DataWritingCommandExec(_, WriteFilesExec)whose child isCometNativeWriteExecIcebergWriteExecwhose child isCometIcebergWriteExecoriginalPlanis now theDataWritingCommandExecorIcebergWriteExecthat 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
revertToSparkcoverage inRevertNativeForTransitionHeavyStagesSuite:SparkToColumnar, and over a unary Comet child (no duplicatedFilter)SparkToColumnar(C2R)under a native writeWriteFilesExecoriginalPlan == childaliases leave the stage unchangedEnd-to-end, with AQE on and off:
transitionRevert.enabled=trueandmaxTransitions=0restoresIcebergWriteExecand still writes the rowsDELETE ... 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 withColumnarBatch 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.DataWritingCommandExec→WriteFilesExecand write the expected row counts