andygrove commented on code in PR #6664:
URL: https://github.com/apache/datafusion-comet/pull/6664#discussion_r4219518071
##########
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:
Done in 70bb11db2. The test now runs the same write twice. The first run
keeps the default threshold of two transitions, so its stage is not reverted,
and it has to plan `IcebergCommit` over `CometIcebergWriteExec`. The second
sets `maxTransitions=0` and has to plan `IcebergCommit` over a JVM
`IcebergWriteExec` with no `CometIcebergWriteExec` left. Each run has to commit
once, and the table ends with both runs' rows. With the Comet scan turned off
inside the test, which makes the write ineligible, the first run now fails.
##########
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:
I'd rather not keep fanout writes on iceberg-java. On Iceberg 1.5+ an
unsorted partitioned table gets the fanout writer by default whatever its
distribution mode, since `useFanoutWriter` is true whenever the write has no
required ordering, so that would take most partitioned writes off the native
path.
#6773 takes your second suggestion instead. When the pool refuses the
writer's reservation, it writes out and closes the partitions holding the most
memory, counting both a partition's open file and the rows it has not handed to
the file yet, until what is left fits. A closed partition's next rows open a
new file, so the write finishes with more, smaller files rather than failing,
and a `files closed early to free memory` metric on `CometIcebergWrite` counts
them. More files is the same kind of difference as the roll points already
listed under accepted divergences, and the PR adds it there.
The suite's 64-partition fanout test now interleaves its partitions and
passes under the same 4 MiB pool, writing 219 files where the full pool writes
64. A write that keeps one file open, unpartitioned or clustered, can still
fail, but only when that one file's row group outgrows the pool. I've added
#6773 to the blockers in the description.
--
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]