Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
15 commits
Select commit Hold shift + click to select a range
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
19 changes: 11 additions & 8 deletions .ai/skills/review-comet-iceberg-write-pr/SKILL.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
8 changes: 4 additions & 4 deletions docs/source/about/faq.md
Original file line number Diff line number Diff line change
Expand Up @@ -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?
Expand Down
9 changes: 5 additions & 4 deletions docs/source/about/gluten_comparison.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
16 changes: 7 additions & 9 deletions docs/source/contributor-guide/iceberg-spark-tests.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down
45 changes: 25 additions & 20 deletions docs/source/contributor-guide/iceberg-writes.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -70,13 +73,13 @@ IcebergCommit driver: collect task commit messages, BatchWrite.c
+- <input query> 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:

Expand All @@ -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
Expand Down Expand Up @@ -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
Expand Down
16 changes: 8 additions & 8 deletions docs/source/contributor-guide/roadmap.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Loading
Loading