Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
21 commits
Select commit Hold shift + click to select a range
f476dfe
fix: restore Spark write execs when reverting transition-heavy stages…
Sep 15, 2026
625f804
fix: drop unused string interpolator in transition-revert warning
sam-1112 Sep 15, 2026
11658a8
fix: remove unused import
sam-1112 Sep 15, 2026
b453c74
fix: remove unused import
sam-1112 Sep 16, 2026
713dc49
test: wait for write-plan listener callback before unregister
sam-1112 Sep 16, 2026
7798597
fix: drop duplicate Parquet write helpers that weaken base-class access
sam-1112 Sep 19, 2026
86e518f
fix: drop duplicate Parquet write helpers that weaken base-class access
sam-1112 Sep 19, 2026
fbe1d17
Merge branch 'main' into fix-5719-write-exec-revert
sam-1112 Sep 19, 2026
51d6568
Merge branch 'main' into fix-5719-write-exec-revert
sam-1112 Sep 22, 2026
4dbbd63
Merge branch 'main' into fix-5719-write-exec-revert
sam-1112 Sep 26, 2026
a3bdbec
fix: keep shuffle-stage transitions when reverting native writes (#6152)
sam-1112 Sep 26, 2026
4aafb06
Merge remote-tracking branch 'upstream/main' into fix-5719-write-exec…
sam-1112 Sep 27, 2026
82df177
fix: keep a native stage when transition stripping exposes an exchang…
sam-1112 Sep 28, 2026
d80e55a
Merge branch 'main' into fix-5719-write-exec-revert
sam-1112 Sep 28, 2026
0a2d5c3
fix: restore row output after stage reversion
sam-1112 Sep 29, 2026
a76e9a7
Merge branch 'main' into fix-5719-write-exec-revert
sam-1112 Sep 29, 2026
219c92f
Merge remote-tracking branch 'upstream/main' into fix-5719-write-exec…
sam-1112 Oct 2, 2026
f07c283
fix: preserve Arrow input for reverted native shuffle stages
sam-1112 Oct 2, 2026
d455691
Merge remote-tracking branch 'upstream/main' into fix-5719-write-exec…
sam-1112 Oct 5, 2026
b21e36e
refactor: move operator-specific restoration into sparkFallback
sam-1112 Oct 5, 2026
ee75eef
Merge remote-tracking branch 'upstream/main' into fix-5719-write-exec…
sam-1112 Oct 6, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
28 changes: 28 additions & 0 deletions docs/source/contributor-guide/adding_a_new_operator.md
Original file line number Diff line number Diff line change
Expand Up @@ -731,6 +731,34 @@ Use `QueryPlanSerde.exprToProto` to convert Spark expressions to protobuf:
val protoExpr = exprToProto(sparkExpr, inputSchema)
```

### Restoring the Spark operator (`sparkFallback`)

`CometExec.originalPlan` is the Spark operator this node replaced. `CometExecRule` copies
`originalPlan.logicalLink` onto the Comet node, which is how AQE finds the node again when it
re-plans a stage. `RevertNativeForTransitionHeavyStages` calls `sparkFallback(newChildren)` to

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 says the rule rebuilds the Spark operator through sparkFallback, but revertToSpark handles CometLocalTopKExec, CometNativeScanExec and CometIcebergNativeScanExec itself before it gets there (RevertNativeForTransitionHeavyStages.scala:263-281). So sparkFallback isn't the whole contract yet, and on CometLocalTopKExec the default would return a second TakeOrderedAndProjectExec that applies the TopK twice. Could those three become sparkFallback overrides, with the local TopK returning its child and the two scans carrying their live DPP filters across? Then the rule only ever calls sparkFallback, and this section holds for the next operator whose live state differs from its originalPlan. If you'd rather not move the code, could the section name the exceptions?

rebuild that Spark operator with the children of the reverted stage.

The default implementation is `originalPlan.withNewChildren(newChildren)`. It refuses a null
`originalPlan`, an `originalPlan` that is one of the node's own children, or a different number of
children than the Spark operator has.

Override `sparkFallback` when conversion changes the plan shape, so the restored node is not that
Spark operator with the same children. `CometNativeWriteExec` replaces a `DataWritingCommandExec`
and drops the `WriteFilesExec` under it; its override puts that wrapper back around the restored
input. `CometIcebergWriteExec` keeps the same shape as `IcebergWriteExec`, so the default is
enough.

Also override `sparkFallback` when the operator's live state differs from `originalPlan`.
`CometNativeScanExec` restores its current partition and data filters, and
`CometIcebergNativeScanExec` restores its current runtime filters, so AQE's executable DPP
subqueries survive reversion. `CometLocalTopKExec` returns the restored child directly: Comet
inserted that local candidate selection, and only the outer TopK restores Spark's offset and
projection. Rebuilding the original TopK at both nodes would apply it twice.

Do not point `originalPlan` at a child. If that child is a shuffle or query stage, the copied
logical link puts this node inside the stage's `LogicalQueryStage`. AQE then re-plans a second
copy of the operator around the one that is already there.

### Handling Fallback

Use `withInfo` to tag operators with fallback reasons:
Expand Down
13 changes: 8 additions & 5 deletions docs/source/contributor-guide/iceberg-writes.md
Original file line number Diff line number Diff line change
Expand Up @@ -107,9 +107,9 @@ serializer requires that tag on Spark 4.1+, leaving stock writers on Spark throu
## From `IcebergWrite` to `CometIcebergWrite`

`CometExecRule` converts an `IcebergWriteExec` with the `CometIcebergNativeWrite` operator serde
when `spark.comet.write.iceberg.enabled` is on. Two arms in `CometExecRule` handle it: one unwraps
the double conversion AQE can produce when it re-fires write planning over a sub-tree that already
contains a `CometIcebergWriteExec`, and the other calls `convertToComet`.
when `spark.comet.write.iceberg.enabled` is on. A single arm in `CometExecRule` handles the
conversion by calling `convertToComet`. The converted node keeps the `IcebergWriteExec` as its
`originalPlan`, so AQE re-plans the write from that node.

`CometIcebergNativeWrite.requiresNativeChildren` is `true`. The native writer consumes Arrow
batches from its child over FFI, so the conversion is declined unless the child is already a Comet
Expand Down Expand Up @@ -421,8 +421,11 @@ Each of these has caused a bug on this path:
([#5691](https://github.com/apache/datafusion-comet/issues/5691),
[#5693](https://github.com/apache/datafusion-comet/issues/5693),
[#6141](https://github.com/apache/datafusion-comet/issues/6141)).
- **Plan rewrites must keep the write node.** Rules that restore Spark operators from a Comet
node's `originalPlan` have to handle the write execs, whose `originalPlan` today is their child
- **Plan rewrites must keep the write node.** `CometIcebergWriteExec.originalPlan` is the
`IcebergWriteExec` it replaced. Restoring Spark execution goes through
[`CometExec.sparkFallback`](adding_a_new_operator.md#restoring-the-spark-operator-sparkfallback),
which rebuilds that node around the reverted children. Pointing `originalPlan` at the child
drops the write, and AQE then treats the write as part of that child's stage
([#5719](https://github.com/apache/datafusion-comet/issues/5719)).
- **The kill switch must still work.** Code that runs for Iceberg writes has to respect
`spark.comet.enabled`, so that disabling Comet restores Spark's own plan
Expand Down
6 changes: 6 additions & 0 deletions docs/source/contributor-guide/native_shuffle.md
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,12 @@ Compare this to JVM shuffle's data path:
Comet Native (columnar) → ColumnarToRowExec → rows → JVM Shuffle → Arrow IPC → columnar
```

When `RevertNativeForTransitionHeavyStages` restores the map stage to Spark execution, the
native exchange stays in place. Its input still needs Arrow-backed Comet vectors, even if the
restored Spark operator supports columnar output. The rule adds `CometSparkToColumnarExec` to
convert either Spark rows or Spark columnar batches to Arrow before the native shuffle consumes
them. Spark's `RowToColumnarExec` alone does not satisfy this input contract.

## When Native Shuffle is Used

Native shuffle (`CometExchange`) is selected when all of the following conditions are met:
Expand Down
20 changes: 6 additions & 14 deletions spark/src/main/scala/org/apache/comet/rules/CometExecRule.scala
Original file line number Diff line number Diff line change
Expand Up @@ -542,23 +542,15 @@ case class CometExecRule(session: SparkSession, queryStagePrep: Boolean = false)
// DataWritingCommandExec and re-implement the write framework inside CometNativeWriteExec.
// This path is retained only for 3.4/3.5 and goes away with them.
//
// AQE reoptimization looks for `DataWritingCommandExec` or `WriteFilesExec`
// if there is none it would reinsert write nodes, and since Comet remap those nodes
// to Comet counterparties the write nodes are twice to the plan.
// Checking if AQE inserted another write Command on top of existing write command
case _ @DataWritingCommandExec(_, w: WriteFilesExec)
if !isSpark40Plus && w.child.isInstanceOf[CometNativeWriteExec] =>
w.child

// `originalPlan` is that command. This rule copies `originalPlan.logicalLink` onto the
// Comet node, so AQE re-plans the write with the command rather than with whatever child
// happened to sit under it. No second DataWritingCommandExec is inserted on top.
case op: DataWritingCommandExec if !isSpark40Plus =>
convertToComet(op, CometDataWritingCommand).getOrElse(op)

// AQE re-fires the Iceberg write planning on every stage materialisation, so a
// partitioned write's physical sub-tree may already contain a `CometIcebergWriteExec`.
// Unwrap to avoid a double conversion.
case op: IcebergWriteExec if op.child.isInstanceOf[CometIcebergWriteExec] =>
op.child

// `originalPlan` is this IcebergWriteExec, so AQE re-plans the write as this node.
// A shuffle directly under the native write stays in the child stage and is not wrapped
// again.
case op: IcebergWriteExec if CometConf.COMET_ICEBERG_NATIVE_WRITE_ENABLED.get(op.conf) =>
convertToComet(op, CometIcebergNativeWrite).getOrElse(op)

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,8 @@ import org.apache.spark.internal.Logging
import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.catalyst.expressions.aggregate.{Final, Partial, PartialMerge}
import org.apache.spark.sql.catalyst.rules.Rule
import org.apache.spark.sql.comet.{CometBaseAggregateExec, CometColumnarToRowExec, CometExec, CometIcebergNativeScanExec, CometLocalTopKExec, CometNativeColumnarToRowExec, CometNativeScanExec, CometSparkToColumnarExec}
import org.apache.spark.sql.comet.{CometBaseAggregateExec, CometColumnarToRowExec, CometExec, CometNativeColumnarToRowExec, CometSparkToColumnarExec}
import org.apache.spark.sql.comet.execution.shuffle.{CometNativeShuffle, CometShuffleExchangeExec}
import org.apache.spark.sql.execution.{ColumnarToRowExec, ColumnarToRowTransition, RowToColumnarExec, SparkPlan}
import org.apache.spark.sql.execution.adaptive.QueryStageExec
import org.apache.spark.sql.execution.exchange.{BroadcastExchangeLike, ShuffleExchangeLike}
Expand Down Expand Up @@ -62,32 +63,41 @@ case class RevertNativeForTransitionHeavyStages(session: SparkSession, wholePlan
plan match {
case _: BroadcastExchangeLike => plan
case exchange: ShuffleExchangeLike =>
revertStageIfNeeded(exchange.child, exchange.supportsColumnar)
revertShuffleStageIfNeeded(exchange)
.map(reverted => exchange.withNewChildren(Seq(reverted)))
.getOrElse(plan)
case _ =>
// Result stage: its output is collected as rows, so no consumer requires columnar input
// and the reverted stage needs no trailing R2C.
// Result stage: its output is collected as rows.
revertStageIfNeeded(plan, outputColumnar = false).getOrElse(plan)

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 result stage still assumes it has to produce rows, and that isn't true for a cached query on Spark 4.0+. When the cache serializer accepts columnar input, which Comet's ArrowCachedBatchSerializer does, Spark marks the cached AdaptiveSparkPlanExec columnar and AQE plans the final stage with outputsColumnar = true. With the Comet cache on, AQE on, transitionRevert.enabled=true, maxTransitions=0 and spark.comet.exec.filter.enabled=false, caching SELECT _2, s FROM (SELECT _2, sum(_1) AS s FROM tbl GROUP BY _2) WHERE s > 10 and collecting it fails with FilterExec has column support mismatch. main fails the same way, so it isn't a regression, but it's the same output-format restoration this PR is fixing. Could both result-stage call sites (here and line 81) pass outputColumnar = plan.supportsColumnar && !plan.supportsRowBased instead of false? The root already says which format the consumer wants. I tried that locally and the repro passes, along with all of RevertNativeForTransitionHeavyStagesSuite. The test would need a suite that installs the cache serializer the way CometInMemoryCacheSuite does. If you'd rather keep this PR to the write fix, could we file an issue for it instead?

}
}

private def applyForNonAQE(plan: SparkPlan): SparkPlan = {
val withRevertedStages = plan.transformUp { case exchange: ShuffleExchangeLike =>
revertStageIfNeeded(exchange.child, exchange.supportsColumnar)
revertShuffleStageIfNeeded(exchange)
.map(reverted => exchange.withNewChildren(Seq(reverted)))
.getOrElse(exchange)
}
revertStageIfNeeded(withRevertedStages, outputColumnar = false)
.getOrElse(withRevertedStages)
}

private def revertShuffleStageIfNeeded(exchange: ShuffleExchangeLike): Option[SparkPlan] = {
val outputArrow = exchange match {
case comet: CometShuffleExchangeExec => comet.shuffleType == CometNativeShuffle
case _ => false
}
revertStageIfNeeded(exchange.child, exchange.supportsColumnar, outputArrow)
}

/**
* Reverts the stage if C2R count exceeds threshold. Wraps in R2C if exchange needs columnar.
* Reverts the stage if C2R count exceeds threshold, restoring the stage's output format when
* the reverted root does not satisfy it.
*/
private def revertStageIfNeeded(
stagePlan: SparkPlan,
outputColumnar: Boolean): Option[SparkPlan] = {
outputColumnar: Boolean,
outputArrow: Boolean = false): Option[SparkPlan] = {
val transitionCount = countTransitions(stagePlan)
if (transitionCount <= maxTransitions) return None

Expand All @@ -101,11 +111,27 @@ case class RevertNativeForTransitionHeavyStages(session: SparkSession, wholePlan
val reason =
s"Stage reverted: $transitionCount C2R transitions exceed threshold $maxTransitions"

val reverted = revertToSpark(stagePlan)
val result = if (outputColumnar && !reverted.supportsColumnar) {
RowToColumnarExec(withFallbackReason(reverted, reason))
val reverted =
try {
revertToSpark(stagePlan)
} catch {
case e: CometExec.InvalidSparkFallbackException =>
logWarning(
"Skipping transition-heavy stage reversion because a Comet operator could not " +
s"restore its Spark plan: ${e.getMessage}")
return None
}
val revertedWithReason = withFallbackReason(reverted, reason)
val result = if (outputArrow) {
// Native shuffle consumes Arrow-backed Comet vectors, not arbitrary Spark columnar
// batches. This bridge converts both row-based and vectorized Spark fallback roots.
CometSparkToColumnarExec(revertedWithReason)
} else if (outputColumnar && !reverted.supportsColumnar) {
RowToColumnarExec(revertedWithReason)
} else if (!outputColumnar && reverted.supportsColumnar) {
ColumnarToRowExec(revertedWithReason)
} else {
withFallbackReason(reverted, reason)
revertedWithReason
}
Some(result)
}
Expand Down Expand Up @@ -146,16 +172,31 @@ case class RevertNativeForTransitionHeavyStages(session: SparkSession, wholePlan
}

/**
* Like `transformDown`, never descends stage-boundary children.
* Like `transformDown`, never descends stage-boundary children. If the rule rewrites the
* current node, re-apply it to the result so stacked transitions such as
* `CometSparkToColumnarExec(CometNativeColumnarToRowExec(x))` are fully unwrapped before
* children are visited. Spark's `transformDown` does not do this; leaving the inner C2R in
* place later calls `CometNativeColumnarToRowExec.withNewChildren` with a reverted row-based
* child, which asserts `child.supportsColumnar`.
*
* A rewrite can itself be the stage boundary. Unwrapping a transition that sits directly on a
* shuffle yields that shuffle, and descending into it strips transitions in the next stage.
* `transformStageUp` and `insertTransitions` do not cross the exchange, so those transitions
* would not be restored (#6152). Return the boundary unchanged.
*/
private def transformStageDown(plan: SparkPlan)(
rule: PartialFunction[SparkPlan, SparkPlan]): SparkPlan = {
val transformed = rule.applyOrElse(plan, identity[SparkPlan])
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.

else transformStageDown(transformed)(rule)
} else {
val newChildren = transformed.children.map { child =>
if (isStageBoundary(child)) child else transformStageDown(child)(rule)
Comment on lines +190 to +195

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.

}
if (newChildren == transformed.children) transformed
else transformed.withNewChildren(newChildren)
}
if (newChildren == transformed.children) transformed
else transformed.withNewChildren(newChildren)
}

/** Like `transformUp`, never descends stage-boundary children. */
Expand Down Expand Up @@ -184,48 +225,40 @@ case class RevertNativeForTransitionHeavyStages(session: SparkSession, wholePlan
count
}

/**
* Checks for Comet operators whose original Spark plan is also their input. This must run
* before any bottom-up rewrite replaces children. Otherwise a stable `originalPlan` reference
* can keep pointing at the old Comet child after `withNewChildren`, hiding the alias and
* causing fallback to reconstruct that Comet child instead of a Spark operator.
*/
private def validateOriginalPlanAliases(plan: SparkPlan): Unit = plan match {
case _ if isStageBoundary(plan) => ()
case cometExec: CometExec =>
val sparkPlan = cometExec.originalPlan
if (sparkPlan != null && cometExec.children.exists(_ eq sparkPlan)) {
throw new CometExec.InvalidSparkFallbackException(
s"${cometExec.getClass.getSimpleName} aliases its original Spark plan with a child")
}
cometExec.children.foreach(validateOriginalPlanAliases)
case _ =>
plan.children.foreach(validateOriginalPlanAliases)
}

private[rules] def revertToSpark(plan: SparkPlan): SparkPlan = {
validateOriginalPlanAliases(plan)
val stripped = transformStageDown(plan) {
case CometNativeColumnarToRowExec(child) => child
case CometColumnarToRowExec(child) => child
case ColumnarToRowExec(child) => child
case sparkToColumnar: CometSparkToColumnarExec => sparkToColumnar.child
case RowToColumnarExec(child) => child
}
val reverted = transformStageUp(stripped) {
// Local candidate selection was inserted by Comet. Only the outer TopK owns
// the original Spark operator's offset and projection.
case local: CometLocalTopKExec => local.child
case cometExec: CometExec =>
if (cometExec.originalPlan.children.size == cometExec.children.size) {
val originalWithCurrentExpressions = cometExec match {
case scan: CometNativeScanExec =>
// AQE's query-stage optimizer rewrites DPP placeholders in the live Comet scan
// before this post-columnar rule runs. The frozen FileSourceScanExec in
// originalPlan still contains SubqueryAdaptiveBroadcastExec, which cannot execute.
// Preserve the rewritten filters when reverting the scan to Spark.
val originalScan = scan.originalPlan.copy(
partitionFilters = scan.partitionFilters,
dataFilters = scan.dataFilters)
scan.originalPlan.logicalLink.foreach(originalScan.setLogicalLink)
originalScan
case scan: CometIcebergNativeScanExec =>
// Iceberg's native scan has the same split between live and frozen filters.
// serializedPartitionData rebuilds originalPlan from runtimeFilters before
// execution, but transition reversion skips that path and executes the restored
// BatchScanExec directly. Carry the executable DPP filters across here as well.
val originalScan = scan.originalPlan.copy(runtimeFilters = scan.runtimeFilters)
scan.originalPlan.logicalLink.foreach(originalScan.setLogicalLink)
originalScan
case _ => cometExec.originalPlan
}
originalWithCurrentExpressions.withNewChildren(cometExec.children)
} else {
logWarning(
"Comet plan and original have different child count for " +
s"${cometExec.getClass.getSimpleName}, using originalPlan as-is.")
cometExec.originalPlan
}
if (isStageBoundary(stripped)) {
throw new CometExec.InvalidSparkFallbackException(
"Cannot revert a stage whose stripped root is a stage boundary")
}
val reverted = transformStageUp(stripped) { case cometExec: CometExec =>
cometExec.sparkFallback(cometExec.children)
}
insertTransitions(reverted)
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -211,7 +211,7 @@ object CometDataWritingCommand extends CometOperatorSerde[DataWritingCommandExec
throw new SparkException(s"Could not instantiate FileCommitProtocol: ${e.getMessage}")
}

CometNativeWriteExec(nativeOp, childPlan, outputPath, cmd.mode, committer, jobId)
CometNativeWriteExec(nativeOp, op, childPlan, outputPath, cmd.mode, committer, jobId)
}

}
Original file line number Diff line number Diff line change
Expand Up @@ -647,7 +647,9 @@ object CometIcebergNativeWrite extends CometOperatorSerde[IcebergWriteExec] {
"Native Iceberg write conversion: SparkWrite.outputSpecId reflection failed"))
CometIcebergWriteExec(
nativeOp,
op,
op.child,
op.output,
op.batchWrite,
table.asInstanceOf[AnyRef],
outputSpecId)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,7 @@ import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.catalyst.expressions.{Attribute, Expression, SortOrder}
import org.apache.spark.sql.catalyst.plans.QueryPlan
import org.apache.spark.sql.catalyst.plans.physical.{Partitioning, UnknownPartitioning}
import org.apache.spark.sql.execution.SQLExecution
import org.apache.spark.sql.execution.{SparkPlan, SQLExecution}
import org.apache.spark.sql.execution.datasources.v2.BatchScanExec
import org.apache.spark.sql.execution.metric.{SQLMetric, SQLMetrics}
import org.apache.spark.sql.vectorized.ColumnarBatch
Expand Down Expand Up @@ -65,6 +65,14 @@ case class CometIcebergNativeScanExec(
@transient nativeIcebergScanMetadata: CometIcebergNativeScanMetadata)
extends CometLeafExec {

override def sparkFallback(newChildren: Seq[SparkPlan]): SparkPlan = {
// Native execution rebuilds originalPlan from the live runtimeFilters during partition
// serialization. Reversion skips that path, so carry the executable DPP filters across here.
val restoredScan = originalPlan.copy(runtimeFilters = runtimeFilters)
originalPlan.logicalLink.foreach(restoredScan.setLogicalLink)
restoredScan
}

override val supportsColumnar: Boolean = true

override val nodeName: String = "CometIcebergNativeScan"
Expand Down
Loading
Loading