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.
With only the split plan enabled, 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 foundation for the second toggle: when
spark.comet.iceberg.write.enabled=true and the write passes the eligibility check below, the
IcebergWrite operator’s per-task Parquet write is delegated to
iceberg-rust via Comet’s native execution pipeline
(#5308).
How the native write works#
The JVM-side planner marshals everything iceberg-rust needs — the write schema and partition
spec as JSON, the data location, the resolved parquet writer settings, the writer mode
(unpartitioned / fanout / clustered, mirroring SparkWrite’s own choice), object-store
configuration (the table’s FileIO properties, e.g. REST-vended credentials, merged over
fs.s3a.* settings translated from the session Hadoop configuration — the same translation
the native scan uses, since HadoopFileIO carries its S3 configuration in the Hadoop
Configuration rather than in FileIO properties), and per-task IDs — into the serialized
native plan. On each task, iceberg-rust writes the Parquet files and
returns its DataFile metadata packed as a single in-memory Iceberg V2 data manifest; the JVM
decodes those bytes with Iceberg’s own ManifestFiles.read, re-derives each file’s manifest
metrics from the written Parquet footer with Iceberg’s MetricsConfig logic (so metrics modes,
truncation, and bounds decisions are iceberg-java’s by construction), and wraps the result in
the same TaskCommit message the JVM writer would have produced. Everything iceberg-java does
post-write — snapshot assignment, manifest-list aggregation, commit validation and retries —
is untouched: IcebergCommit performs the normal BatchWrite.commit.
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
# Native-write eligibility detection (experimental, off by default; requires the split plan)
spark.comet.iceberg.write.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.
Native Parquet write eligibility#
When spark.comet.iceberg.write.enabled=true
(#5308), the IcebergWrite operator’s
per-task Parquet write is delegated to iceberg-rust.
The native writer must produce the same outcome as iceberg-java — the same Parquet features,
statistics, and manifest metadata — so a write is only eligible when every table property it
depends on is one the native path reproduces exactly, and additionally only when the plan
feeding the write is fully Comet-native. For a partitioned table that plan includes the hash
distribution and local sort Iceberg requests on its partition transforms; those stay native
because the transforms themselves have native implementations (see
Iceberg system functions). Ineligible writes run through iceberg-java unchanged,
with the reason reported as a fall-back reason in Comet’s extended EXPLAIN output.
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
else — any other write-affecting property, any key added by a future Iceberg version, any
value outside the supported set, any reflection failure while inspecting the write — falls
back to iceberg-java with a reason reported in extended EXPLAIN. Checks run on the effective
configuration: table properties overlaid with SparkWrite.writeProperties, which is where
iceberg-java resolves per-write options and spark.sql.iceberg.* session overrides.
A write is eligible only when ALL of the following hold:
Setting |
Supported values |
|---|---|
resolved write format ( |
|
|
|
|
the sizes/limit must be positive Java ints ( |
|
unset or |
|
unset or |
|
unset or |
|
unset or |
|
any value (only meaningful when shredding, which is gated) |
|
unset or |
|
any value (manifest metrics are re-derived on the JVM with Iceberg’s own logic) |
|
any value (the native writer implements both clustered and fanout modes) |
|
any value (the two writers can choose different roll points; see accepted divergences) |
data location URI scheme |
|
partition spec |
any |
column types |
any except |
Within the namespaces that shape data-file bytes — write.parquet.* and parquet.* —
everything not listed above must be absent: unvetted write.parquet.* keys (e.g.
bloom-filter-max-bytes, stats-enabled.column.*, keys added by future Iceberg versions),
any parquet.* table property (including parquet.enable.dictionary), and any parquet.*
key in the session Hadoop configuration (with HadoopFileIO-backed output those reach
iceberg-java’s writer but not the native one). Also gated explicitly: any encryption.* key,
write.object-storage.enabled=true, write.location-provider.impl, and io-impl.
Two checks look past properties at the table’s instantiated state, because both can be
configured at the catalog level (or by a custom TableOperations) without any table or write
property changing: table.io() must be a recognized FileIO (the same allowlist the native
scan uses, minus the EncryptingFileIO family — the native writer produces plaintext files,
so an encrypting FileIO is rejected on the write side), and table.encryption() must be
Iceberg’s PlaintextEncryptionManager. Anything else falls back.
Other write.* properties are intentionally not gated because they cannot make the native
writer produce different data files: distribution and ordering settings shape the Spark plan
identically on both paths, WAP / branch / snapshot properties act on the JVM committer,
write.avro.* / write.orc.* apply only to formats already excluded, and merge-on-read
settings route the write through WriteDelta, which the split plan never intercepts. Every
rule is pinned by CometIcebergWriteDetectionSuite.
Manifest DataFile metrics are assembled on the JVM before commit: each written file’s
metrics are re-derived from its parquet footer through the version-matched
ParquetUtil.footerMetrics and MetricsConfig.forTable, with float/double NaN counts and
bounds carried over from the native writer’s tracked state. iceberg-java’s metadata decisions
— 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.
Failure handling#
Eligibility is decided entirely at plan time. That includes the reflection surface: every iceberg-java class, method, and constructor the executor-side commit-message assembly uses is eagerly resolved by the eligibility gate on the driver, so an Iceberg release that moves any of them declines the native path with a fall-back reason instead of failing tasks mid-write. Once planned, the physical plan is fixed — there is no per-task re-decision or runtime switch back to the JVM writer.
When a native write fails partway through a task (an object-store error, a data-dependent cast failure), the error propagates as an ordinary Spark task failure and Spark’s task retry re-executes it — through the native writer again. Retries cannot collide: each attempt’s task attempt id is embedded in its data file names.
Partial results are never committed. The commit set is exactly the commit messages returned by
successful tasks — a failed task contributes none — and if the job fails, the driver-side
commit operator aborts without committing anything. A failed task attempt also deletes the
data files it created, as iceberg-java’s writer abort does. The native writer records every
location it hands to a file writer, and exactly one side owns deleting them at any moment: the
native writer owns them until its output batch reaches the JVM (so it cleans up a failed write,
a task torn down before the write completed — for example because the operator feeding it threw
— and a failure encoding the manifest or building that batch), and the JVM owns them from then
on through a task failure listener. The handoff does not depend on decoding the manifest: the
native side reports the locations in the output batch alongside it, and the listener is handed
them before the manifest is decoded, so a failure in that decode still cleans up. Both
deletions are best-effort and never mask the original failure; anything they miss is invisible
to every reader, since readers resolve files through committed manifests only, and is reclaimed
by Iceberg’s normal remove_orphan_files maintenance.
When one task fails, the tasks that had already completed leave committed-nothing data files
too. The committer collects each task’s commit message as that task finishes, so on a job
failure it aborts with the completed messages and deletes their data files through the table
FileIO. (Iceberg’s own SparkWrite.abort skips cleanup unless a commit failed with a
cleanable error, so on the stock path those files are left for remove_orphan_files.) A
failure during the driver-side commit itself behaves exactly as on the stock path: the commit
messages carry genuine SparkWrite$TaskCommit objects, so Iceberg’s own SparkWrite.abort
cleanup (which deletes the files listed in the commit messages for cleanable failures) applies
unchanged.
Accepted divergences behind the toggle#
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.iceberg.write.enabled accepts 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
below), and one operational path-layout caveat.
Physical file layout only (cosmetic)#
No Iceberg reader bases a planning or correctness decision on these; they change the bytes of a data file but not what any reader computes from it:
Footer key-value metadata differs: native files carry an
ARROW:schemaentry and noiceberg.schemaentry; iceberg-java files are the opposite.The Parquet root schema element is named
arrow_schema(iceberg-java:table).created_byidentifies parquet-rs, not parquet-mr.No page CRC checksums and no page-header statistics (parquet-mr writes both by default). Absent page-header statistics can only make a reader scan more pages, never skip pages it should have read; page pruning uses the column index, which the native writer does produce.
Dictionary-encoded pages are labeled
RLE_DICTIONARY(parquet-mr v1 files:PLAIN_DICTIONARY).Fixed-length binary columns (
uuid,fixed, decimals with precision > 18) are not dictionary-encoded (parquet-mr dictionary-encodes them).Row-group boundaries: parquet-mr flushes by byte size at a record-count check cadence, parquet-rs buffers by row count. File naming follows the same cadence-style difference (iceberg-java names files
<partition>-<task>-<operation>-<count>; iceberg-rust uses a process-local counter).Partition directory names match iceberg-java 1.8+’s
PartitionSpec.partitionToPathfor every partition type exceptfloatanddouble, where the value is rendered with Rust’s shortest representation instead ofFloat.toString/Double.toString(f=1where iceberg-java writesf=1.0). On Iceberg 1.5.x, which the Spark 3.4 profile pins, iceberg-java itself spelledtimestampandtimestamptzdirectories withLocalDateTime.toString()/OffsetDateTime.toString()(ts=1969-12-31T23:59:58.500Z) and left the partition field name unescaped; Comet uses the 1.8+ spelling on every profile. Distinct partition values still get distinct directories in all cases, and no reader parses these names — files are resolved through committed manifests. Iceberg deprecated float and double partitioning in 1.3.File rolling lands on the same row grid as iceberg-java but not necessarily on the same row. Both writers re-check the current file’s size against
write.target-file-size-bytesonce every 1000 rows of that file (iceberg-java’sRollingFileWriter.ROWS_DIVISOR; Comet hands the iceberg-rust writer rows in 1000-row units, per partition file, to get the same grid), so each writer rolls only on a 1000-row boundary of its own file. The shared grid is all that is shared. What each writer compares against the target differs — flushed bytes plus parquet-rs’s estimate of the open row group, versus parquet-mr’s file position plus its buffered size — and the two use different threshold comparisons. These are independent size estimates, so nothing bounds how far apart the two writers’ roll points are: they may cross the target several grid steps apart, and the resulting files can differ in row count by an arbitrary number of 1000-row blocks. Do not rely on file-layout parity between the two writers; rely only on each file rolling on its own 1000-row boundary.A fanout write lists a task’s data files in file-path order, where iceberg-java lists them in its own
StructLikeMapiteration order. Both are stable across runs, and neither is a documented ordering, but the manifest entry order becomes the scan-task order and so the row order of an unorderedSELECT *. Only the sorted order is reproducible on the native path: iceberg-rust’sFanoutWritercloses its per-partition writers out of aHashMap, which under Rust’s per-processRandomStatewould otherwise give a different order on every run. Clustered and unpartitioned writes append in creation order on both paths and are unaffected.Compressed page bytes are implementation-defined: the codec and any explicit level are translated, but parquet-rs and parquet-mr embed different encoder implementations and defaults (zstd default levels, LZ4 framing), so byte-identical output is not achievable even for a default
zstdtable. The decompressed data is identical. For the same reason, codec-level side channels (zlib.compress.level,compression.brotli.quality,io.compression.codec.zstd.level— the last is present in every Hadoop configuration by default) are not gated: they can only shift compressed bytes, which are already accepted as divergent.
Manifest metadata visible to later readers#
Manifest metrics drive partition- and file-level pruning for every future reader of the table,
so a divergence here would outlive the write. This class is deliberately kept almost empty:
DataFile metrics are not taken from the native writer’s manifest but re-derived on the JVM
from each written file’s parquet footer through iceberg-java’s own ParquetUtil.footerMetrics
and MetricsConfig.forTable (see above). Metrics modes, lower/upper bound truncation, the
null-count conventions, and list/map bounds suppression are therefore iceberg-java’s code
making iceberg-java’s decisions, and the parity suite compares committed manifests
byte-for-byte against JVM-written ones. Two footer-derived values can still differ from what
iceberg-java’s writer-tracked state would have recorded, and both are analyzed safe:
Float/double bounds involving zero may differ in sign: parquet-rs normalises footer statistics to min
-0.0/ max+0.0(the parquet-format recommendation), while iceberg-java’s writer-tracked bounds preserve the exact sign it saw. The native path’s manifest bounds inherit the normalised values — a strictly conservative widening that cannot change pruning decisions.On Iceberg 1.10+, manifest
value_counts/null_value_countsfor float/double columns nested under a nullable struct count rows whose parent struct is null (they come from the parquet footer), while iceberg-java’s writer-tracked counts do not. Both counts inflate by the same amount, so the derived null ratios andIS NULL/IS NOT NULLpruning decisions are unaffected.
All content not listed above — the logical data, encodings for non-FLBA columns, statistics values, and manifest metadata — must match iceberg-java exactly, or the write falls back.