Iceberg Writes#
This document describes how Comet writes Iceberg tables. It covers the split-operator plan that replaces Spark’s single write operator, how a write is admitted to the native writer, what crosses the JVM/native boundary in each direction, who owns cleanup when a task fails, and how to test a change to any of it.
For the user-facing view (configuration, the eligibility table, and the accepted differences from iceberg-java’s output), see Iceberg Writes. This page 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 |
|---|---|
|
The plan shape. Spark’s single V2 write operator becomes |
|
Who writes the data files. An eligible |
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 under the epic
#5649.
One rule runs through the whole native path: the native writer must produce the outcome iceberg-java would have produced, or decline. iceberg-java is the reference for every data file, manifest entry, partition value and failure behaviour. Where Comet cannot show that it reproduces iceberg-java for some table or setting, the write falls back to iceberg-java at plan time. The consequences of getting this wrong are asymmetric. A decline costs performance on one write. A native write that differs from iceberg-java writes a table that every later reader, Comet or not, has to live with.
A second rule follows from the first: decide at plan time, never mid-write. Once the plan
contains a CometIcebergWrite, there is no runtime switch back to the JVM writer. Anything the
native side could reject at execution time (a storage scheme, a type, a partition spec, a
reflective accessor the executor needs) has to be rejected by the planning gate first, or the user
gets a failed query instead of a fallback.
The Split-Operator Plan#
Spark plans an Iceberg write as one physical operator (AppendDataExec, ReplaceDataExec and so
on) that writes the files and commits. It is a V2CommandExec, so Spark’s
InsertAdaptiveSparkPlan wraps its input query in AQE but leaves the operator itself outside. The
split plan replaces it with two operators so that data-file writing moves inside AQE, apart from
the commit.
IcebergCommit driver: collect task commit messages, BatchWrite.commit; outside AQE
+- IcebergWrite executors: write data files, emit one commit message per task; inside AQE
+- <input query> scans, projects, exchanges, sorts; inside AQE with or without the split
Component |
Location |
Role |
|---|---|---|
|
|
Planner strategy. Matches |
|
same |
Logical anchor for the writer, so AQE re-plans re-emit only the writer and not a second committer. |
|
|
JVM writer. Runs iceberg-java’s |
|
same |
Driver committer. A |
|
|
Version differences: Spark 4.x operation-coded |
Things to know before changing this layer:
One
BatchWriteis shared.buildTwoOpcallswrite.toBatchonce and hands the same instance to both operators. Iceberg’s commit-time validation must see the instance the writer wrote through, andtoBatch()returns a new one per call.The commit message column is the contract between the two operators. Both
IcebergWriteExecandCometIcebergWriteExecemit a single non-nullBINARYcolumn namediceberg_commit_message, holding a Java-serializedWriterCommitMessage.IcebergCommitExecdoes not know which writer produced it.Messages are collected per task as tasks finish (
sparkContext.runJobwith a result handler), not withexecuteCollect. That is how a failed job still knows which tasks completed, so it can delete their files.Commit-coordinator writes are not intercepted. Iceberg’s
SparkWritenever asks for one; the check inbuildTwoOpis defensive.What is not intercepted: merge-on-read (
WriteDelta), streaming writes, and CTAS/RTAS on Spark 3.4. Those keep Spark’s plan.
From IcebergWrite to CometIcebergWrite#
CometExecRule converts an IcebergWriteExec with the CometIcebergNativeWrite operator serde
when spark.comet.iceberg.write.enabled is on. Two arms in CometExecRule handle it: one unwraps
the double conversion AQE can produce when it re-fires write planning over a sub-tree that already
contains a CometIcebergWriteExec, and the other calls convertToComet.
CometIcebergNativeWrite.requiresNativeChildren is true. The native writer consumes Arrow
batches from its child over FFI, so the conversion is declined unless the child is already a Comet
native operator. This is why a write fed by a LocalTableScanExec (INSERT ... VALUES, a local
DataFrame) stays on the JVM writer unless spark.comet.exec.localTableScan.enabled is also set.
CometIcebergWriteExec is row-based (it emits the commit message row) over a columnar child, so
Spark inserts a columnar-to-row transition beneath it. EliminateRedundantTransitions removes that
transition in a dedicated pass that runs before its main transformUp, and the scaladoc on
stripIcebergWriteInputTransition explains why the order matters. Do not try to suppress the
transition by making the write a ColumnarToRowTransition: Spark then skips the whole subtree
below it and never inserts the transitions the rest of the plan needs
(#5689, visible only with AQE off).
CometIcebergWriteExec carries its own serialized plan (serializedPlanOpt serializes nativeOp
on demand), so CometExecRule resets firstNativeOp at the write and lets its child start a
separate native block.
The Eligibility Gate#
CometIcebergNativeWrite.getSupportLevel evaluates an ordered list of TriggerRules against a
TriggerContext and reports the first reason it finds as Unsupported. A rule is a function from
the context to Option[String]: None means the rule has no objection.
The gate has these properties, and a new rule has to keep them:
It is an allow list. It accepts only configurations shown to produce iceberg-java’s output. Unknown
write.parquet.*keys, anyparquet.*table property, and anyparquet.*key in the session Hadoop configuration are declined, including keys a future Iceberg version adds. When you support a new setting, add it to the vetted set and translate it; do not widen a prefix.It reads the effective configuration.
TriggerContext.propertiesis the table’s properties overlaid withSparkWrite.writeProperties, which is where iceberg-java resolves per-write options andspark.sql.iceberg.*session overrides. Reading table properties alone misses those.It checks instantiated state as well as properties. A catalog or a custom
TableOperationscan install aLocationProvider,FileIO, orEncryptionManagerwithout any property changing, sorequireDefaultLocationProvider,requireRecognizedTableFileIO, andrequirePlaintextEncryptionManagerinspecttable.locationProvider(),table.io(), andtable.encryption()themselves.It fails closed. A reflection lookup that cannot answer returns a reason, not
None, andgetSupportLevelturns any non-fatal exception intoUnsupported.convertdoes the same for failures while building the proto.It covers what the executor will need.
requireExecutorReflectionResolvableresolves every iceberg-java class, method and constructor that the executor-side commit-message assembly calls reflectively. An Iceberg release that moves one of them becomes a plan-time fallback rather than a task failure after the data is written.It agrees with the native side. When the gate and the native code both interpret the same input (a location’s scheme, a partition spec, a column type), they must interpret it the same way. A gate that is more permissive than the native code turns a fallback into a failed query (#6140 is an example: the gate and
scheme_ofsplit a location’s scheme differently). Prefer sharing one implementation, or pin the pair with a test that feeds both the same inputs.
Every rule is pinned by CometIcebergWriteDetectionSuite. Add a case there for any new rule, for
both the accepted and the declined side. Update the eligibility table in the user guide in the same
change. Whether each existing restriction is permanent is tracked in
#5643.
What Crosses the Boundary#
JVM to native: the IcebergWrite proto#
buildIcebergWriteProto builds an IcebergWrite message (native/proto/src/proto/operator.proto)
on the driver. Almost all of it is the per-write IcebergWriteCommon:
the write schema and partition spec as Iceberg JSON (
SchemaParser/PartitionSpecParser), taken from theSparkWriterather than the table, so a concurrent schema change cannot alter what this write producesthe data location, the operation id (the
SparkWritequery id, used in file names), the target file size, and the writer mode (unpartitioned, fanout or clustered, resolved the waySparkWritechooses its own writer)IcebergParquetWriteSettings, translated from the effective properties byIcebergWriteProtoTranslationthe sort order id, which the native side ignores and the JVM stamps onto the files afterwards
catalog_propertiesfor the nativeFileIO: the table’sFileIO.properties()merged over thefs.s3a.*settings translated from the Hadoop configuration, the same translation the native scan uses
The per-task partition_id and task_attempt_id are stamped onto a copy of the proto inside the
task closure in CometIcebergWriteExec.doExecuteColumnar. The native side refuses to run without
them, since defaulting them would make every task write the same file names.
catalog_properties can carry credentials. CometIcebergWriteExec.stringArgs deliberately prints
only the data location and writer mode, because the default argString would put the protobuf’s
text dump, secrets included, into explain(), the SQL UI and the event log. Keep it that way when
adding fields.
On Spark 4.x, copy-on-write DELETE/UPDATE/MERGE rows arrive with an operation code and file
metadata columns around the data columns. dropNonDataColumns inserts a native Projection that
keeps only the write schema’s columns. This is only equivalent to the JVM writer while format
version 3 is declined, because on v3 iceberg-java reads row-lineage fields from those metadata
columns.
Native execution#
The planner builds IcebergWriteExec (native/core/src/execution/operators/iceberg_write.rs) over
the FFI scan of the child’s batches. run_write_task builds iceberg-rust’s writer stack once per
task:
ParquetWriterBuilder -> RollingFileWriterBuilder -> DataFileWriterBuilder
-> UnpartitionedWriter | FanoutWriter | ClusteredWriter
Points where Comet adapts iceberg-rust to match iceberg-java:
Location generation.
CometLocationGenerator(iceberg_partition_path.rs) replaces iceberg-rust’sDefaultLocationGeneratorand renders partition directories the way iceberg-java’sPartitionSpec.partitionToPathdoes. iceberg-rust’s own rendering differs for several types, and panics for pre-1970timestamptzand for dates pastchrono’s calendar. Float and double values use Comet’s JavaDouble.toStringrenderer (java_float_string), shared withcast(float as string). The partition type is resolved once when the generator is built, becauseLocationGenerator::generate_locationcannot return an error.Partition values.
PartitionValueCalculator(iceberg_partition_value.rs) replaces iceberg-rust’s calculator of the same name, andPartitionSplitterreplaces itsRecordBatchPartitionSplitter.yearandmonthgo through Comet’siceberg_years/iceberg_monthskernels, the ones the sort in front of a clustered write runs. iceberg-rust computes them with Arrow’sdate_part, which returns NULL pastchrono’s calendar (#6145). Every other transform stays on iceberg-rust.File names.
file_name_prefixembeds the partition id, the task attempt id and the operation id, so a retried or speculative attempt never reuses another attempt’s file names.Row pacing. iceberg-java’s rolling writer checks the target file size every 1000 rows of the current file. iceberg-rust checks once per
writecall.RowPacerhands the writer rows inROWS_DIVISOR(1000) row units per file, so both writers roll on the same grid.Field ids and casting.
decorate_batch_with_field_idscasts each batch to the field-id-annotated Arrow schema derived from the Iceberg schema, withsafe: false, so a type mismatch fails the task instead of writing NULLs.Deterministic output order. iceberg-rust’s
FanoutWritercloses its writers out of aHashMap, so the fanout path sorts itsDataFiles by path before returning them (#5776). Manifest order becomes the row order of an unordered read, so any map iteration that reaches the output needs the same care.
FileIO comes from load_file_io in iceberg_common.rs, shared with the native scan. It picks
the storage backend from the data location’s scheme and wires in Comet’s S3 credential bridge when
one is configured. For writes the bridge fails closed: if a configured provider cannot initialize,
the task fails rather than writing with the default credential chain.
Native to JVM: the task payload#
Each task emits exactly one batch with one row and two BINARY columns (build_output_schema):
The task’s
DataFiles encoded as an Iceberg v2 data manifest, written by iceberg-rust’sManifestWriterinto an in-processmemory:FileIO. Empty when the task wrote no files.Every location the task’s writers were handed (
encode_locations: a big-endian count, then a length-prefixed UTF-8 string per location).
CometIcebergWriteExec.doExecute then, per task:
takes cleanup ownership of the locations from column 2 (see below)
decodes the manifest with Iceberg’s own
ManifestFiles.read(decodeManifestToDataFiles)rebuilds each
DataFile’s metrics with iceberg-java’sParquetUtil.footerMetricsandMetricsConfig.forTable, reading each written file’s footer (rebuildDataFilesWithJavaMetrics). Only float/double NaN counts and bounds are carried over from the native writer, because the footer does not have them.stamps the write’s sort order id, which iceberg-rust does not set
builds a genuine
SparkWrite$TaskCommitreflectively, reports Spark output metrics the way iceberg-java’sTaskCommitdoes, and serializes it as the commit message
Rebuilding metrics on the JVM is the main reason manifest parity holds: metrics modes, truncation, the inferred-column cap and list/map bounds suppression are decided by iceberg-java’s code. When touching the native writer, do not try to make iceberg-rust’s own metrics authoritative; the JVM discards them. The price is one footer read per written file.
Failure Handling and Cleanup Ownership#
A failed attempt must leave no data files behind, as iceberg-java’s DataWriter.abort() does, and a
failed job must not commit anything. The files are always owned by exactly one side:
Phase |
Owner |
Mechanism |
|---|---|---|
Writing, closing writers, encoding the manifest, building the batch |
Native |
|
After the batch reaches the JVM, until the task succeeds |
JVM task |
|
Job failure after some tasks completed |
Driver ( |
Calls |
Commit failure |
Iceberg |
|
The handoff between the first two rows is why the payload carries the locations separately from the
manifest: a failure decoding the manifest would otherwise lose the list of files to delete. All
deletion is best effort and logged. It must never replace the original exception, and anything it
misses is unreferenced and reclaimed by Iceberg’s remove_orphan_files.
The same planning-time rule applies to failures: a failure that happens because the native side rejected something the gate admitted is a bug in the gate, even if cleanup works.
Matching iceberg-java#
Changes to the native path are measured against iceberg-java’s output in three tiers, which are the same tiers the user guide’s “accepted divergences” section uses:
Logical content and manifest metadata must match. Row data, partition values, record counts, metrics, spec ids and sort order ids. A difference here outlives the write, because later readers prune on it. There are two documented exceptions, both analyzed as unable to change a pruning decision.
Physical file layout may differ where no reader bases a decision on it: footer key-value metadata,
created_by, encodings, compressed bytes, row-group and file roll points, file names, and fanout file order.Anything else is a fallback. If the native writer cannot match tier 1 for some table or setting, add a gate rule rather than documenting a new divergence.
A new accepted divergence belongs in the user guide’s list with the reasoning for why no reader can observe it, not only in a code comment.
Useful places to look when checking parity:
CometIcebergWriteActionSuitewrites the same data through both writers into sibling tables and compares rows,readable_metricsand partition paths.iceberg-rust’s own code, at the pinned revision. iceberg-rust makes different choices from iceberg-java in places that matter to the table’s contents, for example grouping partition keys with
OrderedFloat, which treats-0.0and0.0as equal (#6138). Check how iceberg-rust compares, hashes and renders values, not only what it writes.
The iceberg-rust pin#
native/Cargo.toml pins iceberg and iceberg-storage-opendal to a git revision, not a release.
Some Comet code depends on iceberg-rust behaviour that is not a stable API. For example,
clustered_write_err recognizes the clustered writer’s unsorted-input error by its message text and
restates it as iceberg-java’s error, which applications match on. A pin bump can change file bytes, manifests or partition layout, so treat it as a change to
the writer: run the write suites and the Iceberg Spark tests. The pin policy is tracked in
#5645.
Testing#
Suite |
What it covers |
|---|---|
|
End-to-end writes through the split plan and the native writer: parity with iceberg-java, row-level DML, partition evolution, file order, cleanup on task and job failure, AQE re-planning. |
|
One case per eligibility rule, accepted and declined. |
|
Native |
|
Iceberg’s |
|
Translation of properties into |
Rust tests in |
File rolling on the 1000-row grid, fanout order, clustered input checks, cleanup guard, manifest round trip, partition path rendering, partition values past |
|
Native versus iceberg-java for unpartitioned, clustered, fanout and copy-on-write delete writes. It checks each arm’s plan before timing it. |
The Comet suites run against the Iceberg version each Spark profile pins in spark/pom.xml: 1.5.2
for Spark 3.4, 1.8.1 for 3.5, 1.10.0 for 4.0 and 4.2, and 1.11.0 for 4.1. Only the default profile
runs on every pull request. Behaviour that differs between Iceberg versions, such as reflection
targets, metrics conventions and partition path spelling, needs a test that runs on each.
When writing a native-write test:
Assert that the native writer ran. Every fallback is silent to the query, so a passing comparison can mean both sides used iceberg-java. Use
assertNativeWriteEngages, or collectCometIcebergWriteExecfrom the captured plans. Enable the writer withwithNativeEnabled, which also turns onlocalTableScansoINSERT ... VALUESinput is native.Compare against iceberg-java, not against expected literals. Write the same rows through the JVM writer into a sibling table and compare.
Run plan-shape and transition changes with AQE off as well as on. The suites run with AQE on by default, and several write-path bugs only appeared with it off. Iceberg’s own extension tests pick AQE on or off at random per session, so a bug that depends on AQE shows up there as intermittent.
Use enough partitions to catch ordering bugs. With two partitions, a random order is right half the time. The fanout order test uses eight.
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
Running Iceberg Spark Tests). 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
(#6148).
Pitfalls#
Each of these has caused a bug on this path:
A permissive gate is a failed query. Anything the native side rejects at execution has to be rejected at plan time, using the same interpretation of the input.
Dropped settings are silent. A table,
FileIOor Hadoop setting the native writer does not read produces a successful write that ignores it. Fail closed on unknown settings in every namespace that reaches the writer, as thegs://path does (#5637, #6139).Map iteration order leaks. It reaches file names, manifest order and read order.
Values that compare equal in Rust may not in Java. Signed zeros and NaN under
OrderedFloat, and string or float rendering in partition paths.chronostops at year 262142. A Spark date reaches year 5881580 and a timestamp year 294247, and iceberg-java handles all of them. Code that goes throughchrono, including Arrow’sdate_part, panics or returns NULL past that (#6145).Partition evolution leaves
voidfields behind. A v1 spec keeps a dropped partition field as avoidtransform, whose source column may later be dropped from the schema. Resolving the spec against the schema then fails (#5691, #5693, #6141).Plan rewrites must keep the write node. Rules that restore Spark operators from a Comet node’s
originalPlanhave to handle the write execs, whoseoriginalPlantoday is their child (#5719).The kill switch must still work. Code that runs for Iceberg writes has to respect
spark.comet.enabled, so that disabling Comet restores Spark’s own plan (#6142).