diff --git a/.ai/skills/review-comet-iceberg-write-pr/SKILL.md b/.ai/skills/review-comet-iceberg-write-pr/SKILL.md index a835f45441a..5b193b16e6b 100644 --- a/.ai/skills/review-comet-iceberg-write-pr/SKILL.md +++ b/.ai/skills/review-comet-iceberg-write-pr/SKILL.md @@ -43,14 +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 | - -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. +| 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/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? diff --git a/docs/source/about/gluten_comparison.md b/docs/source/about/gluten_comparison.md index 56fb5340cc9..fbf1a6b7faf 100644 --- a/docs/source/about/gluten_comparison.md +++ b/docs/source/about/gluten_comparison.md @@ -117,10 +117,11 @@ 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 and 1.8 through 1.11 and supports Iceberg spec v1, v2, and v3, schema evolution, time travel and branch reads, positional and equality -deletes and deletion vectors, encrypted v3 tables, REST catalogs, and S3-compatible object storage. Comet can also -write Iceberg data files natively through iceberg-rust, as an experimental feature that is disabled by default (see -[Comet Iceberg writes]). 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. +deletes and deletion vectors, encrypted v3 tables, REST catalogs, and S3-compatible object storage. For Iceberg +writes, Comet writes the data files of eligible writes natively through iceberg-rust by default, and iceberg-java +writes the rest and commits every write (see [Comet Iceberg writes]). 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. [Comet Iceberg guide]: /user-guide/latest/iceberg.md [Comet Iceberg writes]: /user-guide/latest/iceberg-writes.md diff --git a/docs/source/contributor-guide/iceberg-spark-tests.md b/docs/source/contributor-guide/iceberg-spark-tests.md index 87ef47427d2..da60288e57f 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 @@ -171,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 1ef7539329a..bbcf5a6ac83 100644 --- a/docs/source/contributor-guide/iceberg-writes.md +++ b/docs/source/contributor-guide/iceberg-writes.md @@ -30,16 +30,19 @@ does not repeat those lists; it explains the code that implements them. ## Overview -Two flags, each of which builds on the one before it: - -| 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. | - -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 +Two layers, the second built on the first, both switched by `spark.comet.write.iceberg.enabled`: + +| 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 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 @@ -70,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: @@ -94,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 @@ -444,7 +449,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/contributor-guide/roadmap.md b/docs/source/contributor-guide/roadmap.md index d318c07a539..aca1712afdd 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 -([#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. +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 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. [#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 f351aff5b6f..8ac564dba73 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 a04f6be61cf..38fdaf4b24d 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 and Native Writer -**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, 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). + +`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 @@ -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.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 @@ -41,13 +46,12 @@ 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` 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 @@ -82,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 @@ -92,12 +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 (experimental, off by default) -spark.comet.write.iceberg.splitOperator.enabled=true - -# Native Parquet writer (experimental, off by default; 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 # native writer; see "Native Parquet write eligibility" below spark.comet.exec.localTableScan.enabled=true @@ -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 @@ -138,10 +139,13 @@ 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.enabled` is set to `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 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 @@ -154,7 +158,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, @@ -172,8 +176,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 @@ -290,7 +294,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 @@ -344,11 +355,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 4c3abe00216..9a247520f5b 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 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 3ce3e88d1da..1742fd97ee9 100644 --- a/docs/source/user-guide/latest/migration-guide.md +++ b/docs/source/user-guide/latest/migration-guide.md @@ -83,6 +83,36 @@ it. `spark.comet.convert.oneRowRelation.enabled` is on by default, so such a que without either setting, and logs no warning. The list remains the way to convert other leaf operators, such as the scan of a Data Source V2 connector. +### Iceberg Writes + +`spark.comet.write.iceberg.enabled` now defaults to `true`. Comet plans an Iceberg `INSERT INTO`, +`INSERT OVERWRITE`, and copy-on-write `DELETE`, `UPDATE` or `MERGE` as two operators, `IcebergWrite` +under `IcebergCommit`, in place of Spark's single V2 write operator, and writes the data files of +each eligible write natively with iceberg-rust. A write that is not eligible still uses +iceberg-java's writer, and iceberg-java still commits every write. The table holds the same rows, +but natively written files differ from iceberg-java's in the ways listed under +[Accepted divergences](iceberg-writes.md#accepted-divergences), and explain output and the Spark UI +show the new operators. Set `spark.comet.write.iceberg.enabled=false` to plan Spark's own operator, +as in Comet 1.1.0. Comet 1.1.0 named this setting `spark.comet.iceberg.write.enabled`, in the +testing category, and Comet now ignores that name, so a deployment that set it to `false` gets the +new default unless it sets `spark.comet.write.iceberg.enabled=false`. +`spark.comet.write.iceberg.splitOperator.enabled` is now a testing setting, and setting it to +`false` does not turn the split operator off. See [Iceberg Writes](iceberg-writes.md). + +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 +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 diff --git a/docs/source/user-guide/latest/operators.md b/docs/source/user-guide/latest/operators.md index acefd4f0890..364e82ac51b 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` | ⚠️ | 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.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..f77cd1b5165 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 @@ -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/CometConf.scala b/spark/src/main/scala/org/apache/comet/CometConf.scala index a187f8fa5f8..76e47991180 100644 --- a/spark/src/main/scala/org/apache/comet/CometConf.scala +++ b/spark/src/main/scala/org/apache/comet/CometConf.scala @@ -125,21 +125,25 @@ object CometConf extends ShimCometConf { conf("spark.comet.write.iceberg.splitOperator.enabled") .category(CATEGORY_TESTING) .doc( - "Whether to rewrite Iceberg V2 writes from Spark's combined V2 write/commit operator " + - "into Comet's two-operator shape: a file writer exec (inside AQE) and a committer " + - "(outside AQE).") + "Whether to plan Iceberg writes as Comet's file writer and committer even when " + + "`spark.comet.write.iceberg.enabled` is false, so that Iceberg's own writer writes " + + "every data file inside Comet's plan. 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) val COMET_ICEBERG_NATIVE_WRITE_ENABLED: ConfigEntry[Boolean] = conf("spark.comet.write.iceberg.enabled") - .category(CATEGORY_TESTING) + .category(CATEGORY_EXEC) .doc( - "Whether to delegate the executor-side Parquet write to Comet's native (iceberg-rust) " + - "writer when the table's properties allow it. Requires " + - "`spark.comet.write.iceberg.splitOperator.enabled = true`. Off by default.") + "Whether Comet plans Iceberg writes and writes their data files natively. Comet " + + "replaces Spark's combined V2 write operator with a file writer (inside AQE) under a " + + "committer (outside AQE), and writes the data files of each eligible write with its " + + "native (iceberg-rust) writer. Other writes use Iceberg's own writer, and Iceberg " + + "commits every write. Merge-on-read writes keep Spark's operator. Set this to false " + + "to plan Spark's own V2 write operator.") .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/main/scala/org/apache/comet/iceberg/IcebergWriteStrategy.scala b/spark/src/main/scala/org/apache/comet/iceberg/IcebergWriteStrategy.scala index dbb19fba797..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,16 +31,21 @@ 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 { 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)) { + // 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. @@ -78,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 => @@ -106,6 +115,7 @@ case class IcebergWriteStrategy(session: SparkSession) extends SparkStrategy { } } .toList + case _ => Nil } } 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 45a4fe597fd..0bc1abf591d 100644 --- a/spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala +++ b/spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala @@ -94,7 +94,10 @@ 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") + .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]") @@ -117,25 +120,176 @@ class CometIcebergWriteActionSuite } } - test("spark.comet.enabled=false keeps Spark's own write plan with the split flag on") { + // Turning off Comet, or only its native execution as an application that uses Comet just for + // scans or shuffle does, keeps Spark's own write operator. + Seq(CometConf.COMET_ENABLED, CometConf.COMET_EXEC_ENABLED).foreach { flag => + test(s"${flag.key}=false keeps Spark's own write plan with the split flag on") { + assume(icebergAvailable, "Iceberg not available in classpath") + withIcebergCatalog { warehouseDir => + val table = flag.key.replace('.', '_') + createTable(warehouseDir, table, partitionSpec = "") + // withSQLConf returns Unit before Spark 4.0, so the assertions run inside it. + withSQLConf(flag.key -> "false") { + val snapshot = captureWrite(table) { + spark.sql(s"INSERT INTO cat.db.$table VALUES " + + "(1, 'us-east', 10.5), (2, 'us-west', 20.3), (3, 'eu', 30.7)") + } + assert( + snapshot.snapshotDelta == 1L, + s"expected 1 commit, got ${snapshot.snapshotDelta}") + val (commits, writes) = collectIcebergWriteOps(snapshot.plans) + assert( + commits.isEmpty && writes.isEmpty, + s"expected Spark's own write plan with ${flag.key}=false. Plans:\n" + + snapshot.plans.mkString("\n--\n")) + } + assertRows(table, expectedIds = Seq(1, 2, 3)) + } + } + } + + // The suite pins both write flags, so this test unsets them to see what an application that + // sets neither gets. A Parquet scan is a native input, so the write is eligible. + // https://github.com/apache/datafusion-comet/issues/5644 + test("an eligible Iceberg write runs natively under the split plan when no flag is set") { assume(icebergAvailable, "Iceberg not available in classpath") withIcebergCatalog { warehouseDir => - createTable(warehouseDir, "comet_disabled", partitionSpec = "") - // withSQLConf returns Unit before Spark 4.0, so the assertions run inside it. - withSQLConf(CometConf.COMET_ENABLED.key -> "false") { - val snapshot = captureWrite("comet_disabled") { - spark.sql( - "INSERT INTO cat.db.comet_disabled VALUES " + - "(1, 'us-east', 10.5), (2, 'us-west', 20.3), (3, 'eu', 30.7)") + withTempPath { dir => + spark + .range(10) + .selectExpr("CAST(id AS INT) AS id", "'eu' AS region", "CAST(id AS DOUBLE) AS amount") + .write + .parquet(dir.getCanonicalPath) + createTable(warehouseDir, "write_defaults", partitionSpec = "") + val snapshot = captureWrite("write_defaults") { + withSessionConf( + CometConf.COMET_ICEBERG_WRITE_SPLIT_OPERATOR_ENABLED.key -> None, + CometConf.COMET_ICEBERG_NATIVE_WRITE_ENABLED.key -> None) { + spark.read + .parquet(dir.getCanonicalPath) + .writeTo(s"$catalog.$ns.write_defaults") + .append() + } } assert(snapshot.snapshotDelta == 1L, s"expected 1 commit, got ${snapshot.snapshotDelta}") - val (commits, writes) = collectIcebergWriteOps(snapshot.plans) + val (commits, _) = collectIcebergWriteOps(snapshot.plans) + val nativeWrites = snapshot.plans.flatMap { plan => + collectWithSubqueries(plan) { case e: CometIcebergWriteExec => e } + } assert( - commits.isEmpty && writes.isEmpty, - "expected Spark's own write plan with Comet disabled. Plans:\n" + + commits.nonEmpty && nativeWrites.nonEmpty, + "expected an IcebergCommitExec over a CometIcebergWriteExec. Plans:\n" + snapshot.plans.mkString("\n--\n")) + assertRows("write_defaults", expectedIds = 0 until 10) + } + } + } + + // 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 + // 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") { + 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 = "") + // 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 + } + plans + } + + // 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)) + } } - assertRows("comet_disabled", expectedIds = Seq(1, 2, 3)) } } @@ -890,13 +1044,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)") } } @@ -4135,17 +4291,27 @@ class CometIcebergWriteActionSuite * version / session-state combinations where `withSQLConf` loses the override before the rule * fires. */ - 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") - 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) + 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 + 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 922b5d6073f..fc418ee0808 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") } @@ -1454,10 +1453,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))