viirya commented on code in PR #6664:
URL: https://github.com/apache/datafusion-comet/pull/6664#discussion_r4224875529


##########
spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala:
##########
@@ -117,25 +120,114 @@ class CometIcebergWriteActionSuite
     }
   }
 
-  test("spark.comet.enabled=false keeps Spark's own write plan with the split 
flag on") {
+  // Turning off Comet, or only its native execution as an application that 
uses Comet just for
+  // scans or shuffle does, keeps Spark's own write operator.
+  Seq(CometConf.COMET_ENABLED, CometConf.COMET_EXEC_ENABLED).foreach { flag =>
+    test(s"${flag.key}=false keeps Spark's own write plan with the split flag 
on") {
+      assume(icebergAvailable, "Iceberg not available in classpath")
+      withIcebergCatalog { warehouseDir =>
+        val table = flag.key.replace('.', '_')
+        createTable(warehouseDir, table, partitionSpec = "")
+        // withSQLConf returns Unit before Spark 4.0, so the assertions run 
inside it.
+        withSQLConf(flag.key -> "false") {
+          val snapshot = captureWrite(table) {
+            spark.sql(s"INSERT INTO cat.db.$table VALUES " +
+              "(1, 'us-east', 10.5), (2, 'us-west', 20.3), (3, 'eu', 30.7)")
+          }
+          assert(
+            snapshot.snapshotDelta == 1L,
+            s"expected 1 commit, got ${snapshot.snapshotDelta}")
+          val (commits, writes) = collectIcebergWriteOps(snapshot.plans)
+          assert(
+            commits.isEmpty && writes.isEmpty,
+            s"expected Spark's own write plan with ${flag.key}=false. 
Plans:\n" +
+              snapshot.plans.mkString("\n--\n"))
+        }
+        assertRows(table, expectedIds = Seq(1, 2, 3))
+      }
+    }
+  }
+
+  // The suite pins both write flags, so this test unsets them to see what an 
application that
+  // sets neither gets. A Parquet scan is a native input, so the write is 
eligible.
+  // https://github.com/apache/datafusion-comet/issues/5644
+  test("an eligible Iceberg write runs natively under the split plan when no 
flag is set") {
     assume(icebergAvailable, "Iceberg not available in classpath")
     withIcebergCatalog { warehouseDir =>
-      createTable(warehouseDir, "comet_disabled", partitionSpec = "")
-      // withSQLConf returns Unit before Spark 4.0, so the assertions run 
inside it.
-      withSQLConf(CometConf.COMET_ENABLED.key -> "false") {
-        val snapshot = captureWrite("comet_disabled") {
-          spark.sql(
-            "INSERT INTO cat.db.comet_disabled VALUES " +
-              "(1, 'us-east', 10.5), (2, 'us-west', 20.3), (3, 'eu', 30.7)")
+      withTempPath { dir =>
+        spark
+          .range(10)
+          .selectExpr("CAST(id AS INT) AS id", "'eu' AS region", "CAST(id AS 
DOUBLE) AS amount")
+          .write
+          .parquet(dir.getCanonicalPath)
+        createTable(warehouseDir, "write_defaults", partitionSpec = "")
+        val snapshot = captureWrite("write_defaults") {
+          withSessionConf(
+            CometConf.COMET_ICEBERG_WRITE_SPLIT_OPERATOR_ENABLED.key -> None,
+            CometConf.COMET_ICEBERG_NATIVE_WRITE_ENABLED.key -> None) {
+            spark.read
+              .parquet(dir.getCanonicalPath)
+              .writeTo(s"$catalog.$ns.write_defaults")
+              .append()
+          }
         }
         assert(snapshot.snapshotDelta == 1L, s"expected 1 commit, got 
${snapshot.snapshotDelta}")
-        val (commits, writes) = collectIcebergWriteOps(snapshot.plans)
+        val (commits, _) = collectIcebergWriteOps(snapshot.plans)
+        val nativeWrites = snapshot.plans.flatMap { plan =>
+          collectWithSubqueries(plan) { case e: CometIcebergWriteExec => e }
+        }
         assert(
-          commits.isEmpty && writes.isEmpty,
-          "expected Spark's own write plan with Comet disabled. Plans:\n" +
+          commits.nonEmpty && nativeWrites.nonEmpty,
+          "expected an IcebergCommitExec over a CometIcebergWriteExec. 
Plans:\n" +
             snapshot.plans.mkString("\n--\n"))
+        assertRows("write_defaults", expectedIds = 0 until 10)
+      }
+    }
+  }
+
+  // With no flag set, an eligible write becomes a CometIcebergWriteExec, so 
reverting a
+  // transition-heavy stage has to put a JVM writer back rather than drop the 
write.
+  // https://github.com/apache/datafusion-comet/issues/5719
+  for (adaptive <- Seq(false, true)) {
+    test(s"transition-heavy fallback keeps the write when no flag is set with 
AQE=$adaptive") {
+      assume(icebergAvailable, "Iceberg not available in classpath")
+      withIcebergCatalog { warehouseDir =>
+        withTempPath { dir =>
+          spark
+            .range(3)
+            .selectExpr(
+              "CAST(id + 1 AS INT) AS id",
+              "'eu' AS region",
+              "CAST(id AS DOUBLE) AS amount")
+            .write
+            .parquet(dir.getCanonicalPath)
+          val table = s"transition_defaults_${if (adaptive) "aqe" else 
"no_aqe"}"
+          createTable(warehouseDir, table, partitionSpec = "")
+          // withSQLConf returns Unit before Spark 4.0, so the assertions run 
inside it.
+          withSQLConf(
+            SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> adaptive.toString,
+            CometConf.COMET_EXEC_TRANSITION_REVERT_ENABLED.key -> "true",
+            CometConf.COMET_EXEC_TRANSITION_REVERT_MAX_TRANSITIONS.key -> "0") 
{
+            val snapshot = captureWrite(table) {
+              withSessionConf(
+                CometConf.COMET_ICEBERG_WRITE_SPLIT_OPERATOR_ENABLED.key -> 
None,
+                CometConf.COMET_ICEBERG_NATIVE_WRITE_ENABLED.key -> None) {
+                
spark.read.parquet(dir.getCanonicalPath).writeTo(s"$catalog.$ns.$table").append()
+              }
+            }
+            assertExactlyOneCommit(snapshot)
+            // Also shows the stage was reverted: otherwise the native writer 
would still be here.
+            val nativeWrites = snapshot.plans.flatMap { plan =>
+              collectWithSubqueries(plan) { case e: CometIcebergWriteExec => e 
}
+            }
+            assert(
+              nativeWrites.isEmpty,

Review Comment:
   Thanks, this reads well. Running the same write under the default threshold 
first makes the reverted run meaningful, and checking for the JVM 
`IcebergWriteExec` too closes the gap. Resolving.



##########
docs/source/user-guide/latest/migration-guide.md:
##########
@@ -83,6 +83,31 @@ it. `spark.comet.convert.oneRowRelation.enabled` is on by 
default, so such a que
 without either setting, and logs no warning. The list remains the way to 
convert other leaf
 operators, such as the scan of a Data Source V2 connector.
 
+### Iceberg Writes
+
+`spark.comet.write.iceberg.enabled` now defaults to `true`. Comet plans an 
Iceberg `INSERT INTO`,
+`INSERT OVERWRITE`, and copy-on-write `DELETE`, `UPDATE` or `MERGE` as two 
operators, `IcebergWrite`
+under `IcebergCommit`, in place of Spark's single V2 write operator, and 
writes the data files of
+each eligible write natively with iceberg-rust. A write that is not eligible 
still uses
+iceberg-java's writer, and iceberg-java still commits every write. The table 
holds the same rows,
+but natively written files differ from iceberg-java's in the ways listed under
+[Accepted divergences](iceberg-writes.md#accepted-divergences), and explain 
output and the Spark UI
+show the new operators. Set `spark.comet.write.iceberg.enabled=false` to plan 
Spark's own operator,
+as in Comet 1.1.0. Comet 1.1.0 named this setting 
`spark.comet.iceberg.write.enabled`, in the
+testing category, and Comet now ignores that name, so a deployment that set it 
to `false` gets the
+new default unless it sets `spark.comet.write.iceberg.enabled=false`.
+`spark.comet.write.iceberg.splitOperator.enabled` is now a testing setting, 
and setting it to
+`false` does not turn the split operator off. See [Iceberg 
Writes](iceberg-writes.md).
+
+The native writer's buffers count against Comet's off-heap memory pool, where 
iceberg-java's buffers
+sit on the JVM heap. A fanout write keeps a data file open for every partition 
a task writes to, and
+each open file holds the row group it is writing, so a task that writes to 
many partitions can need
+more memory than the pool grants it. The task then fails with
+`Additional allocation failed for IcebergWriteExec`. Disabling the fanout 
writer

Review Comment:
   You're right about `useFanoutWriter`, and I had the scope too narrow. Since 
apache/iceberg#8621, `writeOrdering` drops the local sort whenever the table is 
unsorted and fanout is on by default, so any unsorted partitioned table gets 
the fanout writer whatever its distribution mode. Keeping fanout on 
iceberg-java would take most partitioned writes off the native path, so closing 
partitions early is the better fix. I'll review #6773 separately, and I'd like 
to keep this thread open until it lands, since it's listed as a blocker.



-- 
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]

Reply via email to