TT Lab
Get started
Learn Learning paths Courses

Lakehouse Table Format — Understanding Apache Iceberg Through Its Metadata

Writing concurrently without locks — what optimistic concurrency protects, and what it doesn't

Continue in TT Lab

In one line

Iceberg writers, without locks, each build a new metadata and then try a conditional swap in the catalog, and if they lose, they read the new state and layer their change on top again. Changes that can always be layered on again, like appends, are resolved automatically, while changes that may overlap with what happened in between, like overwrites, are stopped by validation. But a result that your code computed from a value it read earlier is something nobody validates for you.

Why locks are not used

Usually several jobs write to one table. A streaming load appends every few minutes, a nightly batch fixes yesterday's partition, and a maintenance job compacts small files. If you lock the whole table, the slowest job stops everyone. On top of that, there is no dependable distributed lock on object storage.

The optimistic concurrency section of the spec takes a different path. A writer builds its metadata assuming that current will not change before its own commit, and then swaps the pointer from the base version to the new version. If the base snapshot is no longer current, it starts over based on the new current. Readers see the snapshot they opened all the way through, so there is no need to lock.

How it works — assumptions and actions

The Reliability docs describe a commit as an assumption and an action. When there is a conflict, the writer checks whether its assumption still holds in the current state, and if it does, it applies the action again and commits. The conflict resolution section of the spec defines that assumption for each operation.

Operation What is checked before layering again
append Nothing — it can always be layered again
replace (compaction, etc.) Whether the files it meant to delete are still in the table
delete (specific files) Whether the files it meant to delete are still in the table
Schema or spec change Whether the schema changed in the meantime

An append can always be layered again because the action "add new files" remains correct no matter what came in meanwhile. The docs also say the cost of retries has been reduced — an append writes its new manifest once and does not rewrite it on every attempt.

The number and spacing of retries are table properties. commit.retry.num-retries defaults to 4, commit.retry.min-wait-ms to 100, commit.retry.max-wait-ms to 60000, and the overall limit commit.retry.total-timeout-ms to 30 minutes. When the retries are used up, a commit failure exception is raised.

Isolation levels — what counts as a conflict

The spec says "which conditions you validate is the isolation level". Spark's DELETE, UPDATE, and MERGE choose it with the table property write.delete.isolation-level (update and merge have the same name), and the default is serializable. Serializable fails if anyone added or deleted something in the range I read and changed in the meantime, while snapshot counts only deletions as conflicts. The DataFrame overwrite in the Spark write options has the same two levels. The looser the level, the less it fails, but you may overwrite without knowing about rows that came in meanwhile.

What the library does not guard for you

Think of a counter. Two workers read hits = 10, each add 5, and overwrite with 15. The one that goes first wins. The other loses the conditional swap, and its retry is stopped by the validation "a file matching the condition I was going to overwrite (name = hits) came in meanwhile". Up to here, the library does the work for you.

The problem comes next. If you catch the exception and write the 15 you computed earlier as it is, the commit succeeds this time. The table says 15 and one side's +5 has disappeared — a lost update. The library only validates file-level assumptions; it does not know which read your value came from. The correct retry is to read anew, compute again, and commit again on the premise of that read.

What it looks like in the field

Compaction and loading collide. If the files compaction (replace) meant to delete are still there, the compaction is layered again even if a load slips in. Conversely, if a MERGE modified the same files, the compaction fails and you just do it again in the next cycle.

Traces of a failed commit. The data files that a writer that lost the commit already wrote remain without entering any snapshot. They do not affect the table but take up space, and clearing such files is orphan file cleanup.

Retries set to 0. If you turn off retries believing there is exactly one job writing to a table, on the day it collides with the maintenance job that runs now and then, the job fails for no apparent reason.

What really matters in practice

What you will do in the next lab

You make two writers as two Table objects in one process with pyiceberg, have them read the same metadata and then append in turn, and see that the losing side is layered again automatically so the snapshots line up in one chain. You turn retries off to run the same race again and find the exception name and the orphan files the losing side left. Finally, with a single counter you race overwrites to see the validation exception, then read anew and recompute so that the counter ends up exactly 20.