andygrove opened a new issue, #5719:
URL: https://github.com/apache/datafusion-comet/issues/5719
### Describe the bug
`RevertNativeForTransitionHeavyStages.revertToSpark` erases the write node
when it reverts a stage containing `CometIcebergWriteExec` or
`CometNativeWriteExec`, because both report their own child as `originalPlan`.
The revert arm is:
```scala
val reverted = transformStageUp(stripped) { case cometExec: CometExec =>
if (cometExec.originalPlan.children.size == cometExec.children.size) {
cometExec.originalPlan.withNewChildren(cometExec.children)
} else {
logWarning(...)
cometExec.originalPlan
}
}
```
For every other `CometExec` this is a like-for-like swap, because
`originalPlan` is the Spark operator the Comet node replaced. The two write
execs are different:
- `CometIcebergWriteExec.originalPlan` is `child`
(`CometIcebergWriteExec.scala:73`)
- `CometNativeWriteExec.originalPlan` is `child`
(`CometNativeWriteExec.scala:80`)
So `originalPlan.withNewChildren(cometExec.children)` reduces to
`child.withNewChildren(Seq(child))`. The write node disappears, and whichever
branch is taken depends only on the child's arity:
- **Leaf or multi-child child** (arity mismatch): the warning branch returns
the child as-is. The write is gone.
- **Unary child** (sizes both 1): the child is grafted onto itself, so the
write is gone *and* the child operator is duplicated.
Downstream, `IcebergCommitExec` calls `child.executeCollect()` and
deserialises each row's binary column as a `WriterCommitMessage`. After the
revert it is handed data rows with the table's own schema instead of the single
`iceberg_commit_message` binary column, so no data is written and the commit
either fails to deserialise or commits garbage.
The duplication is a second, independent defect: it double-applies the
operator. Harmless for an idempotent `Filter`, not harmless for a projection
carrying `monotonically_increasing_id()`, a sample, or a limit.
### Steps to reproduce
Calling `revertToSpark` directly on a write plan shows both shapes. With a
leaf child:
```
=== INPUT ===
CometIcebergWrite [iceberg_commit_message#0], , ICEBERG_WRITER_UNPARTITIONED
+- CometSparkRowToColumnar
+- <row source>
=== REVERTED ===
<row source>
write present: false
```
With a unary Comet child:
```
=== INPUT ===
CometIcebergWrite [iceberg_commit_message#6], , ICEBERG_WRITER_UNPARTITIONED
+- CometFilter [_1#4], (isnotnull(_1#4) AND (_1#4 > 2))
+- CometNativeScan parquet [_1#4] ...
=== REVERTED ===
Filter (isnotnull(_1#4) AND (_1#4 > 2))
+- Filter (isnotnull(_1#4) AND (_1#4 > 2))
+- ColumnarToRow
+- FileScan parquet [_1#4] ...
write present: false
```
Note the `iceberg_commit_message` output column is gone in both, and the
`Filter` is duplicated in the second.
To reach this through a real query rather than a direct call, the write's
stage has to exceed the transition threshold:
- `spark.comet.exec.transitionRevert.enabled=true` (default `false`, so this
is opt-in)
- more than `spark.comet.exec.transitionRevert.maxTransitions` (default 2)
columnar-to-row transitions between the write and the first exchange below it,
for instance with `spark.comet.sparkToColumnar.enabled` and a row-based `Union`
of Spark-columnar scans directly under the write
### Expected behavior
Reverting a stage should either leave the write node in place (only the
operators below it are candidates for reverting) or replace it with the JVM
write operator it stands in for, which for `CometIcebergWriteExec` is
`IcebergWriteExec`. It should never erase the node that produces the
commit-message rows, and it should never duplicate a child.
Two candidate fixes:
1. Stop `revertToSpark` from reverting a node whose `originalPlan` is one of
its own children. That is a general guard: `originalPlan = child` means "I have
no distinct Spark original", and grafting a child onto itself is never correct.
2. Give the write execs a real `originalPlan`. `CometIcebergWriteExec` was
built from an `IcebergWriteExec`, so it could carry that node instead of its
child.
### Additional context
Found while reviewing #5696 (credit to @jordepic for spotting it there). It
is pre-existing and unrelated to that PR's transition fix: `countTransitions`
counts `ColumnarToRowTransition` nodes, and the `ColumnarToRowExec` #5696
causes Spark to insert below the write simply takes the place the write node
itself used to occupy in that count, so the revert decision is unchanged either
way.
Part of the native Iceberg writes epic #5649.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]