Two Spark jobs commit to one Iceberg table at the same moment — what happens?
answer
- nobody holds a lock while writing
- one value decides the winner
- compare against the base you read
- the loser rebuilds metadata, not data
basics
~20 sBoth write their new metadata files, then each asks the catalog to swap the table pointer from the base metadata they read to their own. The swap is an atomic compare-and-swap, so exactly one wins; the loser refreshes, re-applies its changes on the new base, and retries.
solid answer
~50 sIceberg uses **optimistic concurrency on a single atomic pointer**. Each writer reads the current metadata file (its *base*), writes its data files, manifest, manifest list and a new metadata JSON, then asks the catalog to move the table pointer from that exact base to the new file. The catalog performs this as a compare-and-swap — a Hive Metastore parameter update under lock, a Glue conditional update, a REST catalog conditional commit, a JDBC row update — so exactly one writer wins. The loser gets a commit failure, refreshes to the winner's metadata, re-applies its own change on top, and retries, governed by `commit.retry.num-retries` and the related backoff properties. Appends almost always succeed on retry, since adding files conflicts with nothing. Deletes, overwrites and `MERGE` run validations against the snapshot they planned on; if the winner touched the same files or the same partitions, the retry throws a validation error and the job fails rather than silently losing an update.
code
sql · 4 linesALTER TABLE db.events SET TBLPROPERTIES (
'commit.retry.num-retries' = '8',
'write.merge.isolation-level' = 'snapshot'
);go deeper
Recall that only one writer's commit succeeds at a time, because committing means swapping a single pointer, and the other writer retries.
Explain optimistic concurrency: read a base metadata file, write new metadata, conditional swap, and refresh-and-retry on failure without rewriting data files.
Show the judgment: which operation pairs retry cleanly, when validation must fail the job, how isolation level changes that, and the orphan files a failed commit leaves behind.
Own commit-rate capacity across the platform — writer batching, table or branch partitioning of workloads, catalog selection for genuine atomicity, and a cleanup policy whose age threshold cannot race in-flight commits.
## The commit protocol in one paragraph An Iceberg writer never holds a long lock. It reads the table's current metadata file — call that the **base** — plans and writes its data files, then builds a full new metadata file describing base-plus-its-change. The final step is a request to the catalog: *set the table's current metadata pointer to N, but only if it is still B*. That conditional swap is the commit. Everything written before it is invisible; everything after it is visible to every reader at once. There is no partial state. ## Why exactly one writer wins Because the swap is a **compare-and-swap on a single value**, concurrency reduces to a standard CAS race. Each catalog implements it with its own primitive: Hive Metastore updates the table's metadata-location parameter while holding a lock and checking the previous value; AWS Glue uses a conditional update keyed on the version id; a REST catalog performs a conditional commit server-side; JDBC catalogs use a conditional UPDATE row; a Hadoop catalog relies on an atomic rename or a version-hint update, which is exactly why Hadoop catalogs are discouraged on object stores where atomic rename does not exist. The loser observes that the pointer no longer equals its base and receives a commit failure (`CommitFailedException` in the Java implementation). Nothing it wrote is visible; its data files are still on storage but unreferenced. ## What the loser does next It does not simply fail. Iceberg **refreshes and retries**: it reloads the winner's metadata as a new base and re-applies its own operation on top, up to a configured number of attempts with backoff — `commit.retry.num-retries`, `commit.retry.min-wait-ms`, `commit.retry.max-wait-ms` and `commit.retry.total-timeout-ms`. Crucially, the *data files are not rewritten*; only the metadata is rebuilt, so a retry is cheap. Whether the retry succeeds depends on the operation: - **Append vs append** — always compatible. Adding files does not depend on which other files exist, so the retry re-references the winner's manifests and commits. This is why many concurrent streaming writers can share one Iceberg table. - **Append vs delete/overwrite** — usually fine, and sequence numbers make it correct: the appended file has a higher sequence number than the earlier delete, so the delete does not apply to it. - **Delete/overwrite/MERGE vs a conflicting change** — the retry runs the operation's **validation**. If the operation was planned against a snapshot and the winner added or removed files in the affected partitions, the validation throws (`ValidationException`) and the job fails loudly. Failing is the correct outcome: silently committing would lose or double-apply an update. ## Isolation levels Row-level operations expose an isolation setting per operation type — `write.delete.isolation-level`, `write.update.isolation-level`, `write.merge.isolation-level` — with **serializable** as the strict choice and **snapshot** as the more permissive one. Serializable rejects a commit if any conflicting data was added to the affected partitions since planning; snapshot isolation only rejects when the specific files the operation read were changed. Serializable produces more retry failures and fewer surprises; snapshot lets more concurrent work through and can miss rows that arrived mid-operation. ## What this costs operationally - **Orphan files.** Every failed commit leaves written-but-unreferenced data files. They cost storage and are cleaned only by an orphan-file maintenance job, which must use a conservative age threshold so it never removes files an in-flight commit is about to reference. - **Retry storms.** Many writers committing to the same table serialize on one pointer. Throughput is bounded by commit rate, not by data volume, so the fix for a hot table is fewer, larger commits — batching a streaming writer's trigger interval — not more retries. - **Catalog choice matters.** The atomicity of the whole table depends on the catalog's CAS being genuinely atomic. A catalog that cannot do it correctly turns concurrent commits into lost updates, which is why directory-based catalogs are avoided on object storage. ## The sentence to say in an interview "Iceberg serializes commits through one atomic pointer swap in the catalog; writers are optimistic, the loser refreshes and re-applies, and only a genuine logical conflict — validated against the files or partitions the operation touched — makes it fail." Then add the appends-always-retry-cleanly nuance, and the interviewer has what they were listening for.
- Why do concurrent appends almost always succeed on retry while concurrent MERGEs often do not?An append's outcome does not depend on which other files exist, so re-applying it on the winner's metadata is always valid. A MERGE was planned against a specific snapshot and must validate that the files or partitions it read were not changed underneath it. If they were, committing anyway would lose or duplicate updates, so it throws instead.
- What is left on storage after a failed Iceberg commit, and why is that a problem?The data files, manifest and manifest list the loser wrote remain, referenced by nothing. They are invisible to queries but consume storage and object counts. Only an orphan-file cleanup removes them, and it must use a generous age threshold so it never deletes files an in-flight commit is about to reference.
- How does the choice of catalog affect concurrent-commit safety?Entirely — the catalog supplies the compare-and-swap that makes commits atomic. Hive, Glue, JDBC, Nessie and REST catalogs each implement a genuine conditional update. A directory-based Hadoop catalog depends on atomic rename, which object stores do not reliably provide, so concurrent commits there can lose updates.
- A table with many streaming writers commits slowly. What is the structural fix?Reduce commit frequency rather than raising retry counts. All writers serialize on one pointer, so throughput is bounded by commits per second, not bytes. Longer trigger intervals, batching several micro-batches per commit, or splitting the workload across tables or branches all attack the real bottleneck.
saying these in an interview costs you the question
- Says Iceberg locks the table for the duration of a write
- Claims both writers' commits are merged automatically
- Thinks the losing writer must rewrite all its data files
- Believes conflicting MERGEs silently pick a last-writer-wins result
- Assumes any catalog works equally well for concurrent commits