Iceberg Writes: Comet’s Split-Operator Plan (Experimental)#
This feature is experimental and disabled by default. Enable it only after validating it against your own workloads.
Overview#
Spark writes an Iceberg table through a single physical operator that combines data-file writing with metadata writing, committing, and catalog validation. Because that operator sits outside Spark’s Adaptive Query Execution (AQE), the sub-query feeding the write — the scans, projects, sorts, and exchanges producing the rows — cannot be re-planned at runtime.
When spark.comet.write.iceberg.splitOperator.enabled=true, Comet rewrites eligible Iceberg
writes into two operators:
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 feeding it run inside AQE.IcebergCommit— collects the commit messages on the driver and performs the normal Iceberg commit (including commit-time validation), outside AQE, exactly once.
Data files are still written by iceberg-java; only the plan shape changes. The split makes the write’s input visible to AQE and to Comet’s columnar rules, and it is the groundwork for a planned follow-up in which Comet writes the data files natively via iceberg-rust.
Configuration#
Standard Comet + Iceberg setup (see iceberg.md) plus the write-side toggle:
# Standard Comet / Iceberg wiring
spark.plugins=org.apache.spark.CometPlugin
spark.sql.extensions=org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions
spark.sql.catalog.<name>=org.apache.iceberg.spark.SparkCatalog
spark.sql.catalog.<name>.type=hadoop # or hive / glue / rest / ...
spark.sql.catalog.<name>.warehouse=...
# Split-operator plan (experimental, off by default)
spark.comet.write.iceberg.splitOperator.enabled=true
Supported operations#
The split-operator plan supports the following operations on every Spark version Comet supports:
INSERT INTO/ DataFrameappend(AppendData)INSERT OVERWRITE, static and dynamic (OverwriteByExpression,OverwritePartitionsDynamic)Copy-on-write
DELETE/UPDATE/MERGE(ReplaceData)
The mechanism behind row-level DML differs by Spark version: on Spark 4.0+ the analyzer emits
operation-coded rows that Comet’s writer dispatches through ReplaceData’s 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 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 cached catalog tables are recached by name after a write so cache entries survive schema changes.
When Comet falls back to Spark’s write operator#
The rewrite is skipped — and the write runs through Spark’s stock combined operator — when:
spark.comet.write.iceberg.splitOperator.enabledisfalse(the default);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 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 not use;
Comet cannot reflect the Iceberg internals needed to build the two-operator plan (for example an unrecognised write class or a
ReplaceDataprojection it cannot map).
In every fallback case the write is planned as if Comet were absent; there is no correctness trade-off, only no plan change.