From cff9738460f089600ffe377e9a6ba8b949ed2151 Mon Sep 17 00:00:00 2001 From: Andy Grove Date: Mon, 5 Oct 2026 08:31:05 -0600 Subject: [PATCH 01/10] feat: enable the Iceberg split-operator write plan by default spark.comet.write.iceberg.splitOperator.enabled now defaults to true, so Comet plans an Iceberg append, overwrite, or copy-on-write DELETE, UPDATE or MERGE as IcebergWrite under IcebergCommit instead of Spark's single V2 write operator. iceberg-java still writes and commits the data files, and the native writer (spark.comet.write.iceberg.enabled) stays off. The setting moves from the testing config category to query execution, since it is now the switch that restores Spark's own operator. The user guide, the upgrade guide, the operators table, the Iceberg contributor guide and the Iceberg write review skill describe the new default. This is step 1 of #5644. --- .../review-comet-iceberg-write-pr/SKILL.md | 4 ++-- .../contributor-guide/iceberg-writes.md | 7 +++--- .../user-guide/latest/iceberg-writes.md | 20 ++++++++++------ docs/source/user-guide/latest/iceberg.md | 5 ++-- .../user-guide/latest/migration-guide.md | 9 ++++++++ docs/source/user-guide/latest/operators.md | 12 +++++----- .../scala/org/apache/comet/CometConf.scala | 8 ++++--- .../comet/CometIcebergWriteActionSuite.scala | 23 +++++++++++++++++++ 8 files changed, 65 insertions(+), 23 deletions(-) diff --git a/.ai/skills/review-comet-iceberg-write-pr/SKILL.md b/.ai/skills/review-comet-iceberg-write-pr/SKILL.md index 9752f865a5d..3e5ac8a3fea 100644 --- a/.ai/skills/review-comet-iceberg-write-pr/SKILL.md +++ b/.ai/skills/review-comet-iceberg-write-pr/SKILL.md @@ -49,8 +49,8 @@ not, inherits it. There is no query that fails to tell you. | Native writer | `spark.comet.write.iceberg.enabled` | `CometIcebergNativeWrite` (gate and serde), `IcebergWriteProtoTranslation`, `CometIcebergWriteExec`, `iceberg_write.rs` | | Shared with the scan | both | `iceberg_common.rs` (`load_file_io`, `storage_factory_for`, `scheme_of`), `IcebergReflection`, `NativeConfig` S3 translation | -A change to the split plan affects every Iceberg write once that flag is on, including writes that -never reach the native writer. A change to a shared file affects the native Iceberg scan too. +A change to the split plan affects every Iceberg write, since that flag is on by default, including +writes that never reach the native writer. A change to a shared file affects the native Iceberg scan too. ## 2. Direction of the Gate Change diff --git a/docs/source/contributor-guide/iceberg-writes.md b/docs/source/contributor-guide/iceberg-writes.md index 98901b474e4..f2cbefbd707 100644 --- a/docs/source/contributor-guide/iceberg-writes.md +++ b/docs/source/contributor-guide/iceberg-writes.md @@ -38,9 +38,10 @@ Two flags, each of which builds on the one before it: | `spark.comet.write.iceberg.enabled` | Who writes the data files. An eligible `IcebergWrite` becomes `CometIcebergWrite`, which writes Parquet with iceberg-rust. | The native flag does nothing without the split flag, because it converts a node only the split plan -creates. Both default to `false`. The roadmap for making them the default, and the criteria for it, -are tracked in [#5644](https://github.com/apache/datafusion-comet/issues/5644) under the epic -[#5649](https://github.com/apache/datafusion-comet/issues/5649). +creates. The split flag defaults to `true` since Comet 1.2.0, so a change to the split plan reaches +every Iceberg write. The native flag defaults to `false`. The roadmap for making it the default, and +the criteria for it, are tracked in [#5644](https://github.com/apache/datafusion-comet/issues/5644) +under the epic [#5649](https://github.com/apache/datafusion-comet/issues/5649). One rule runs through the whole native path: **the native writer must produce the outcome iceberg-java would have produced, or decline.** iceberg-java is the reference for every data file, diff --git a/docs/source/user-guide/latest/iceberg-writes.md b/docs/source/user-guide/latest/iceberg-writes.md index 0b340dbbe27..16f356ded22 100644 --- a/docs/source/user-guide/latest/iceberg-writes.md +++ b/docs/source/user-guide/latest/iceberg-writes.md @@ -17,10 +17,15 @@ under the License. --> -# Iceberg Writes: Comet's Split-Operator Plan (Experimental) +# Iceberg Writes: Comet's Split-Operator Plan -**This feature is experimental and disabled by default.** Enable it only after validating it -against your own workloads. +Since Comet 1.2.0, Comet plans Iceberg writes with the split-operator plan described below by +default. The plan changes only how Spark runs the write, and iceberg-java still writes the data +files. Set `spark.comet.write.iceberg.splitOperator.enabled=false` to plan Spark's own write +operator instead. + +**The native Parquet writer (`spark.comet.write.iceberg.enabled`) is experimental and disabled by +default.** Enable it only after validating it against your own workloads. ## Overview @@ -32,8 +37,8 @@ data-file writing cannot be re-planned in response to how its input ran. And bec writing is bundled with the metadata and commit steps, there is no separate step for Comet to replace. -When `spark.comet.write.iceberg.splitOperator.enabled=true`, Comet rewrites eligible Iceberg -writes into two operators: +When `spark.comet.write.iceberg.splitOperator.enabled=true`, the default, Comet rewrites eligible +Iceberg writes into two operators: 1. **`IcebergWrite`** — writes the data files on the executors, exactly as iceberg-java does today, and returns each task's serialized commit message. This operator and the sub-query @@ -79,7 +84,7 @@ spark.sql.catalog.=org.apache.iceberg.spark.SparkCatalog spark.sql.catalog..type=hadoop # or hive / glue / rest / ... spark.sql.catalog..warehouse=... -# Split-operator plan (experimental, off by default) +# Split-operator plan (on by default since Comet 1.2.0) spark.comet.write.iceberg.splitOperator.enabled=true # Native Parquet writer (experimental, off by default; requires the split plan) @@ -120,7 +125,8 @@ changes. The rewrite is skipped — and the write runs through Spark's stock combined operator — when: -- `spark.comet.write.iceberg.splitOperator.enabled` is `false` (the default); +- `spark.comet.write.iceberg.splitOperator.enabled` is set to `false`; +- Comet is disabled (`spark.comet.enabled=false`); - the write is not an Iceberg `SparkWrite` (any other V2 data source); - the table uses merge-on-read: delta writes (Iceberg `WriteDelta`) are not intercepted; - the statement is CTAS / RTAS on Spark 3.4, where the staged exec writes inline; on Spark diff --git a/docs/source/user-guide/latest/iceberg.md b/docs/source/user-guide/latest/iceberg.md index 4c3abe00216..457e9730507 100644 --- a/docs/source/user-guide/latest/iceberg.md +++ b/docs/source/user-guide/latest/iceberg.md @@ -197,8 +197,9 @@ The following scenarios will fall back to the JVM Iceberg reader: change. The check uses the table's schema history, so the fallback stays after those files are rewritten -Writes are not covered by this list. By default Iceberg writes use Spark's own writer; see -[Iceberg Writes](iceberg-writes.md) for the experimental native writer and when it applies. +Writes are not covered by this list. By default Comet plans an Iceberg write with its split-operator +plan, and iceberg-java still writes the data files; see [Iceberg Writes](iceberg-writes.md) for the +plan, the experimental native writer, and when each applies. ### Iceberg UDFs diff --git a/docs/source/user-guide/latest/migration-guide.md b/docs/source/user-guide/latest/migration-guide.md index 0a9731a2fa5..89247071b36 100644 --- a/docs/source/user-guide/latest/migration-guide.md +++ b/docs/source/user-guide/latest/migration-guide.md @@ -80,6 +80,15 @@ to use instead. Before Spark 4.1, Spark plans a query without a `FROM` clause as those versions `RDDScan` in the list also converted it. The list remains the way to convert other leaf operators, such as the scan of a Data Source V2 connector. +### Iceberg Write Plan + +`spark.comet.write.iceberg.splitOperator.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. iceberg-java +still writes the data files and commits them, so the written table is the same, but explain output +and the Spark UI show the two operators. Set `spark.comet.write.iceberg.splitOperator.enabled=false` +to plan Spark's own operator, as in Comet 1.1.0. See [Iceberg Writes](iceberg-writes.md). + ## Upgrading to Comet 1.1.0 Comet `1.1.0` makes no behavior changes that need a `spark.comet.legacy.*` key. The changes below diff --git a/docs/source/user-guide/latest/operators.md b/docs/source/user-guide/latest/operators.md index aea4fc8a884..6899c380aa6 100644 --- a/docs/source/user-guide/latest/operators.md +++ b/docs/source/user-guide/latest/operators.md @@ -131,12 +131,12 @@ natively on `BroadcastHashJoinExec` and `ShuffledHashJoinExec`. Existence sort-m ## Writes -| Operator | Status | Notes | -| ------------------------------------------------------------------------------------------------------------------ | ------ | ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | -| `WriteFilesExec` | ⚠️ | Spark 4.0+. Experimental native Parquet writes, disabled by default (opt-in). Non-partitioned, non-bucketed writes only, and not when `spark.sql.files.maxRecordsPerFile` is set. | -| `DataWritingCommandExec` | ⚠️ | Spark 3.4/3.5 only. Experimental native Parquet writes, disabled by default (opt-in). Replaced by `WriteFilesExec` on Spark 4.0+ and removed with Spark 3.x support. | -| Iceberg writes: `AppendDataExec`, `OverwriteByExpressionExec`, `OverwritePartitionsDynamicExec`, `ReplaceDataExec` | ⚠️ | Experimental, disabled by default. `spark.comet.write.iceberg.splitOperator.enabled=true` plans the write as `IcebergWrite` and `IcebergCommit`, and Iceberg's Java writer still writes the data files. `spark.comet.write.iceberg.enabled=true` then writes eligible data files natively (`CometIcebergWrite`). Covers `INSERT INTO`, `INSERT OVERWRITE`, and copy-on-write `DELETE` / `UPDATE` / `MERGE`. Merge-on-read writes (`WriteDeltaExec`) fall back. See [Iceberg Writes](iceberg-writes.md). | -| `MergeRowsExec` | ⚠️ | Spark 3.5+. Experimental, disabled by default (`spark.comet.exec.mergeRows.enabled`). On Spark 4.1+, stock V2 writers retain Spark MergeRowsExec; Comet's split Iceberg path can run it natively while preserving MergeSummary. See [MERGE INTO](compatibility/operators.md#merge-into-mergerowsexec). | +| Operator | Status | Notes | +| ------------------------------------------------------------------------------------------------------------------ | ------ | -------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | +| `WriteFilesExec` | ⚠️ | Spark 4.0+. Experimental native Parquet writes, disabled by default (opt-in). Non-partitioned, non-bucketed writes only, and not when `spark.sql.files.maxRecordsPerFile` is set. | +| `DataWritingCommandExec` | ⚠️ | Spark 3.4/3.5 only. Experimental native Parquet writes, disabled by default (opt-in). Replaced by `WriteFilesExec` on Spark 4.0+ and removed with Spark 3.x support. | +| Iceberg writes: `AppendDataExec`, `OverwriteByExpressionExec`, `OverwritePartitionsDynamicExec`, `ReplaceDataExec` | ⚠️ | Planned as `IcebergWrite` and `IcebergCommit` by default (`spark.comet.write.iceberg.splitOperator.enabled`), and Iceberg's Java writer still writes the data files. Experimental, disabled by default: `spark.comet.write.iceberg.enabled=true` writes eligible data files natively (`CometIcebergWrite`). Covers `INSERT INTO`, `INSERT OVERWRITE`, and copy-on-write `DELETE` / `UPDATE` / `MERGE`. Merge-on-read writes (`WriteDeltaExec`) fall back. See [Iceberg Writes](iceberg-writes.md). | +| `MergeRowsExec` | ⚠️ | Spark 3.5+. Experimental, disabled by default (`spark.comet.exec.mergeRows.enabled`). On Spark 4.1+, stock V2 writers retain Spark MergeRowsExec; Comet's split Iceberg path can run it natively while preserving MergeSummary. See [MERGE INTO](compatibility/operators.md#merge-into-mergerowsexec). | ## Python and UDF diff --git a/spark/src/main/scala/org/apache/comet/CometConf.scala b/spark/src/main/scala/org/apache/comet/CometConf.scala index e966a367e42..da750ca58ab 100644 --- a/spark/src/main/scala/org/apache/comet/CometConf.scala +++ b/spark/src/main/scala/org/apache/comet/CometConf.scala @@ -123,13 +123,15 @@ object CometConf extends ShimCometConf { val COMET_ICEBERG_WRITE_SPLIT_OPERATOR_ENABLED: ConfigEntry[Boolean] = conf("spark.comet.write.iceberg.splitOperator.enabled") - .category(CATEGORY_TESTING) + .category(CATEGORY_EXEC) .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).") + "(outside AQE). Iceberg's own writer still writes the data files unless " + + "`spark.comet.write.iceberg.enabled` is also set. Set this to false to plan " + + "Spark's own V2 write operator.") .booleanConf - .createWithDefault(false) + .createWithDefault(true) val COMET_ICEBERG_NATIVE_WRITE_ENABLED: ConfigEntry[Boolean] = conf("spark.comet.write.iceberg.enabled") diff --git a/spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala b/spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala index 7bf57481e5d..7234aad40ee 100644 --- a/spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala +++ b/spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala @@ -133,6 +133,29 @@ class CometIcebergWriteActionSuite } } + // The suite turns the split plan on explicitly, so this test drops that setting to see what an + // application that never sets it gets. + // https://github.com/apache/datafusion-comet/issues/5644 + test("an Iceberg write plans the split operator when the flag is not set") { + assume(icebergAvailable, "Iceberg not available in classpath") + withIcebergCatalog { warehouseDir => + createTable(warehouseDir, "split_default", partitionSpec = "") + val key = CometConf.COMET_ICEBERG_WRITE_SPLIT_OPERATOR_ENABLED.key + val conf = spark.sessionState.conf + conf.unsetConf(key) + try { + val snapshot = captureWrite("split_default") { + spark.sql( + "INSERT INTO cat.db.split_default VALUES (1, 'us-east', 10.5), (2, 'eu', 20.3)") + } + assertExactlyOneCommit(snapshot) + } finally { + conf.setConfString(key, "true") + } + assertRows("split_default", expectedIds = Seq(1, 2)) + } + } + test("AppendData partitioned INSERT INTO routes through two-op") { assume(icebergAvailable, "Iceberg not available in classpath") withIcebergCatalog { warehouseDir => From 6f56b09f98c695d44bbdfd458f62700e39f45ee7 Mon Sep 17 00:00:00 2001 From: Andy Grove Date: Mon, 5 Oct 2026 09:03:36 -0600 Subject: [PATCH 02/10] feat: write eligible Iceberg data files natively by default spark.comet.write.iceberg.enabled now defaults to true, so an eligible Iceberg write is written by iceberg-rust and every other write still falls back to iceberg-java. The setting moves to the query execution category, and the docs and the 1.2.0 upgrade-guide entry describe the new default. CometIcebergWriteActionSuite pins the flag off, so the tables it writes outside withNativeEnabled stay an iceberg-java baseline, and withNativeEnabled now restores the previous value instead of unsetting it, which would turn the native writer back on. --- .../review-comet-iceberg-write-pr/SKILL.md | 5 +- .../contributor-guide/iceberg-writes.md | 8 +-- .../user-guide/latest/iceberg-writes.md | 28 ++++---- docs/source/user-guide/latest/iceberg.md | 4 +- .../user-guide/latest/migration-guide.md | 21 +++--- docs/source/user-guide/latest/operators.md | 12 ++-- .../scala/org/apache/comet/CometConf.scala | 14 ++-- .../comet/CometIcebergWriteActionSuite.scala | 64 +++++++++++++++++-- 8 files changed, 108 insertions(+), 48 deletions(-) diff --git a/.ai/skills/review-comet-iceberg-write-pr/SKILL.md b/.ai/skills/review-comet-iceberg-write-pr/SKILL.md index 3e5ac8a3fea..c689a6e596b 100644 --- a/.ai/skills/review-comet-iceberg-write-pr/SKILL.md +++ b/.ai/skills/review-comet-iceberg-write-pr/SKILL.md @@ -49,8 +49,9 @@ not, inherits it. There is no query that fails to tell you. | Native writer | `spark.comet.write.iceberg.enabled` | `CometIcebergNativeWrite` (gate and serde), `IcebergWriteProtoTranslation`, `CometIcebergWriteExec`, `iceberg_write.rs` | | Shared with the scan | both | `iceberg_common.rs` (`load_file_io`, `storage_factory_for`, `scheme_of`), `IcebergReflection`, `NativeConfig` S3 translation | -A change to the split plan affects every Iceberg write, since that flag is on by default, including -writes that never reach the native writer. A change to a shared file affects the native Iceberg scan too. +Both flags are on by default, so a change to the split plan affects every Iceberg write, including +writes that never reach the native writer, and a change to the native writer affects every eligible +one. A change to a shared file affects the native Iceberg scan too. ## 2. Direction of the Gate Change diff --git a/docs/source/contributor-guide/iceberg-writes.md b/docs/source/contributor-guide/iceberg-writes.md index f2cbefbd707..59e12dac12a 100644 --- a/docs/source/contributor-guide/iceberg-writes.md +++ b/docs/source/contributor-guide/iceberg-writes.md @@ -38,10 +38,10 @@ Two flags, each of which builds on the one before it: | `spark.comet.write.iceberg.enabled` | Who writes the data files. An eligible `IcebergWrite` becomes `CometIcebergWrite`, which writes Parquet with iceberg-rust. | The native flag does nothing without the split flag, because it converts a node only the split plan -creates. The split flag defaults to `true` since Comet 1.2.0, so a change to the split plan reaches -every Iceberg write. The native flag defaults to `false`. The roadmap for making it the default, and -the criteria for it, are tracked in [#5644](https://github.com/apache/datafusion-comet/issues/5644) -under the epic [#5649](https://github.com/apache/datafusion-comet/issues/5649). +creates. Both default to `true` since Comet 1.2.0, so a change to the split plan reaches every +Iceberg write and a change to the native writer reaches every eligible one. The rollout and its +criteria are tracked in [#5644](https://github.com/apache/datafusion-comet/issues/5644) under the +epic [#5649](https://github.com/apache/datafusion-comet/issues/5649). One rule runs through the whole native path: **the native writer must produce the outcome iceberg-java would have produced, or decline.** iceberg-java is the reference for every data file, diff --git a/docs/source/user-guide/latest/iceberg-writes.md b/docs/source/user-guide/latest/iceberg-writes.md index 16f356ded22..5ee05aae8c5 100644 --- a/docs/source/user-guide/latest/iceberg-writes.md +++ b/docs/source/user-guide/latest/iceberg-writes.md @@ -17,15 +17,15 @@ under the License. --> -# Iceberg Writes: Comet's Split-Operator Plan +# Iceberg Writes: Comet's Split-Operator Plan and Native Writer -Since Comet 1.2.0, Comet plans Iceberg writes with the split-operator plan described below by -default. The plan changes only how Spark runs the write, and iceberg-java still writes the data -files. Set `spark.comet.write.iceberg.splitOperator.enabled=false` to plan Spark's own write -operator instead. +Since Comet 1.2.0, Comet plans Iceberg writes with the split-operator plan described below, and +writes the data files of each eligible write natively with iceberg-rust. A write that is not +eligible falls back to iceberg-java's writer. Natively written files differ from iceberg-java's in +the ways listed under [Accepted divergences](#accepted-divergences). -**The native Parquet writer (`spark.comet.write.iceberg.enabled`) is experimental and disabled by -default.** Enable it only after validating it against your own workloads. +Set `spark.comet.write.iceberg.enabled=false` to write every data file with iceberg-java, and +`spark.comet.write.iceberg.splitOperator.enabled=false` as well to plan Spark's own write operator. ## Overview @@ -49,7 +49,8 @@ Iceberg writes into two operators: With only the split plan enabled, data files are still written by iceberg-java; only the plan shape changes. The split moves data-file writing inside AQE and separates it from the commit, and it is the foundation for the second toggle: when -`spark.comet.write.iceberg.enabled=true` and the write passes the eligibility check below, the +`spark.comet.write.iceberg.enabled=true`, also the default, and the write passes the eligibility +check below, the `IcebergWrite` operator's per-task Parquet write is delegated to [iceberg-rust](https://github.com/apache/iceberg-rust) via Comet's native execution pipeline ([#5361](https://github.com/apache/datafusion-comet/pull/5361)). @@ -87,7 +88,7 @@ spark.sql.catalog..warehouse=... # Split-operator plan (on by default since Comet 1.2.0) spark.comet.write.iceberg.splitOperator.enabled=true -# Native Parquet writer (experimental, off by default; requires the split plan) +# Native Parquet writer (on by default since Comet 1.2.0; requires the split plan) spark.comet.write.iceberg.enabled=true # Lets writes whose input is a local relation (INSERT ... VALUES, a local DataFrame) use the @@ -141,7 +142,7 @@ trade-off, only no plan change. ## Native Parquet write eligibility -When `spark.comet.write.iceberg.enabled=true` +When `spark.comet.write.iceberg.enabled=true`, the default ([#5361](https://github.com/apache/datafusion-comet/pull/5361)), the `IcebergWrite` operator's per-task Parquet write is delegated to [iceberg-rust](https://github.com/apache/iceberg-rust). The native writer must produce the same outcome as iceberg-java — the same Parquet features, @@ -280,11 +281,12 @@ messages carry genuine `SparkWrite$TaskCommit` objects, so Iceberg's own `SparkW cleanup (which deletes the files listed in the commit messages for cleanable failures) applies unchanged. -## Accepted divergences behind the toggle +## Accepted divergences Some differences between parquet-mr and the pinned parquet-rs / iceberg-rust are unconditional — -they apply to every native write and cannot be configured away. Enabling -`spark.comet.write.iceberg.enabled` accepts them. They fall into three classes with very +they apply to every native write and cannot be configured away. Leaving +`spark.comet.write.iceberg.enabled` at its default of `true` accepts them, and setting it to +`false` avoids them. They fall into three classes with very different blast radius: differences confined to the physical bytes of a data file (cosmetic — no reader decision is based on them), differences visible in manifest metadata (these outlive the write and feed later readers' pruning decisions, so each one is analyzed individually diff --git a/docs/source/user-guide/latest/iceberg.md b/docs/source/user-guide/latest/iceberg.md index 457e9730507..9a247520f5b 100644 --- a/docs/source/user-guide/latest/iceberg.md +++ b/docs/source/user-guide/latest/iceberg.md @@ -198,8 +198,8 @@ The following scenarios will fall back to the JVM Iceberg reader: rewritten Writes are not covered by this list. By default Comet plans an Iceberg write with its split-operator -plan, and iceberg-java still writes the data files; see [Iceberg Writes](iceberg-writes.md) for the -plan, the experimental native writer, and when each applies. +plan and writes eligible data files natively; see [Iceberg Writes](iceberg-writes.md) for the plan, +the native writer, and when each applies. ### Iceberg UDFs diff --git a/docs/source/user-guide/latest/migration-guide.md b/docs/source/user-guide/latest/migration-guide.md index 89247071b36..75fd12076eb 100644 --- a/docs/source/user-guide/latest/migration-guide.md +++ b/docs/source/user-guide/latest/migration-guide.md @@ -80,14 +80,19 @@ to use instead. Before Spark 4.1, Spark plans a query without a `FROM` clause as those versions `RDDScan` in the list also converted it. The list remains the way to convert other leaf operators, such as the scan of a Data Source V2 connector. -### Iceberg Write Plan - -`spark.comet.write.iceberg.splitOperator.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. iceberg-java -still writes the data files and commits them, so the written table is the same, but explain output -and the Spark UI show the two operators. Set `spark.comet.write.iceberg.splitOperator.enabled=false` -to plan Spark's own operator, as in Comet 1.1.0. See [Iceberg Writes](iceberg-writes.md). +### Iceberg Writes + +`spark.comet.write.iceberg.splitOperator.enabled` and `spark.comet.write.iceberg.enabled` now default +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 write every data file with +iceberg-java, and `spark.comet.write.iceberg.splitOperator.enabled=false` as well to plan Spark's own +operator, as in Comet 1.1.0. See [Iceberg Writes](iceberg-writes.md). ## Upgrading to Comet 1.1.0 diff --git a/docs/source/user-guide/latest/operators.md b/docs/source/user-guide/latest/operators.md index 6899c380aa6..9e4ede449cb 100644 --- a/docs/source/user-guide/latest/operators.md +++ b/docs/source/user-guide/latest/operators.md @@ -131,12 +131,12 @@ natively on `BroadcastHashJoinExec` and `ShuffledHashJoinExec`. Existence sort-m ## Writes -| Operator | Status | Notes | -| ------------------------------------------------------------------------------------------------------------------ | ------ | -------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | -| `WriteFilesExec` | ⚠️ | Spark 4.0+. Experimental native Parquet writes, disabled by default (opt-in). Non-partitioned, non-bucketed writes only, and not when `spark.sql.files.maxRecordsPerFile` is set. | -| `DataWritingCommandExec` | ⚠️ | Spark 3.4/3.5 only. Experimental native Parquet writes, disabled by default (opt-in). Replaced by `WriteFilesExec` on Spark 4.0+ and removed with Spark 3.x support. | -| Iceberg writes: `AppendDataExec`, `OverwriteByExpressionExec`, `OverwritePartitionsDynamicExec`, `ReplaceDataExec` | ⚠️ | Planned as `IcebergWrite` and `IcebergCommit` by default (`spark.comet.write.iceberg.splitOperator.enabled`), and Iceberg's Java writer still writes the data files. Experimental, disabled by default: `spark.comet.write.iceberg.enabled=true` writes eligible data files natively (`CometIcebergWrite`). Covers `INSERT INTO`, `INSERT OVERWRITE`, and copy-on-write `DELETE` / `UPDATE` / `MERGE`. Merge-on-read writes (`WriteDeltaExec`) fall back. See [Iceberg Writes](iceberg-writes.md). | -| `MergeRowsExec` | ⚠️ | Spark 3.5+. Experimental, disabled by default (`spark.comet.exec.mergeRows.enabled`). On Spark 4.1+, stock V2 writers retain Spark MergeRowsExec; Comet's split Iceberg path can run it natively while preserving MergeSummary. See [MERGE INTO](compatibility/operators.md#merge-into-mergerowsexec). | +| Operator | Status | Notes | +| ------------------------------------------------------------------------------------------------------------------ | ------ | ---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | +| `WriteFilesExec` | ⚠️ | Spark 4.0+. Experimental native Parquet writes, disabled by default (opt-in). Non-partitioned, non-bucketed writes only, and not when `spark.sql.files.maxRecordsPerFile` is set. | +| `DataWritingCommandExec` | ⚠️ | Spark 3.4/3.5 only. Experimental native Parquet writes, disabled by default (opt-in). Replaced by `WriteFilesExec` on Spark 4.0+ and removed with Spark 3.x support. | +| Iceberg writes: `AppendDataExec`, `OverwriteByExpressionExec`, `OverwritePartitionsDynamicExec`, `ReplaceDataExec` | ⚠️ | Planned as `IcebergWrite` and `IcebergCommit` by default (`spark.comet.write.iceberg.splitOperator.enabled`), and eligible data files are written natively by default (`CometIcebergWrite`, `spark.comet.write.iceberg.enabled`). Other writes use Iceberg's Java writer. Covers `INSERT INTO`, `INSERT OVERWRITE`, and copy-on-write `DELETE` / `UPDATE` / `MERGE`. Merge-on-read writes (`WriteDeltaExec`) fall back. See [Iceberg Writes](iceberg-writes.md). | +| `MergeRowsExec` | ⚠️ | Spark 3.5+. Experimental, disabled by default (`spark.comet.exec.mergeRows.enabled`). On Spark 4.1+, stock V2 writers retain Spark MergeRowsExec; Comet's split Iceberg path can run it natively while preserving MergeSummary. See [MERGE INTO](compatibility/operators.md#merge-into-mergerowsexec). | ## Python and UDF diff --git a/spark/src/main/scala/org/apache/comet/CometConf.scala b/spark/src/main/scala/org/apache/comet/CometConf.scala index da750ca58ab..c5867493bbb 100644 --- a/spark/src/main/scala/org/apache/comet/CometConf.scala +++ b/spark/src/main/scala/org/apache/comet/CometConf.scala @@ -127,21 +127,23 @@ object CometConf extends ShimCometConf { .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). Iceberg's own writer still writes the data files unless " + - "`spark.comet.write.iceberg.enabled` is also set. Set this to false to plan " + - "Spark's own V2 write operator.") + "(outside AQE). The data files are written by Comet's native writer when " + + "`spark.comet.write.iceberg.enabled` allows it, and by Iceberg's own writer " + + "otherwise. Set this to false to plan Spark's own V2 write operator.") .booleanConf .createWithDefault(true) 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.") + "`spark.comet.write.iceberg.splitOperator.enabled = true`. A write the native " + + "writer cannot reproduce falls back to Iceberg's own writer. Set this to false to " + + "write every data file with Iceberg's own writer.") .booleanConf - .createWithDefault(false) + .createWithDefault(true) val COMET_ICEBERG_DATA_FILE_CONCURRENCY_LIMIT: ConfigEntry[Int] = conf("spark.comet.scan.icebergNative.dataFileConcurrencyLimit") diff --git a/spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala b/spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala index 7234aad40ee..a7a3c8ea55c 100644 --- a/spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala +++ b/spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala @@ -89,6 +89,9 @@ class CometIcebergWriteActionSuite override protected def sparkConf: SparkConf = { super.sparkConf .set(CometConf.COMET_ICEBERG_WRITE_SPLIT_OPERATOR_ENABLED.key, "true") + // A table written outside `withNativeEnabled` is the iceberg-java baseline that the native + // writer's output is compared against. + .set(CometConf.COMET_ICEBERG_NATIVE_WRITE_ENABLED.key, "false") // local[N,M] sets task max failures to M; the retry test needs one retry, and // spark.task.maxFailures does not override this part of a local master URL. .setMaster("local[5,2]") @@ -156,6 +159,47 @@ class CometIcebergWriteActionSuite } } + // The suite pins the native writer off for its iceberg-java baselines, so this test drops that + // setting to see what an application that never sets it 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 when the native flag is not set") { + assume(icebergAvailable, "Iceberg not available in classpath") + withIcebergCatalog { warehouseDir => + 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, "native_default", partitionSpec = "") + val key = CometConf.COMET_ICEBERG_NATIVE_WRITE_ENABLED.key + val conf = spark.sessionState.conf + conf.unsetConf(key) + try { + val snapshot = captureWrite("native_default") { + spark.read + .parquet(dir.getCanonicalPath) + .writeTo(s"$catalog.$ns.native_default") + .append() + } + assert( + snapshot.snapshotDelta == 1L, + s"expected 1 commit, got ${snapshot.snapshotDelta}") + val nativeWrites = snapshot.plans.flatMap { plan => + collectWithSubqueries(plan) { case e: CometIcebergWriteExec => e } + } + assert( + nativeWrites.nonEmpty, + "expected a CometIcebergWriteExec. Plans:\n" + snapshot.plans.mkString("\n--\n")) + } finally { + conf.setConfString(key, "false") + } + assertRows("native_default", expectedIds = 0 until 10) + } + } + } + test("AppendData partitioned INSERT INTO routes through two-op") { assume(icebergAvailable, "Iceberg not available in classpath") withIcebergCatalog { warehouseDir => @@ -3341,17 +3385,23 @@ class CometIcebergWriteActionSuite * (rather than `withSQLConf`) keeps the override visible to the columnar rule across some Spark * version / session-state combinations where `withSQLConf` loses the override before the rule * fires. + * + * Each setting is put back to the value it had before, not unset: unsetting the native write + * flag would turn it on, because it defaults to true while this suite pins it off. */ private def withNativeEnabled[T](action: => T): T = { - val session = spark - session.sessionState.conf - .setConfString(CometConf.COMET_ICEBERG_NATIVE_WRITE_ENABLED.key, "true") - session.sessionState.conf - .setConfString(CometConf.COMET_EXEC_LOCAL_TABLE_SCAN_ENABLED.key, "true") + val conf = spark.sessionState.conf + val keys = Seq( + CometConf.COMET_ICEBERG_NATIVE_WRITE_ENABLED.key, + CometConf.COMET_EXEC_LOCAL_TABLE_SCAN_ENABLED.key) + val previous = keys.map(key => key -> Option(conf.getConfString(key, null))) + keys.foreach(conf.setConfString(_, "true")) try action finally { - session.sessionState.conf.unsetConf(CometConf.COMET_EXEC_LOCAL_TABLE_SCAN_ENABLED.key) - session.sessionState.conf.unsetConf(CometConf.COMET_ICEBERG_NATIVE_WRITE_ENABLED.key) + previous.foreach { + case (key, Some(value)) => conf.setConfString(key, value) + case (key, None) => conf.unsetConf(key) + } } } From fd5cfb51bf13772506fee751dbe0c68bb2c5c6ca Mon Sep 17 00:00:00 2001 From: Andy Grove Date: Tue, 6 Oct 2026 12:05:26 -0600 Subject: [PATCH 03/10] docs: note the native Iceberg writer's off-heap memory use in the 1.2.0 upgrade guide With spark.comet.write.iceberg.enabled on by default, every eligible Iceberg write draws its buffers from Comet's off-heap memory pool (#6247) instead of the JVM heap. The 1.2.0 upgrade-guide entry now says so, names the error a task fails with when a fanout write to many partitions outgrows the pool, and lists what lets such a write fit. --- docs/source/user-guide/latest/migration-guide.md | 9 +++++++++ 1 file changed, 9 insertions(+) diff --git a/docs/source/user-guide/latest/migration-guide.md b/docs/source/user-guide/latest/migration-guide.md index 41c3dd02c9f..496c79c3411 100644 --- a/docs/source/user-guide/latest/migration-guide.md +++ b/docs/source/user-guide/latest/migration-guide.md @@ -97,6 +97,15 @@ show the new operators. Set `spark.comet.write.iceberg.enabled=false` to write e iceberg-java, and `spark.comet.write.iceberg.splitOperator.enabled=false` as well to plan Spark's own operator, as in Comet 1.1.0. 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 +(`write.spark.fanout.enabled=false`), a smaller `write.parquet.row-group-size-bytes` or a larger +`spark.memory.offHeap.size` lets such a write fit, and `spark.comet.write.iceberg.enabled=false` +writes it with iceberg-java as before. See [Failure handling](iceberg-writes.md#failure-handling). + ## Upgrading to Comet 1.1.0 Comet `1.1.0` makes no behavior changes that need a `spark.comet.legacy.*` key. The changes below From 1dc0d67df2ee8b15dc5b8704598dfb71e2d22831 Mon Sep 17 00:00:00 2001 From: Andy Grove Date: Tue, 6 Oct 2026 13:05:20 -0600 Subject: [PATCH 04/10] feat: make spark.comet.write.iceberg.enabled the only Iceberg write switch spark.comet.write.iceberg.enabled now plans the split operator by itself, so one setting turns Comet's Iceberg write path on or off: on plans IcebergCommit over IcebergWrite and writes eligible data files natively, off plans Spark's own V2 write operator. Separate flags for the two layers gave four combinations but only three behaviours: the native writer flag did nothing without the split flag, and the split plan with the native writer off gives a user nothing over Spark's own operator. spark.comet.write.iceberg.splitOperator.enabled goes back to the testing category, off by default. It plans the split operator with the native writer off, which CometIcebergWriteActionSuite uses for its iceberg-java baselines. The Iceberg Spark test diffs set both flags to true, so they run the same plan as before. --- .../review-comet-iceberg-write-pr/SKILL.md | 20 ++-- .../contributor-guide/iceberg-spark-tests.md | 12 +- .../contributor-guide/iceberg-writes.md | 24 ++-- .../user-guide/latest/iceberg-writes.md | 32 +++-- .../user-guide/latest/migration-guide.md | 19 ++- docs/source/user-guide/latest/operators.md | 12 +- .../latest/understanding-comet-plans.md | 4 +- .../scala/org/apache/comet/CometConf.scala | 23 ++-- .../comet/iceberg/IcebergWriteStrategy.scala | 6 +- .../CometIcebergRewriteActionSuite.scala | 3 +- .../CometIcebergSystemFunctionSuite.scala | 1 - .../comet/CometIcebergWriteActionSuite.scala | 109 +++++++----------- .../CometIcebergWriteDetectionSuite.scala | 6 +- .../CometIcebergWriteBenchmark.scala | 2 - 14 files changed, 125 insertions(+), 148 deletions(-) diff --git a/.ai/skills/review-comet-iceberg-write-pr/SKILL.md b/.ai/skills/review-comet-iceberg-write-pr/SKILL.md index 9334b0899ca..5b193b16e6b 100644 --- a/.ai/skills/review-comet-iceberg-write-pr/SKILL.md +++ b/.ai/skills/review-comet-iceberg-write-pr/SKILL.md @@ -43,15 +43,17 @@ not, inherits it. There is no query that fails to tell you. ## 1. Which Layer -| Layer | Flag | Code | -| -------------------- | ------------------------------------------------- | ---------------------------------------------------------------------------------------------------------------------------- | -| Split-operator plan | `spark.comet.write.iceberg.splitOperator.enabled` | `IcebergWriteStrategy`, `IcebergWriteLogical`, `IcebergWriteExec`, `IcebergCommitExec`, the `spark-*/.../iceberg/` shims | -| Native writer | `spark.comet.write.iceberg.enabled` | `CometIcebergNativeWrite` (gate and serde), `IcebergWriteProtoTranslation`, `CometIcebergWriteExec`, `iceberg_write.rs` | -| Shared with the scan | both | `iceberg_common.rs` (`load_file_io`, `storage_factory_for`, `scheme_of`), `IcebergReflection`, `NativeConfig` S3 translation | - -Both flags are on by default, so a change to the split plan affects every Iceberg write, including -writes that never reach the native writer, and a change to the native writer affects every eligible -one. A change to a shared file affects the native Iceberg scan too. +| Layer | Code | +| -------------------- | ---------------------------------------------------------------------------------------------------------------------------- | +| Split-operator plan | `IcebergWriteStrategy`, `IcebergWriteLogical`, `IcebergWriteExec`, `IcebergCommitExec`, the `spark-*/.../iceberg/` shims | +| Native writer | `CometIcebergNativeWrite` (gate and serde), `IcebergWriteProtoTranslation`, `CometIcebergWriteExec`, `iceberg_write.rs` | +| Shared with the scan | `iceberg_common.rs` (`load_file_io`, `storage_factory_for`, `scheme_of`), `IcebergReflection`, `NativeConfig` S3 translation | + +`spark.comet.write.iceberg.enabled` switches both write layers and is on by default, so a change to +the split plan affects every Iceberg write, including writes that never reach the native writer, and +a change to the native writer affects every eligible one. The testing-only +`spark.comet.write.iceberg.splitOperator.enabled` plans the split operator with the native writer +off. A change to a shared file affects the native Iceberg scan too. ## 2. Direction of the Gate Change diff --git a/docs/source/contributor-guide/iceberg-spark-tests.md b/docs/source/contributor-guide/iceberg-spark-tests.md index 87ef47427d2..812c53fca50 100644 --- a/docs/source/contributor-guide/iceberg-spark-tests.md +++ b/docs/source/contributor-guide/iceberg-spark-tests.md @@ -31,13 +31,11 @@ Here is an overview of the changes that the diffs make to Iceberg: uses a native Iceberg scan, these classes fail to compile and must be removed. - Configure test base classes (`TestBase`, `ExtensionsTestBase`, `ScanTestBase`, etc.) to load the Comet Spark plugin and shuffle manager -- Enable the Iceberg write split-operator plan (`spark.comet.write.iceberg.splitOperator.enabled`) alongside the - native scan in every Comet-configured session. The flag is off by default for users, so Iceberg's own suites - are the only place the split plan (`IcebergCommit -> IcebergWrite`) is exercised against Iceberg's write, - commit, and row-level-operation tests. See [#5259] -- Enable Comet's native (iceberg-rust) Parquet writer (`spark.comet.write.iceberg.enabled`) in the same sessions. - The native writer is experimental and off by default for users, so this is where it runs against Iceberg's - write, commit, and row-level-operation tests. +- Enable Comet's Iceberg write path, the split-operator plan (`IcebergCommit -> IcebergWrite`) and the native + (iceberg-rust) Parquet writer, alongside the native scan in every Comet-configured session, so that Iceberg's + write, commit, and row-level-operation tests run through it. `spark.comet.write.iceberg.enabled` switches both + and is on by default; the diffs set it explicitly, together with the testing-only + `spark.comet.write.iceberg.splitOperator.enabled`, which it makes redundant. See [#5259] - Enable `spark.comet.exec.localTableScan.enabled` in the same sessions. `CometIcebergNativeWrite` sets `requiresNativeChildren`, so without this flag a write fed by an inline `VALUES` list keeps Spark's row-based `LocalTableScanExec`, the conversion is declined, and the write silently runs on the JVM writer. Many Iceberg diff --git a/docs/source/contributor-guide/iceberg-writes.md b/docs/source/contributor-guide/iceberg-writes.md index d22e486d695..424bf1f4495 100644 --- a/docs/source/contributor-guide/iceberg-writes.md +++ b/docs/source/contributor-guide/iceberg-writes.md @@ -30,18 +30,20 @@ does not repeat those lists; it explains the code that implements them. ## Overview -Two flags, each of which builds on the one before it: +Two layers, the second built on the first, both switched by `spark.comet.write.iceberg.enabled`: -| Flag | What it changes | -| ------------------------------------------------- | --------------------------------------------------------------------------------------------------------------------------------------- | -| `spark.comet.write.iceberg.splitOperator.enabled` | The plan shape. Spark's single V2 write operator becomes `IcebergCommit` over `IcebergWrite`. iceberg-java still writes the data files. | -| `spark.comet.write.iceberg.enabled` | Who writes the data files. An eligible `IcebergWrite` becomes `CometIcebergWrite`, which writes Parquet with iceberg-rust. | +| Layer | What it changes | +| ----------------------- | ----------------------------------------------------------------------------------------------------------------------------------------------- | +| The split-operator plan | The plan shape. Spark's single V2 write operator becomes `IcebergCommit` over `IcebergWrite`, and iceberg-java writes the data files inside it. | +| The native writer | Who writes the data files. An eligible `IcebergWrite` becomes `CometIcebergWrite`, which writes Parquet with iceberg-rust. | -The native flag does nothing without the split flag, because it converts a node only the split plan -creates. Both default to `true` since Comet 1.2.0, so a change to the split plan reaches every -Iceberg write and a change to the native writer reaches every eligible one. The rollout and its -criteria are tracked in [#5644](https://github.com/apache/datafusion-comet/issues/5644) under the -epic [#5649](https://github.com/apache/datafusion-comet/issues/5649). +The flag defaults to `true` since Comet 1.2.0, so a change to the split plan reaches every Iceberg +write and a change to the native writer reaches every eligible one. Set to `false`, it plans Spark's +own write operator. The testing-only `spark.comet.write.iceberg.splitOperator.enabled` plans the +split operator with the native writer off, which the suites use to compare the two writers under +the same plan. The rollout and its criteria are tracked in +[#5644](https://github.com/apache/datafusion-comet/issues/5644) under the epic +[#5649](https://github.com/apache/datafusion-comet/issues/5649). One rule runs through the whole native path: **the native writer must produce the outcome iceberg-java would have produced, or decline.** iceberg-java is the reference for every data file, @@ -442,7 +444,7 @@ When writing a native-write test: - **Check storage state for failure tests**, not only the table: list the files under the data location and compare them with what the manifests reference. -The upstream Iceberg Spark tests also run with both flags and `localTableScan` enabled (see +The upstream Iceberg Spark tests also run with the native writer and `localTableScan` enabled (see [Running Iceberg Spark Tests](iceberg-spark-tests.md)). They are a broad regression net, but they do not assert which writer ran, and Comet's fallback reasons do not appear in their CI logs, so a green run is not evidence that the native writer handled a given test diff --git a/docs/source/user-guide/latest/iceberg-writes.md b/docs/source/user-guide/latest/iceberg-writes.md index cd279d38d8d..6b9d520487f 100644 --- a/docs/source/user-guide/latest/iceberg-writes.md +++ b/docs/source/user-guide/latest/iceberg-writes.md @@ -24,8 +24,8 @@ writes the data files of each eligible write natively with iceberg-rust. A write eligible falls back to iceberg-java's writer. Natively written files differ from iceberg-java's in the ways listed under [Accepted divergences](#accepted-divergences). -Set `spark.comet.write.iceberg.enabled=false` to write every data file with iceberg-java, and -`spark.comet.write.iceberg.splitOperator.enabled=false` as well to plan Spark's own write operator. +`spark.comet.write.iceberg.enabled` switches both. Set it to `false` to plan Spark's own write +operator instead, as Comet 1.1.0 did, so that iceberg-java writes every data file. ## Overview @@ -37,8 +37,8 @@ data-file writing cannot be re-planned in response to how its input ran. And bec writing is bundled with the metadata and commit steps, there is no separate step for Comet to replace. -When `spark.comet.write.iceberg.splitOperator.enabled=true`, the default, Comet rewrites eligible -Iceberg writes into two operators: +When `spark.comet.write.iceberg.enabled=true`, the default, Comet rewrites eligible Iceberg writes +into two operators: 1. **`IcebergWrite`** — writes the data files on the executors, exactly as iceberg-java does today, and returns each task's serialized commit message. This operator and the sub-query @@ -46,14 +46,12 @@ Iceberg writes into two operators: 2. **`IcebergCommit`** — collects the commit messages on the driver and performs the normal Iceberg commit (including commit-time validation), outside AQE, exactly once. -With only the split plan enabled, data files are still written by iceberg-java; only the plan -shape changes. The split moves data-file writing inside AQE and separates it from the commit, -and it is the foundation for the second toggle: when -`spark.comet.write.iceberg.enabled=true`, also the default, and the write passes the eligibility -check below, the -`IcebergWrite` operator's per-task Parquet write is delegated to +The split moves data-file writing inside AQE and separates it from the commit, which gives Comet +a step to replace: when the write passes the eligibility check below, the `IcebergWrite` +operator's per-task Parquet write is delegated to [iceberg-rust](https://github.com/apache/iceberg-rust) via Comet's native execution pipeline -([#5361](https://github.com/apache/datafusion-comet/pull/5361)). +([#5361](https://github.com/apache/datafusion-comet/pull/5361)). A write that does not pass keeps +`IcebergWrite`, and iceberg-java writes its data files; only the plan shape changes. ## How the native write works @@ -98,10 +96,8 @@ spark.sql.catalog.=org.apache.iceberg.spark.SparkCatalog spark.sql.catalog..type=hadoop # or hive / glue / rest / ... spark.sql.catalog..warehouse=... -# Split-operator plan (on by default since Comet 1.2.0) -spark.comet.write.iceberg.splitOperator.enabled=true - -# Native Parquet writer (on by default since Comet 1.2.0; requires the split plan) +# Split-operator plan and native Parquet writer (on by default since Comet 1.2.0; false plans +# Spark's own write operator) spark.comet.write.iceberg.enabled=true # Lets writes whose input is a local relation (INSERT ... VALUES, a local DataFrame) use the @@ -139,7 +135,7 @@ changes. The rewrite is skipped — and the write runs through Spark's stock combined operator — when: -- `spark.comet.write.iceberg.splitOperator.enabled` is set to `false`; +- `spark.comet.write.iceberg.enabled` is set to `false`; - Comet is disabled (`spark.comet.enabled=false`); - the write is not an Iceberg `SparkWrite` (any other V2 data source); - the table uses merge-on-read: delta writes (Iceberg `WriteDelta`) are not intercepted; @@ -173,8 +169,8 @@ The native writer reads its input as Arrow batches from a Comet operator, so the must itself run in Comet. A write whose input is a local relation, such as `INSERT ... VALUES` or `df.writeTo(...).append()` on a DataFrame built from local data, is fed by Spark's `LocalTableScanExec`, which Comet only converts when `spark.comet.exec.localTableScan.enabled=true` -(off by default). Without that setting such writes run through iceberg-java even when both write -flags are on. +(off by default). Without that setting such writes keep `IcebergWrite` and run through +iceberg-java. **Most Iceberg write settings are not supported.** Detection is an allowlist: a write is eligible only when its entire effective configuration matches the table below, and anything diff --git a/docs/source/user-guide/latest/migration-guide.md b/docs/source/user-guide/latest/migration-guide.md index 496c79c3411..2b5479d3e1f 100644 --- a/docs/source/user-guide/latest/migration-guide.md +++ b/docs/source/user-guide/latest/migration-guide.md @@ -85,17 +85,16 @@ operators, such as the scan of a Data Source V2 connector. ### Iceberg Writes -`spark.comet.write.iceberg.splitOperator.enabled` and `spark.comet.write.iceberg.enabled` now default -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 +`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 write every data file with -iceberg-java, and `spark.comet.write.iceberg.splitOperator.enabled=false` as well to plan Spark's own -operator, as in Comet 1.1.0. See [Iceberg Writes](iceberg-writes.md). +show the new operators. Set `spark.comet.write.iceberg.enabled=false` to plan Spark's own operator, +as in Comet 1.1.0. `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 diff --git a/docs/source/user-guide/latest/operators.md b/docs/source/user-guide/latest/operators.md index f2aa2ec84f7..ed5f60c9c17 100644 --- a/docs/source/user-guide/latest/operators.md +++ b/docs/source/user-guide/latest/operators.md @@ -132,12 +132,12 @@ natively on `BroadcastHashJoinExec` and `ShuffledHashJoinExec`. Existence sort-m ## Writes -| Operator | Status | Notes | -| ------------------------------------------------------------------------------------------------------------------ | ------ | ---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | -| `WriteFilesExec` | ⚠️ | Spark 4.0+. Experimental native Parquet writes, disabled by default (opt-in). Non-partitioned, non-bucketed writes only, and not when `spark.sql.files.maxRecordsPerFile` is set. | -| `DataWritingCommandExec` | ⚠️ | Spark 3.4/3.5 only. Experimental native Parquet writes, disabled by default (opt-in). Replaced by `WriteFilesExec` on Spark 4.0+ and removed with Spark 3.x support. | -| Iceberg writes: `AppendDataExec`, `OverwriteByExpressionExec`, `OverwritePartitionsDynamicExec`, `ReplaceDataExec` | ⚠️ | Planned as `IcebergWrite` and `IcebergCommit` by default (`spark.comet.write.iceberg.splitOperator.enabled`), and eligible data files are written natively by default (`CometIcebergWrite`, `spark.comet.write.iceberg.enabled`). Other writes use Iceberg's Java writer. Covers `INSERT INTO`, `INSERT OVERWRITE`, and copy-on-write `DELETE` / `UPDATE` / `MERGE`. Merge-on-read writes (`WriteDeltaExec`) fall back. See [Iceberg Writes](iceberg-writes.md). | -| `MergeRowsExec` | ⚠️ | Spark 3.5+. Experimental, disabled by default (`spark.comet.exec.mergeRows.enabled`). On Spark 4.1+, stock V2 writers retain Spark MergeRowsExec; Comet's split Iceberg path can run it natively while preserving MergeSummary. See [MERGE INTO](compatibility/operators.md#merge-into-mergerowsexec). | +| Operator | Status | Notes | +| ------------------------------------------------------------------------------------------------------------------ | ------ | -------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | +| `WriteFilesExec` | ⚠️ | Spark 4.0+. Experimental native Parquet writes, disabled by default (opt-in). Non-partitioned, non-bucketed writes only, and not when `spark.sql.files.maxRecordsPerFile` is set. | +| `DataWritingCommandExec` | ⚠️ | Spark 3.4/3.5 only. Experimental native Parquet writes, disabled by default (opt-in). Replaced by `WriteFilesExec` on Spark 4.0+ and removed with Spark 3.x support. | +| Iceberg writes: `AppendDataExec`, `OverwriteByExpressionExec`, `OverwritePartitionsDynamicExec`, `ReplaceDataExec` | ⚠️ | Planned as `IcebergWrite` and `IcebergCommit` by default (`spark.comet.write.iceberg.enabled`), and eligible data files are written natively (`CometIcebergWrite`). Other writes use Iceberg's Java writer. Covers `INSERT INTO`, `INSERT OVERWRITE`, and copy-on-write `DELETE` / `UPDATE` / `MERGE`. Merge-on-read writes (`WriteDeltaExec`) fall back. See [Iceberg Writes](iceberg-writes.md). | +| `MergeRowsExec` | ⚠️ | Spark 3.5+. Experimental, disabled by default (`spark.comet.exec.mergeRows.enabled`). On Spark 4.1+, stock V2 writers retain Spark MergeRowsExec; Comet's split Iceberg path can run it natively while preserving MergeSummary. See [MERGE INTO](compatibility/operators.md#merge-into-mergerowsexec). | ## Python and UDF diff --git a/docs/source/user-guide/latest/understanding-comet-plans.md b/docs/source/user-guide/latest/understanding-comet-plans.md index 13733f0b510..655865ff250 100644 --- a/docs/source/user-guide/latest/understanding-comet-plans.md +++ b/docs/source/user-guide/latest/understanding-comet-plans.md @@ -222,8 +222,8 @@ Keep the following in mind when reading the reports: also counts its subqueries, so do not add the reports together. - Only the JVM side of planning runs. Anything that would fail when DataFusion builds the native plan still counts as accelerated, so treat the percentage as an upper bound. -- Comet's split Iceberg V2 write (`spark.comet.write.iceberg.splitOperator.enabled`) is - declined in plan-only mode, so such writes run on, and are reported as, Spark. +- Comet's split Iceberg V2 write (`spark.comet.write.iceberg.enabled`) is declined in + plan-only mode, so such writes run on, and are reported as, Spark. - Under AQE the report describes the plan before any adaptive re-planning, so coverage of the plan that finally executes can differ. In particular, AQE plans subqueries into the outer query only after the report is produced, so the outer report counts their operators as diff --git a/spark/src/main/scala/org/apache/comet/CometConf.scala b/spark/src/main/scala/org/apache/comet/CometConf.scala index c814971e1fe..a2a4725c8d3 100644 --- a/spark/src/main/scala/org/apache/comet/CometConf.scala +++ b/spark/src/main/scala/org/apache/comet/CometConf.scala @@ -123,25 +123,24 @@ object CometConf extends ShimCometConf { val COMET_ICEBERG_WRITE_SPLIT_OPERATOR_ENABLED: ConfigEntry[Boolean] = conf("spark.comet.write.iceberg.splitOperator.enabled") - .category(CATEGORY_EXEC) + .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). The data files are written by Comet's native writer when " + - "`spark.comet.write.iceberg.enabled` allows it, and by Iceberg's own writer " + - "otherwise. Set this to false to plan Spark's own V2 write operator.") + "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(true) + .createWithDefault(false) val COMET_ICEBERG_NATIVE_WRITE_ENABLED: ConfigEntry[Boolean] = conf("spark.comet.write.iceberg.enabled") .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`. A write the native " + - "writer cannot reproduce falls back to Iceberg's own writer. Set this to false to " + - "write every data file with Iceberg's own writer.") + "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(true) diff --git a/spark/src/main/scala/org/apache/comet/iceberg/IcebergWriteStrategy.scala b/spark/src/main/scala/org/apache/comet/iceberg/IcebergWriteStrategy.scala index 85fbcc8ce3a..057c803ad57 100644 --- a/spark/src/main/scala/org/apache/comet/iceberg/IcebergWriteStrategy.scala +++ b/spark/src/main/scala/org/apache/comet/iceberg/IcebergWriteStrategy.scala @@ -38,9 +38,13 @@ case class IcebergWriteStrategy(session: SparkSession) extends SparkStrategy { override def apply(plan: LogicalPlan): Seq[SparkPlan] = { val conf = session.sessionState.conf + // The native write flag plans the split operator on its own. The split flag plans it with + // the native writer off, which only tests do. + val splitEnabled = CometConf.COMET_ICEBERG_NATIVE_WRITE_ENABLED.get(conf) || + CometConf.COMET_ICEBERG_WRITE_SPLIT_OPERATOR_ENABLED.get(conf) // Planner strategies run whether or not Comet is enabled, so check it here too: with Comet // off, Spark must plan its own V2 write operator. - if (!isCometLoaded(conf) || !CometConf.COMET_ICEBERG_WRITE_SPLIT_OPERATOR_ENABLED.get(conf)) { + if (!isCometLoaded(conf) || !splitEnabled) { return Nil } // Planner strategies run before CometRule, so plan-only mode needs its own guard here. diff --git a/spark/src/test/scala/org/apache/comet/CometIcebergRewriteActionSuite.scala b/spark/src/test/scala/org/apache/comet/CometIcebergRewriteActionSuite.scala index 506e9370fc7..5def10d762c 100644 --- a/spark/src/test/scala/org/apache/comet/CometIcebergRewriteActionSuite.scala +++ b/spark/src/test/scala/org/apache/comet/CometIcebergRewriteActionSuite.scala @@ -396,7 +396,7 @@ class CometIcebergRewriteActionSuite extends CometTestBase with CometIcebergTest // -- Plan assertions ------------------------------------------------------- // The per-group rewrite write plans as Spark's AppendData, or as Comet's IcebergCommit when - // COMET_ICEBERG_WRITE_SPLIT_OPERATOR_ENABLED is on. + // COMET_ICEBERG_NATIVE_WRITE_ENABLED is on. private def isRewriteWrite(plan: CapturedPlan): Boolean = plan.hasNode("AppendData") || plan.hasNode("IcebergCommit") @@ -524,7 +524,6 @@ class CometIcebergRewriteActionSuite extends CometTestBase with CometIcebergTest CometConf.COMET_ENABLED.key -> "true", CometConf.COMET_EXEC_ENABLED.key -> "true", CometConf.COMET_ICEBERG_NATIVE_ENABLED.key -> "true", - CometConf.COMET_ICEBERG_WRITE_SPLIT_OPERATOR_ENABLED.key -> "true", CometConf.COMET_ICEBERG_NATIVE_WRITE_ENABLED.key -> "true", CometConf.COMET_SCALA_UDF_CODEGEN_ENABLED.key -> "true")(body) diff --git a/spark/src/test/scala/org/apache/comet/CometIcebergSystemFunctionSuite.scala b/spark/src/test/scala/org/apache/comet/CometIcebergSystemFunctionSuite.scala index 5db20ba52f6..613f726767d 100644 --- a/spark/src/test/scala/org/apache/comet/CometIcebergSystemFunctionSuite.scala +++ b/spark/src/test/scala/org/apache/comet/CometIcebergSystemFunctionSuite.scala @@ -62,7 +62,6 @@ class CometIcebergSystemFunctionSuite override protected def sparkConf: SparkConf = { super.sparkConf - .set(CometConf.COMET_ICEBERG_WRITE_SPLIT_OPERATOR_ENABLED.key, "true") .set(CometConf.COMET_ICEBERG_NATIVE_WRITE_ENABLED.key, "true") } diff --git a/spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala b/spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala index a52c236a49d..39b5a07b911 100644 --- a/spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala +++ b/spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala @@ -93,9 +93,9 @@ class CometIcebergWriteActionSuite override protected def sparkConf: SparkConf = { super.sparkConf + // The split plan with the native writer off: a table written outside `withNativeEnabled` is + // the iceberg-java baseline that the native writer's output is compared against. .set(CometConf.COMET_ICEBERG_WRITE_SPLIT_OPERATOR_ENABLED.key, "true") - // A table written outside `withNativeEnabled` is the iceberg-java baseline that the native - // writer's output is compared against. .set(CometConf.COMET_ICEBERG_NATIVE_WRITE_ENABLED.key, "false") // local[N,M] sets task max failures to M; the retry test needs one retry, and // spark.task.maxFailures does not override this part of a local master URL. @@ -141,34 +141,10 @@ class CometIcebergWriteActionSuite } } - // The suite turns the split plan on explicitly, so this test drops that setting to see what an - // application that never sets it gets. + // 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 Iceberg write plans the split operator when the flag is not set") { - assume(icebergAvailable, "Iceberg not available in classpath") - withIcebergCatalog { warehouseDir => - createTable(warehouseDir, "split_default", partitionSpec = "") - val key = CometConf.COMET_ICEBERG_WRITE_SPLIT_OPERATOR_ENABLED.key - val conf = spark.sessionState.conf - conf.unsetConf(key) - try { - val snapshot = captureWrite("split_default") { - spark.sql( - "INSERT INTO cat.db.split_default VALUES (1, 'us-east', 10.5), (2, 'eu', 20.3)") - } - assertExactlyOneCommit(snapshot) - } finally { - conf.setConfString(key, "true") - } - assertRows("split_default", expectedIds = Seq(1, 2)) - } - } - - // The suite pins the native writer off for its iceberg-java baselines, so this test drops that - // setting to see what an application that never sets it 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 when the native flag is not set") { + 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 => withTempPath { dir => @@ -177,30 +153,27 @@ class CometIcebergWriteActionSuite .selectExpr("CAST(id AS INT) AS id", "'eu' AS region", "CAST(id AS DOUBLE) AS amount") .write .parquet(dir.getCanonicalPath) - createTable(warehouseDir, "native_default", partitionSpec = "") - val key = CometConf.COMET_ICEBERG_NATIVE_WRITE_ENABLED.key - val conf = spark.sessionState.conf - conf.unsetConf(key) - try { - val snapshot = captureWrite("native_default") { + 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.native_default") + .writeTo(s"$catalog.$ns.write_defaults") .append() } - assert( - snapshot.snapshotDelta == 1L, - s"expected 1 commit, got ${snapshot.snapshotDelta}") - val nativeWrites = snapshot.plans.flatMap { plan => - collectWithSubqueries(plan) { case e: CometIcebergWriteExec => e } - } - assert( - nativeWrites.nonEmpty, - "expected a CometIcebergWriteExec. Plans:\n" + snapshot.plans.mkString("\n--\n")) - } finally { - conf.setConfString(key, "false") } - assertRows("native_default", expectedIds = 0 until 10) + assert(snapshot.snapshotDelta == 1L, s"expected 1 commit, got ${snapshot.snapshotDelta}") + val (commits, _) = collectIcebergWriteOps(snapshot.plans) + val nativeWrites = snapshot.plans.flatMap { plan => + collectWithSubqueries(plan) { case e: CometIcebergWriteExec => e } + } + assert( + commits.nonEmpty && nativeWrites.nonEmpty, + "expected an IcebergCommitExec over a CometIcebergWriteExec. Plans:\n" + + snapshot.plans.mkString("\n--\n")) + assertRows("write_defaults", expectedIds = 0 until 10) } } } @@ -828,13 +801,15 @@ class CometIcebergWriteActionSuite } } - test("disabled config falls through to Spark's V2ExistingTableWriteExec") { + // The split flag is a testing setting, so an application turns Comet's Iceberg writes off with + // spark.comet.write.iceberg.enabled alone. The suite already pins that flag off. + test("spark.comet.write.iceberg.enabled=false plans Spark's V2ExistingTableWriteExec") { assume(icebergAvailable, "Iceberg not available in classpath") withIcebergCatalog { warehouseDir => createTable(warehouseDir, "disabled_conf", partitionSpec = "") val snapshot = captureWrite("disabled_conf") { - withSQLConf(CometConf.COMET_ICEBERG_WRITE_SPLIT_OPERATOR_ENABLED.key -> "false") { + withSessionConf(CometConf.COMET_ICEBERG_WRITE_SPLIT_OPERATOR_ENABLED.key -> None) { spark.sql("INSERT INTO cat.db.disabled_conf VALUES (1, 'us-east', 10.5)") } } @@ -3973,24 +3948,28 @@ class CometIcebergWriteActionSuite * (rather than `withSQLConf`) keeps the override visible to the columnar rule across some Spark * version / session-state combinations where `withSQLConf` loses the override before the rule * fires. - * - * Each setting is put back to the value it had before, not unset: unsetting the native write - * flag would turn it on, because it defaults to true while this suite pins it off. */ - private def withNativeEnabled[T](action: => T): T = { + private def withNativeEnabled[T](action: => T): T = + withSessionConf( + CometConf.COMET_ICEBERG_NATIVE_WRITE_ENABLED.key -> Some("true"), + CometConf.COMET_EXEC_LOCAL_TABLE_SCAN_ENABLED.key -> Some("true"))(action) + + /** + * Sets each key to its value, or unsets it for `None`, for the duration of `action`, then puts + * back the value each key had before. It writes the session conf directly; see + * [[withNativeEnabled]] for why. A key is restored rather than unset because unsetting the + * native write flag would turn it on: it defaults to true while this suite pins it off. + */ + private def withSessionConf[T](settings: (String, Option[String])*)(action: => T): T = { val conf = spark.sessionState.conf - val keys = Seq( - CometConf.COMET_ICEBERG_NATIVE_WRITE_ENABLED.key, - CometConf.COMET_EXEC_LOCAL_TABLE_SCAN_ENABLED.key) - val previous = keys.map(key => key -> Option(conf.getConfString(key, null))) - keys.foreach(conf.setConfString(_, "true")) - try action - finally { - previous.foreach { - case (key, Some(value)) => conf.setConfString(key, value) - case (key, None) => conf.unsetConf(key) - } + def set(key: String, value: Option[String]): Unit = value match { + case Some(v) => conf.setConfString(key, v) + case None => conf.unsetConf(key) } + val previous = settings.map { case (key, _) => key -> Option(conf.getConfString(key, null)) } + settings.foreach { case (key, value) => set(key, value) } + try action + finally previous.foreach { case (key, value) => set(key, value) } } /** diff --git a/spark/src/test/scala/org/apache/comet/CometIcebergWriteDetectionSuite.scala b/spark/src/test/scala/org/apache/comet/CometIcebergWriteDetectionSuite.scala index da60c14ff0b..0ddbe8107d4 100644 --- a/spark/src/test/scala/org/apache/comet/CometIcebergWriteDetectionSuite.scala +++ b/spark/src/test/scala/org/apache/comet/CometIcebergWriteDetectionSuite.scala @@ -53,7 +53,6 @@ class CometIcebergWriteDetectionSuite extends CometTestBase with CometIcebergTes override protected def sparkConf: SparkConf = { super.sparkConf - .set(CometConf.COMET_ICEBERG_WRITE_SPLIT_OPERATOR_ENABLED.key, "true") .set(CometConf.COMET_ICEBERG_NATIVE_WRITE_ENABLED.key, "true") } @@ -1440,10 +1439,13 @@ class CometIcebergWriteDetectionSuite extends CometTestBase with CometIcebergTes } } + // With the native writer off, only the testing split flag still plans an `IcebergWriteExec`. test("no fall-back reason is recorded when the iceberg write feature is disabled") { withDetectionCatalog { dir => createTable(dir, "flag_off", partitionSpec = "") - withSQLConf(CometConf.COMET_ICEBERG_NATIVE_WRITE_ENABLED.key -> "false") { + withSQLConf( + CometConf.COMET_ICEBERG_WRITE_SPLIT_OPERATOR_ENABLED.key -> "true", + CometConf.COMET_ICEBERG_NATIVE_WRITE_ENABLED.key -> "false") { val writeExec = insertWriteExec("flag_off") assert( writeExec.getTagValue(CometExplainInfo.FALLBACK_REASONS).isEmpty, diff --git a/spark/src/test/scala/org/apache/spark/sql/benchmark/CometIcebergWriteBenchmark.scala b/spark/src/test/scala/org/apache/spark/sql/benchmark/CometIcebergWriteBenchmark.scala index fbcde2d46d6..75372ead556 100644 --- a/spark/src/test/scala/org/apache/spark/sql/benchmark/CometIcebergWriteBenchmark.scala +++ b/spark/src/test/scala/org/apache/spark/sql/benchmark/CometIcebergWriteBenchmark.scala @@ -147,8 +147,6 @@ object CometIcebergWriteBenchmark extends CometBenchmarkBase { Seq( CometConf.COMET_ENABLED.key -> "true", CometConf.COMET_EXEC_ENABLED.key -> "true", - // The native writer requires the split-operator plan; enabling it alone is a no-op. - CometConf.COMET_ICEBERG_WRITE_SPLIT_OPERATOR_ENABLED.key -> "true", CometConf.COMET_ICEBERG_NATIVE_WRITE_ENABLED.key -> "true"), expectComet = true, expectNativeWrite = true)) From 116bcd027c511f5460e9f4e2699bb501b214eb7f Mon Sep 17 00:00:00 2001 From: Andy Grove Date: Tue, 6 Oct 2026 15:34:40 -0600 Subject: [PATCH 05/10] fix: keep Spark's Iceberg write operator when native execution is off With the new default, an application that turns off spark.comet.exec.enabled and uses Comet only for scans or shuffle got IcebergCommit over a JVM IcebergWrite. The strategy now plans Spark's own operator in that case too. Also brings the remaining docs in line with the single, default-on setting, and notes the 1.1.0 name of the setting in the migration guide. --- docs/source/about/gluten_comparison.md | 3 +- docs/source/contributor-guide/roadmap.md | 14 +++---- docs/source/user-guide/latest/datasources.md | 5 ++- .../user-guide/latest/iceberg-writes.md | 12 +++--- .../user-guide/latest/migration-guide.md | 7 +++- .../latest/understanding-comet-plans.md | 22 +++++------ .../comet/iceberg/IcebergWriteStrategy.scala | 4 +- .../comet/CometIcebergWriteActionSuite.scala | 38 +++++++++++-------- 8 files changed, 58 insertions(+), 47 deletions(-) diff --git a/docs/source/about/gluten_comparison.md b/docs/source/about/gluten_comparison.md index 735662c52c4..bf47853eb1b 100644 --- a/docs/source/about/gluten_comparison.md +++ b/docs/source/about/gluten_comparison.md @@ -113,7 +113,8 @@ covers a broader set of table formats overall. Comet provides a native Iceberg scan built on iceberg-rust. It has been tested with Iceberg 1.5 through 1.10 and supports Iceberg spec v1 and v2, schema evolution, time travel and branch reads, positional and equality deletes -on merge-on-read tables, REST catalogs, and S3-compatible object storage. Iceberg writes still go through Spark. +on merge-on-read tables, REST catalogs, and S3-compatible object storage. For Iceberg writes, Comet +writes the data files of eligible writes natively, and iceberg-java writes the rest and commits every write. Comet does not currently provide native integrations for Delta Lake, Hudi, or Paimon. See the [Comet Iceberg guide] for the full list of supported features and known limitations. diff --git a/docs/source/contributor-guide/roadmap.md b/docs/source/contributor-guide/roadmap.md index d318c07a539..69993597329 100644 --- a/docs/source/contributor-guide/roadmap.md +++ b/docs/source/contributor-guide/roadmap.md @@ -137,14 +137,14 @@ enabled by default ([#1625]). ## Iceberg Table Writes -Comet can now write Iceberg tables natively. The feature is experimental, disabled by default, and controlled by two -settings, the second of which requires the first. The first splits Spark's Iceberg V2 write operator into separate -writer and committer operators, so the query feeding the write becomes visible to AQE and to Comet's columnar rules -([#4658]). The second delegates each task's Parquet write to `iceberg-rust` when the write passes an eligibility -check ([#5361]); writes that don't pass keep Iceberg Java's writer. Merge-on-read writes are not intercepted yet +Since Comet 1.2.0, Comet writes Iceberg tables natively by default, controlled by one setting, +`spark.comet.write.iceberg.enabled` ([#5644]). Comet splits Spark's Iceberg V2 write operator into separate writer +and committer operators, so the query feeding the write becomes visible to AQE and to Comet's columnar rules +([#4658]), and delegates each task's Parquet write to `iceberg-rust` when the write passes an eligibility check +([#5361]); writes that don't pass keep Iceberg Java's writer. Merge-on-read writes are not intercepted yet ([#6240]). The remaining work toward the original goal of [#4322], an ETL job that runs end to end in native code, is -tracked in [#5649]: correctness fixes, failure handling that matches Iceberg Java, broader coverage, and enabling -both settings by default ([#5644]). See [Iceberg Writes](iceberg-writes.md) for how the write path works. +tracked in [#5649]: correctness fixes, failure handling that matches Iceberg Java, and broader coverage. See +[Iceberg Writes](iceberg-writes.md) for how the write path works. [#4322]: https://github.com/apache/datafusion-comet/issues/4322 [#4658]: https://github.com/apache/datafusion-comet/pull/4658 diff --git a/docs/source/user-guide/latest/datasources.md b/docs/source/user-guide/latest/datasources.md index eb56d1e8970..0671fccba9e 100644 --- a/docs/source/user-guide/latest/datasources.md +++ b/docs/source/user-guide/latest/datasources.md @@ -29,8 +29,9 @@ Arrow format, allowing the Comet pipeline to take over after that, but the proce ### Apache Iceberg -Comet accelerates Iceberg scans of Parquet files and has an experimental, opt-in native Iceberg writer. -See the [Iceberg Guide] and [Iceberg Writes](iceberg-writes.md) for more information. +Comet accelerates Iceberg scans of Parquet files and, since Comet 1.2.0, writes the data files of +eligible Iceberg writes natively. See the [Iceberg Guide] and [Iceberg Writes](iceberg-writes.md) +for more information. [iceberg guide]: iceberg.md diff --git a/docs/source/user-guide/latest/iceberg-writes.md b/docs/source/user-guide/latest/iceberg-writes.md index 6b9d520487f..205b0aacfdf 100644 --- a/docs/source/user-guide/latest/iceberg-writes.md +++ b/docs/source/user-guide/latest/iceberg-writes.md @@ -86,7 +86,9 @@ already has, even if they do not fill its first page. ## Configuration -Standard Comet + Iceberg setup (see [`iceberg.md`](iceberg.md)) plus the write-side toggle: +Iceberg writes need only the standard Comet and Iceberg setup (see [`iceberg.md`](iceberg.md)). +`spark.comet.write.iceberg.enabled` is on by default; set it to `false` to plan Spark's own write +operator instead. ``` # Standard Comet / Iceberg wiring @@ -96,10 +98,6 @@ spark.sql.catalog.=org.apache.iceberg.spark.SparkCatalog spark.sql.catalog..type=hadoop # or hive / glue / rest / ... spark.sql.catalog..warehouse=... -# Split-operator plan and native Parquet writer (on by default since Comet 1.2.0; false plans -# Spark's own write operator) -spark.comet.write.iceberg.enabled=true - # Lets writes whose input is a local relation (INSERT ... VALUES, a local DataFrame) use the # native writer; see "Native Parquet write eligibility" below spark.comet.exec.localTableScan.enabled=true @@ -136,7 +134,9 @@ changes. The rewrite is skipped — and the write runs through Spark's stock combined operator — when: - `spark.comet.write.iceberg.enabled` is set to `false`; -- Comet is disabled (`spark.comet.enabled=false`); +- Comet is disabled (`spark.comet.enabled=false`), or its native execution is + (`spark.comet.exec.enabled=false`); +- Comet is in plan-only mode (`spark.comet.explain.planOnly.enabled=true`); - the write is not an Iceberg `SparkWrite` (any other V2 data source); - the table uses merge-on-read: delta writes (Iceberg `WriteDelta`) are not intercepted; - the statement is CTAS / RTAS on Spark 3.4, where the staged exec writes inline; on Spark diff --git a/docs/source/user-guide/latest/migration-guide.md b/docs/source/user-guide/latest/migration-guide.md index 2b5479d3e1f..6f32b56fb6d 100644 --- a/docs/source/user-guide/latest/migration-guide.md +++ b/docs/source/user-guide/latest/migration-guide.md @@ -93,8 +93,11 @@ iceberg-java's writer, and iceberg-java still commits every write. The table hol 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. `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). +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 diff --git a/docs/source/user-guide/latest/understanding-comet-plans.md b/docs/source/user-guide/latest/understanding-comet-plans.md index 655865ff250..f77cd1b5165 100644 --- a/docs/source/user-guide/latest/understanding-comet-plans.md +++ b/docs/source/user-guide/latest/understanding-comet-plans.md @@ -329,22 +329,22 @@ consecutively in a plan, they execute as a single fused block. | `CometTakeOrderedAndProject` | `TakeOrderedAndProjectExec` | | `CometWriteFiles` | `WriteFilesExec` (Spark 4.0 and later, experimental native Parquet writes) | | `CometNativeWrite` | `DataWritingCommandExec` (Spark 3.x, experimental native Parquet writes) | -| `CometIcebergWrite` | `IcebergWrite` (experimental native Iceberg data-file writes) | +| `CometIcebergWrite` | `IcebergWrite` (native Iceberg data-file writes) | ### JVM-Side Operators These keep their data on the JVM but participate in the Comet pipeline. -| Node | Notes | -| ------------------------ | ------------------------------------------------------------------------------------------------------------------------------------------- | -| `CometUnion` | JVM-side union of Comet inputs. The Rust side reads each branch as a separate scan. | -| `CometCoalesce` | JVM-side partition coalesce. | -| `CometCollectLimit` | JVM-side collect limit, equivalent to `CollectLimitExec`. | -| `CometBroadcastExchange` | Broadcast exchange producing serialized Arrow batches that the consumer can decode. | -| `CometSubqueryBroadcast` | Companion to `CometBroadcastExchange` for dynamic partition pruning subqueries. | -| `CometMapInBatch` | Runs `mapInArrow` / `mapInPandas` Python UDFs on Comet's Arrow batches (experimental). See [PyArrow UDF Acceleration](pyarrow-udfs.md). | -| `IcebergWrite` | Executor-side Iceberg data-file write in Comet's split-operator Iceberg write plan (experimental). See [Iceberg Writes](iceberg-writes.md). | -| `IcebergCommit` | Driver-side commit for Comet's split-operator Iceberg write plan (experimental). | +| Node | Notes | +| ------------------------ | --------------------------------------------------------------------------------------------------------------------------------------- | +| `CometUnion` | JVM-side union of Comet inputs. The Rust side reads each branch as a separate scan. | +| `CometCoalesce` | JVM-side partition coalesce. | +| `CometCollectLimit` | JVM-side collect limit, equivalent to `CollectLimitExec`. | +| `CometBroadcastExchange` | Broadcast exchange producing serialized Arrow batches that the consumer can decode. | +| `CometSubqueryBroadcast` | Companion to `CometBroadcastExchange` for dynamic partition pruning subqueries. | +| `CometMapInBatch` | Runs `mapInArrow` / `mapInPandas` Python UDFs on Comet's Arrow batches (experimental). See [PyArrow UDF Acceleration](pyarrow-udfs.md). | +| `IcebergWrite` | Executor-side Iceberg data-file write in Comet's split-operator Iceberg write plan. See [Iceberg Writes](iceberg-writes.md). | +| `IcebergCommit` | Driver-side commit for Comet's split-operator Iceberg write plan. | ### Shuffle Operators diff --git a/spark/src/main/scala/org/apache/comet/iceberg/IcebergWriteStrategy.scala b/spark/src/main/scala/org/apache/comet/iceberg/IcebergWriteStrategy.scala index 057c803ad57..1d262ec6896 100644 --- a/spark/src/main/scala/org/apache/comet/iceberg/IcebergWriteStrategy.scala +++ b/spark/src/main/scala/org/apache/comet/iceberg/IcebergWriteStrategy.scala @@ -43,8 +43,8 @@ case class IcebergWriteStrategy(session: SparkSession) extends SparkStrategy { val splitEnabled = CometConf.COMET_ICEBERG_NATIVE_WRITE_ENABLED.get(conf) || CometConf.COMET_ICEBERG_WRITE_SPLIT_OPERATOR_ENABLED.get(conf) // Planner strategies run whether or not Comet is enabled, so check it here too: with Comet - // off, Spark must plan its own V2 write operator. - if (!isCometLoaded(conf) || !splitEnabled) { + // or its native execution off, Spark must plan its own V2 write operator. + if (!isCometLoaded(conf) || !CometConf.COMET_EXEC_ENABLED.get(conf) || !splitEnabled) { return Nil } // Planner strategies run before CometRule, so plan-only mode needs its own guard here. diff --git a/spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala b/spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala index 39b5a07b911..3dcce123082 100644 --- a/spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala +++ b/spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala @@ -119,25 +119,31 @@ class CometIcebergWriteActionSuite } } - test("spark.comet.enabled=false keeps Spark's own write plan with the split flag on") { - 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 " + + // 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")) } - assert(snapshot.snapshotDelta == 1L, s"expected 1 commit, got ${snapshot.snapshotDelta}") - val (commits, writes) = collectIcebergWriteOps(snapshot.plans) - assert( - commits.isEmpty && writes.isEmpty, - "expected Spark's own write plan with Comet disabled. Plans:\n" + - snapshot.plans.mkString("\n--\n")) + assertRows(table, expectedIds = Seq(1, 2, 3)) } - assertRows("comet_disabled", expectedIds = Seq(1, 2, 3)) } } From 6864040bb05cceb01d58e43f3add18338e97f31c Mon Sep 17 00:00:00 2001 From: Andy Grove Date: Tue, 6 Oct 2026 16:35:55 -0600 Subject: [PATCH 06/10] test: keep the Iceberg write when a transition-heavy stage reverts with default flags With no write flag set, an eligible write becomes a CometIcebergWriteExec, so transition reversion has to restore a JVM writer (#5719, fixed by #5957). --- .../comet/CometIcebergWriteActionSuite.scala | 46 +++++++++++++++++++ 1 file changed, 46 insertions(+) diff --git a/spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala b/spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala index 3dcce123082..96a239f8f93 100644 --- a/spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala +++ b/spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala @@ -184,6 +184,52 @@ class CometIcebergWriteActionSuite } } + // 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, + "transition reversion should restore IcebergWriteExec. Plans:\n" + + snapshot.plans.mkString("\n--\n")) + } + assertRows(table, Seq(1, 2, 3)) + } + } + } + } + test("AppendData partitioned INSERT INTO routes through two-op") { assume(icebergAvailable, "Iceberg not available in classpath") withIcebergCatalog { warehouseDir => From 70bb11db284ab97f1687975781fb648c8a6e786b Mon Sep 17 00:00:00 2001 From: Andy Grove Date: Wed, 7 Oct 2026 21:32:46 -0600 Subject: [PATCH 07/10] test: show the transition-heavy write runs natively before its stage is reverted Run the same write under the default threshold first and require a CometIcebergWriteExec, so the reverted run cannot pass because the write was never eligible. --- .../comet/CometIcebergWriteActionSuite.scala | 67 +++++++++++++------ 1 file changed, 46 insertions(+), 21 deletions(-) diff --git a/spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala b/spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala index c5965cf4af6..61773733464 100644 --- a/spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala +++ b/spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala @@ -186,7 +186,9 @@ class CometIcebergWriteActionSuite } // 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. + // transition-heavy stage has to put a JVM writer back rather than drop the write. The same write + // runs first under the default threshold, which leaves its stage alone, so that the reverted run + // is known to start from a native write rather than from one that was never eligible. // 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") { @@ -203,29 +205,52 @@ class CometIcebergWriteActionSuite .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() + // Appends the source with no write flag set, checks that it committed once, and returns + // the plans it ran. + def append(maxTransitions: Int): Seq[SparkPlan] = { + var plans = Seq.empty[SparkPlan] + // withSQLConf returns Unit before Spark 4.0, so the plans are kept from 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 -> maxTransitions.toString) { + 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() + } } + assert( + snapshot.snapshotDelta == 1L, + s"expected 1 commit, got ${snapshot.snapshotDelta}") + plans = snapshot.plans } - 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, - "transition reversion should restore IcebergWriteExec. Plans:\n" + - snapshot.plans.mkString("\n--\n")) + plans } - assertRows(table, Seq(1, 2, 3)) + + // The default threshold of two transitions leaves the stage alone. + val native = append(maxTransitions = 2) + val (nativeCommits, _) = collectIcebergWriteOps(native) + assert( + nativeCommits.nonEmpty && native.exists { plan => + collectWithSubqueries(plan) { case e: CometIcebergWriteExec => e }.nonEmpty + }, + "expected an IcebergCommitExec over a CometIcebergWriteExec. Plans:\n" + + native.mkString("\n--\n")) + + val reverted = append(maxTransitions = 0) + val (revertedCommits, revertedWrites) = collectIcebergWriteOps(reverted) + assert( + revertedCommits.nonEmpty && revertedWrites.nonEmpty && reverted.forall { plan => + collectWithSubqueries(plan) { case e: CometIcebergWriteExec => e }.isEmpty + }, + "transition reversion should restore IcebergWriteExec. Plans:\n" + + reverted.mkString("\n--\n")) + assertRows(table, Seq(1, 1, 2, 2, 3, 3)) } } } From 1cb149f6e610c722f977a1c94efaaa69f9d17ddc Mon Sep 17 00:00:00 2001 From: Andy Grove Date: Thu, 8 Oct 2026 06:42:34 -0600 Subject: [PATCH 08/10] docs: describe the footer reads native Iceberg writes make on object stores The JVM reads each natively written file's footer back to compute its Iceberg metrics. Say that this is two reads per file, made one after another, which on S3 or GCS are GET requests, so a write of many small files pays two round trips for each. --- docs/source/user-guide/latest/iceberg-writes.md | 9 ++++++++- docs/source/user-guide/latest/migration-guide.md | 5 +++++ 2 files changed, 13 insertions(+), 1 deletion(-) diff --git a/docs/source/user-guide/latest/iceberg-writes.md b/docs/source/user-guide/latest/iceberg-writes.md index 551b788dec9..b609ea45e18 100644 --- a/docs/source/user-guide/latest/iceberg-writes.md +++ b/docs/source/user-guide/latest/iceberg-writes.md @@ -293,7 +293,14 @@ bounds carried over from the native writer's tracked state. iceberg-java's metad — metrics modes, the inferred-column cap (`write.metadata.metrics.max-inferred-column-defaults`), bound truncation, and list/map bounds suppression — are therefore applied by iceberg-java's own code regardless of what the native -writer reports. This costs one footer-sized ranged read per written file at write time. +writer reports. + +iceberg-java's writer takes these metrics from the footer it still holds in memory, but the native +path reads each footer back from storage: two small reads per written file, one for the footer's +length and one for the footer itself, made one file after another before the task finishes. On S3 +or GCS each read is a GET request, so the task waits two request round trips per file. That is +small next to uploading a large file, but a write that produces many small files, such as a fanout +write over many partitions, pays it for every one of them. ## Failure handling diff --git a/docs/source/user-guide/latest/migration-guide.md b/docs/source/user-guide/latest/migration-guide.md index 6f32b56fb6d..1742fd97ee9 100644 --- a/docs/source/user-guide/latest/migration-guide.md +++ b/docs/source/user-guide/latest/migration-guide.md @@ -99,6 +99,11 @@ 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). +To compute the Iceberg metrics of a natively written file, Comet reads the file's Parquet footer +back from storage, where iceberg-java's writer takes them from memory. On S3 or GCS that is two GET +requests per file, which adds up for a write that produces many small files. See +[Native Parquet write eligibility](iceberg-writes.md#native-parquet-write-eligibility). + 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 From 94915161e926b0cd8bf1cf27855cc0f0086ba1a3 Mon Sep 17 00:00:00 2001 From: Andy Grove Date: Thu, 8 Oct 2026 06:51:07 -0600 Subject: [PATCH 09/10] fix: keep Spark's operator for merge-on-read Iceberg writes unless the testing split flag is on #6693 plans Spark 3.5+ WriteDelta through the split plan, with Iceberg's JVM DeltaWriter still writing the rows. With spark.comet.write.iceberg.enabled on by default, every merge-on-read write would take that path although the plan gives it nothing over Spark's own operator. Only the testing split flag plans it now, so users keep Spark's operator until a native delta writer exists (#6240), while the Iceberg Spark tests, which set both flags, still cover it. --- .../contributor-guide/iceberg-spark-tests.md | 4 +- .../contributor-guide/iceberg-writes.md | 20 +++++----- docs/source/contributor-guide/roadmap.md | 4 +- .../user-guide/latest/iceberg-writes.md | 13 ++++--- .../scala/org/apache/comet/CometConf.scala | 7 ++-- .../comet/iceberg/IcebergWriteStrategy.scala | 12 ++++-- .../comet/CometIcebergWriteActionSuite.scala | 37 +++++++++++++++++++ 7 files changed, 72 insertions(+), 25 deletions(-) diff --git a/docs/source/contributor-guide/iceberg-spark-tests.md b/docs/source/contributor-guide/iceberg-spark-tests.md index 812c53fca50..da60288e57f 100644 --- a/docs/source/contributor-guide/iceberg-spark-tests.md +++ b/docs/source/contributor-guide/iceberg-spark-tests.md @@ -169,8 +169,8 @@ records one of three writers: - `jvm`: Comet's split operator planned the write but kept Iceberg's JVM writer (`IcebergWriteExec`). The line includes the reasons Comet recorded for not converting it. - `spark`: Spark's own V2 write operator ran the write, so Comet's split operator never saw it. Examples - are `WriteDelta` for merge-on-read, `WriteToDataSourceV2` for a streaming micro-batch, and on Spark - 3.4 the CTAS and RTAS execs, which write the table themselves. + are `WriteToDataSourceV2` for a streaming micro-batch, and on Spark 3.4 both `WriteDelta` for + merge-on-read and the CTAS and RTAS execs, which write the table themselves. `dev/ci/summarize-iceberg-writes.py` turns these records into a table on the job's summary page. It shows the count and share of each writer, the most common fallback reasons, and the Spark write diff --git a/docs/source/contributor-guide/iceberg-writes.md b/docs/source/contributor-guide/iceberg-writes.md index e0d7cc16266..bbcf5a6ac83 100644 --- a/docs/source/contributor-guide/iceberg-writes.md +++ b/docs/source/contributor-guide/iceberg-writes.md @@ -73,13 +73,13 @@ IcebergCommit driver: collect task commit messages, BatchWrite.c +- scans, projects, exchanges, sorts; inside AQE with or without the split ``` -| Component | Location | Role | -| ------------------------------------------------------------------------------------------------------------------------------------ | -------------------------------------------------- | ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------ | -| `IcebergWriteStrategy` | `spark/src/main/scala/org/apache/comet/iceberg/` | Planner strategy. Matches `AppendData`, `OverwriteByExpression`, `OverwritePartitionsDynamic`, `ReplaceData` (and Iceberg's own `ReplaceIcebergData`, which Iceberg 1.5.2 plans on Spark 3.4), plus Spark 3.5+ Iceberg `WriteDelta`. | -| `IcebergWriteLogical` | same | Logical anchor for the writer, so AQE re-plans re-emit only the writer and not a second committer. | -| `IcebergWriteExec` | `spark/src/main/scala/org/apache/spark/sql/comet/` | JVM writer. Runs iceberg-java's `DataWriter` or, for `WriteDelta`, `DeltaWriter` per task and returns the serialized `WriterCommitMessage` as one binary row. | -| `IcebergCommitExec` | same | Driver committer. A `V2CommandExec`, so `run()` is memoized and the commit happens once. | -| `IcebergReplaceDataShim`, `IcebergDeltaLogicalShim`, `IcebergDeltaWriterShim`, `IcebergRefreshCacheShim`, `IcebergDriverMetricsShim` | `spark/src/main/spark-*/org/apache/comet/iceberg/` | Version differences: operation-coded `ReplaceData` / `WriteDelta` rows, cache refresh by name on 4.1+, driver metric reporting. | +| Component | Location | Role | +| ------------------------------------------------------------------------------------------------------------------------------------ | -------------------------------------------------- | ----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | +| `IcebergWriteStrategy` | `spark/src/main/scala/org/apache/comet/iceberg/` | Planner strategy. Matches `AppendData`, `OverwriteByExpression`, `OverwritePartitionsDynamic`, `ReplaceData` (and Iceberg's own `ReplaceIcebergData`, which Iceberg 1.5.2 plans on Spark 3.4), plus Spark 3.5+ Iceberg `WriteDelta` under the testing split flag. | +| `IcebergWriteLogical` | same | Logical anchor for the writer, so AQE re-plans re-emit only the writer and not a second committer. | +| `IcebergWriteExec` | `spark/src/main/scala/org/apache/spark/sql/comet/` | JVM writer. Runs iceberg-java's `DataWriter` or, for `WriteDelta`, `DeltaWriter` per task and returns the serialized `WriterCommitMessage` as one binary row. | +| `IcebergCommitExec` | same | Driver committer. A `V2CommandExec`, so `run()` is memoized and the commit happens once. | +| `IcebergReplaceDataShim`, `IcebergDeltaLogicalShim`, `IcebergDeltaWriterShim`, `IcebergRefreshCacheShim`, `IcebergDriverMetricsShim` | `spark/src/main/spark-*/org/apache/comet/iceberg/` | Version differences: operation-coded `ReplaceData` / `WriteDelta` rows, cache refresh by name on 4.1+, driver metric reporting. | Things to know before changing this layer: @@ -97,8 +97,10 @@ Things to know before changing this layer: position-delta write do not ask for one; the checks in `buildTwoOp` and `buildDeltaTwoOp` are defensive. - **WriteDelta stays JVM-backed.** Spark 3.5+ merge-on-read commands are intercepted by the split - plan, but `IcebergWriteExec` delegates their rows to Iceberg's JVM `DeltaWriter`; the native - Iceberg writer is explicitly declined. Spark 3.4 `WriteDelta` keeps Spark's plan. + plan only when the testing flag `spark.comet.write.iceberg.splitOperator.enabled` is on, and then + `IcebergWriteExec` delegates their rows to Iceberg's JVM `DeltaWriter`; the native Iceberg writer + is explicitly declined. Since the split plan gives such a write nothing over Spark's operator, + `spark.comet.write.iceberg.enabled` alone leaves it on Spark's plan, as does Spark 3.4. - **What is not intercepted:** streaming writes and CTAS/RTAS on Spark 3.4. Those keep Spark's plan. On Spark 4.1+, `IcebergWriteSummaryShim` finds either Spark's `MergeRowsExec` or diff --git a/docs/source/contributor-guide/roadmap.md b/docs/source/contributor-guide/roadmap.md index 69993597329..aca1712afdd 100644 --- a/docs/source/contributor-guide/roadmap.md +++ b/docs/source/contributor-guide/roadmap.md @@ -141,8 +141,8 @@ Since Comet 1.2.0, Comet writes Iceberg tables natively by default, controlled b `spark.comet.write.iceberg.enabled` ([#5644]). Comet splits Spark's Iceberg V2 write operator into separate writer and committer operators, so the query feeding the write becomes visible to AQE and to Comet's columnar rules ([#4658]), and delegates each task's Parquet write to `iceberg-rust` when the write passes an eligibility check -([#5361]); writes that don't pass keep Iceberg Java's writer. Merge-on-read writes are not intercepted yet -([#6240]). The remaining work toward the original goal of [#4322], an ETL job that runs end to end in native code, is +([#5361]); writes that don't pass keep Iceberg Java's writer. Merge-on-read writes keep Spark's write operator +until the native writer can write them ([#6240]). The remaining work toward the original goal of [#4322], an ETL job that runs end to end in native code, is tracked in [#5649]: correctness fixes, failure handling that matches Iceberg Java, and broader coverage. See [Iceberg Writes](iceberg-writes.md) for how the write path works. diff --git a/docs/source/user-guide/latest/iceberg-writes.md b/docs/source/user-guide/latest/iceberg-writes.md index b609ea45e18..38fdaf4b24d 100644 --- a/docs/source/user-guide/latest/iceberg-writes.md +++ b/docs/source/user-guide/latest/iceberg-writes.md @@ -124,10 +124,11 @@ analyzer emits operation-coded rows that Comet's writer dispatches through `Repl projections, while on Spark 3.4/3.5 the rewritten rows are written as a plain row stream. The supported set of operations is the same either way. -On Spark 3.5+, merge-on-read uses Spark's `WriteDelta`. The split plan intercepts that command so -Comet can keep the same driver commit and reporting path, but task-side row-level writes stay on -Iceberg's JVM `DeltaWriter`; `CometIcebergWriteExec` is never used for position-delta rows. -Spark 3.4 leaves `WriteDelta` on Spark's stock write plan. +Merge-on-read uses Spark's `WriteDelta`, and Iceberg's JVM `DeltaWriter` writes its rows under +either plan, so it keeps Spark's stock write plan. On Spark 3.5+ the testing setting +`spark.comet.write.iceberg.splitOperator.enabled` plans it as the split plan, with the same driver +commit and reporting path, but `CometIcebergWriteExec` is never used for position-delta rows. +Spark 3.4 always leaves `WriteDelta` on Spark's stock write plan. On Spark 4.1+ the split plan matches two further stock-Spark behaviours: MERGE metrics are forwarded to the writer's commit (Iceberg 1.11+ records them in the snapshot summary), and @@ -143,8 +144,8 @@ The rewrite is skipped — and the write runs through Spark's stock combined ope (`spark.comet.exec.enabled=false`); - Comet is in plan-only mode (`spark.comet.explain.planOnly.enabled=true`); - the write is neither an Iceberg `SparkWrite` nor a supported Iceberg position-delta write; -- the table uses merge-on-read on Spark 3.4; Spark 3.5+ `WriteDelta` is intercepted but remains - on Iceberg's JVM `DeltaWriter`; +- the table uses merge-on-read (`WriteDelta`), unless the testing setting + `spark.comet.write.iceberg.splitOperator.enabled` is on with Spark 3.5+; - the statement is CTAS / RTAS on Spark 3.4, where the staged exec writes inline; on Spark 3.5+ those statements re-plan their inner append, which is intercepted normally; - the write requires Spark's commit coordinator, which Comet's per-task commit protocol does diff --git a/spark/src/main/scala/org/apache/comet/CometConf.scala b/spark/src/main/scala/org/apache/comet/CometConf.scala index 76a1a2d1975..0b441eff12c 100644 --- a/spark/src/main/scala/org/apache/comet/CometConf.scala +++ b/spark/src/main/scala/org/apache/comet/CometConf.scala @@ -127,8 +127,8 @@ object CometConf extends ShimCometConf { .doc( "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.") + "every data file inside Comet's plan. It is also what plans merge-on-read writes on " + + "Spark 3.5+ that way. Used by tests to compare the two writers under the same plan.") .booleanConf .createWithDefault(false) @@ -140,7 +140,8 @@ object CometConf extends ShimCometConf { "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.") + "commits every write. Merge-on-read writes keep Spark's operator. Set this to false " + + "to plan Spark's own V2 write operator.") .booleanConf .createWithDefault(true) diff --git a/spark/src/main/scala/org/apache/comet/iceberg/IcebergWriteStrategy.scala b/spark/src/main/scala/org/apache/comet/iceberg/IcebergWriteStrategy.scala index f5764caba05..49cc6519867 100644 --- a/spark/src/main/scala/org/apache/comet/iceberg/IcebergWriteStrategy.scala +++ b/spark/src/main/scala/org/apache/comet/iceberg/IcebergWriteStrategy.scala @@ -31,8 +31,9 @@ import org.apache.comet.CometSparkSessionExtensions.isCometLoaded import org.apache.comet.shims.ShimCometMergeRows /** - * Spark strategy for Comet's split Iceberg V2 writer. WriteDelta is modeled and dispatched here, - * but position-delta rows are still executed by Iceberg's JVM DeltaWriter. + * Spark strategy for Comet's split Iceberg V2 writer. WriteDelta is modeled and dispatched here + * when the testing split flag is on, but position-delta rows are still executed by Iceberg's JVM + * DeltaWriter. */ case class IcebergWriteStrategy(session: SparkSession) extends SparkStrategy { @@ -82,7 +83,11 @@ case class IcebergWriteStrategy(session: SparkSession) extends SparkStrategy { // Hit by AQE. case l @ IcebergWriteLogical(child, batchWrite, dispatch) => Seq(IcebergWriteExec(batchWrite, l.output, planLater(child), dispatch)) - case delta => + // Iceberg's JVM DeltaWriter writes a merge-on-read write's rows under the split plan too, + // so the plan would give it nothing over Spark's own operator. Until the native writer can + // write them, only the testing split flag plans it. + // https://github.com/apache/datafusion-comet/issues/6240 + case delta if CometConf.COMET_ICEBERG_WRITE_SPLIT_OPERATOR_ENABLED.get(conf) => IcebergDeltaLogicalShim .extract(delta) .flatMap { fields => @@ -110,6 +115,7 @@ case class IcebergWriteStrategy(session: SparkSession) extends SparkStrategy { } } .toList + case _ => Nil } } diff --git a/spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala b/spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala index 61773733464..0bc1abf591d 100644 --- a/spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala +++ b/spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala @@ -185,6 +185,43 @@ class CometIcebergWriteActionSuite } } + // Iceberg's DeltaWriter writes a merge-on-read write's rows under either plan, so with no flag + // set the write keeps Spark's own operator. Only the testing split flag plans it as Comet's. + // https://github.com/apache/datafusion-comet/issues/6240 + test("a merge-on-read write keeps Spark's own write plan when no flag is set") { + assume(icebergAvailable, "Iceberg not available in classpath") + assume(isSpark35Plus, "WriteDelta interception starts with Spark 3.5") + withIcebergCatalog { warehouseDir => + createTable( + warehouseDir, + "mor_defaults", + partitionSpec = "PARTITIONED BY (region)", + properties = Some("'format-version'='2', 'write.delete.mode'='merge-on-read'")) + // Ids 1 and 2 share a data file, so deleting id 1 writes a delete file through WriteDelta + // rather than dropping a whole file. + coalesceInsert("mor_defaults", Seq((1, "a", 10.0), (2, "a", 20.0), (3, "c", 30.0))) + val snapshot = captureWrite("mor_defaults") { + withSessionConf( + CometConf.COMET_ICEBERG_WRITE_SPLIT_OPERATOR_ENABLED.key -> None, + CometConf.COMET_ICEBERG_NATIVE_WRITE_ENABLED.key -> None) { + spark.sql(s"DELETE FROM $catalog.$ns.mor_defaults WHERE id = 1") + } + } + assert(snapshot.snapshotDelta == 1L, s"expected 1 commit, got ${snapshot.snapshotDelta}") + val (commits, writes) = collectIcebergWriteOps(snapshot.plans) + val sparkDeltaWrites = snapshot.plans.flatMap { plan => + collectWithSubqueries(plan) { + case p if p.getClass.getSimpleName == "WriteDeltaExec" => p + } + } + assert( + commits.isEmpty && writes.isEmpty && sparkDeltaWrites.nonEmpty, + "expected Spark's WriteDeltaExec for a merge-on-read DELETE. Plans:\n" + + snapshot.plans.mkString("\n--\n")) + assertRows("mor_defaults", expectedIds = Seq(2, 3)) + } + } + // 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. The same write // runs first under the default threshold, which leaves its stage alone, so that the reverted run From b8cadbad483ed9b6c36bfa625ad46122b6886143 Mon Sep 17 00:00:00 2001 From: Andy Grove Date: Thu, 8 Oct 2026 16:47:03 -0600 Subject: [PATCH 10/10] docs: describe native Iceberg writes in the FAQ as on by default The FAQ added on main called Comet's native Iceberg writes experimental. They are now on by default for eligible writes. --- docs/source/about/faq.md | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/docs/source/about/faq.md b/docs/source/about/faq.md index 25edfe2a816..d92976a1e19 100644 --- a/docs/source/about/faq.md +++ b/docs/source/about/faq.md @@ -141,10 +141,10 @@ Common causes of a small speedup are: ### Does Comet support Delta Lake, Apache Hudi, or Apache Paimon tables? Not yet. Comet accelerates [Apache Iceberg](../user-guide/latest/iceberg.md) tables, with a native -reader that is enabled by default and -[experimental native writes](../user-guide/latest/iceberg-writes.md). It does not accelerate scans of -Delta Lake, Hudi, or Paimon tables, so Spark reads those. Native Delta Lake reads are in development; -see the [roadmap](../contributor-guide/roadmap.md#delta-lake-support). Hudi and Paimon are not on the +reader and [native writes](../user-guide/latest/iceberg-writes.md) that are both enabled by default. +It does not accelerate scans of Delta Lake, Hudi, or Paimon tables, so Spark reads those. Native +Delta Lake reads are in development; see the +[roadmap](../contributor-guide/roadmap.md#delta-lake-support). Hudi and Paimon are not on the roadmap. ### Does Comet accelerate PySpark jobs and Python UDFs?