sunchao commented on code in PR #6664:
URL: https://github.com/apache/datafusion-comet/pull/6664#discussion_r4200024207
##########
spark/src/main/scala/org/apache/comet/CometConf.scala:
##########
@@ -125,21 +125,24 @@ object CometConf extends ShimCometConf {
conf("spark.comet.write.iceberg.splitOperator.enabled")
.category(CATEGORY_TESTING)
.doc(
- "Whether to rewrite Iceberg V2 writes from Spark's combined V2
write/commit operator " +
- "into Comet's two-operator shape: a file writer exec (inside AQE)
and a committer " +
- "(outside AQE).")
+ "Whether to plan Iceberg writes as Comet's file writer and committer
even when " +
+ "`spark.comet.write.iceberg.enabled` is false, so that Iceberg's own
writer writes " +
+ "every data file inside Comet's plan. Used by tests to compare the
two writers " +
+ "under the same plan.")
.booleanConf
.createWithDefault(false)
val COMET_ICEBERG_NATIVE_WRITE_ENABLED: ConfigEntry[Boolean] =
conf("spark.comet.write.iceberg.enabled")
- .category(CATEGORY_TESTING)
+ .category(CATEGORY_EXEC)
.doc(
- "Whether to delegate the executor-side Parquet write to Comet's native
(iceberg-rust) " +
- "writer when the table's properties allow it. Requires " +
- "`spark.comet.write.iceberg.splitOperator.enabled = true`. Off by
default.")
+ "Whether Comet plans Iceberg writes and writes their data files
natively. Comet " +
+ "replaces Spark's combined V2 write operator with a file writer
(inside AQE) under a " +
+ "committer (outside AQE), and writes the data files of each eligible
write with its " +
+ "native (iceberg-rust) writer. Other writes use Iceberg's own
writer, and Iceberg " +
+ "commits every write. Set this to false to plan Spark's own V2 write
operator.")
.booleanConf
- .createWithDefault(false)
+ .createWithDefault(true)
Review Comment:
[P2] Preserve the write operator during transition reversion before enabling
this default. With `spark.sql.adaptive.enabled=false`,
`spark.comet.exec.transitionRevert.enabled=true` and
`spark.comet.exec.transitionRevert.maxTransitions=0`, an eligible Iceberg
`INSERT ... SELECT` now selects `CometIcebergWriteExec` without either write
flag being set. Spark inserts a `ColumnarToRowExec` beneath it, which triggers
reversion before redundant transitions are removed. Because the native writer’s
`originalPlan` is its child, reversion removes the writer entirely.
`IcebergCommitExec` then attempts to deserialize table data as commit messages
and the write fails. Previously, these settings retained Spark’s working V2
writer. The underlying defect is tracked in #5719, but this default materially
expands its exposure to users who never enabled experimental Iceberg writes.
Restore an `IcebergWriteExec` when reverting, or safely exclude this
combination from native conversion, and add a re
gression covering the default flags.
Evidence: Bounded reproduction:
`/tmp/6664-current-review-075caeqj/WriteRevertProbe.scala`, with results in
`probe-full-pipeline.log`. Using exact-head classes, Spark’s transition
insertion and `CometRule.postColumnarRules` transformed `IcebergCommit ->
CometIcebergWrite -> ColumnarToRow -> CometProject -> CometLocalTableScan` into
`IcebergCommit -> Project -> Project -> LocalTableScan`. Output: `COUNT=1`,
`WRITERS_LEFT=0`, followed by `java.io.EOFException` for one ordinary binary
row. The equivalent stock `AppendDataExec` control reported `COMMIT_ROWS=1` and
`STOCK_SUCCEEDED=true`. Disabling reversion retained the native writer. This
was a JVM plan/commit harness with a test BatchWrite, without native storage
I/O.
--
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]