viirya commented on code in PR #6664:
URL: https://github.com/apache/datafusion-comet/pull/6664#discussion_r4212967359
##########
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:
This paragraph is clear about the failure mode, but it leaves the user to
discover it after a job fails. On Iceberg 1.5+, a partitioned table with
`write.distribution-mode=none` and no sort order uses the fanout writer by
default. That is a fairly common setup for people avoiding the shuffle. Under
this PR those writes go native, their buffers come out of the off-heap pool,
and the writer can't spill. A write that fit in the executor heap on 1.1.0 can
now fail with `Additional allocation failed for IcebergWriteExec`. Spark's
retry takes the same native path, so the job fails.
Would it be reasonable to keep fanout writes on iceberg-java by default for
now? Another option is for the native writer to close the largest open file
when the pool refuses a reservation, then retry. Roll points are already an
accepted divergence, so that wouldn't add a new kind of difference. If neither
fits, could you share some numbers on how many partitions a fanout task can
hold with a typical `spark.memory.offHeap.size`? Then we can judge whether the
default is safe for that shape.
##########
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:
The only evidence here that the stage was reverted is that no
`CometIcebergWriteExec` remains. That would also be true if this write stopped
being eligible, for example after a gate change or a change to how the Parquet
scan is planned. The test would then keep passing without exercising #5957.
Could the test first run the same write without `maxTransitions=0` and assert
that it plans a `CometIcebergWriteExec`? Then the reverted run would be
compared against a write we know goes native, and the regression test can't
pass vacuously.
--
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]