Lesson 3 · Open table formats

The commit, and time travel

Two writers, one pointer, no locks. How Iceberg resolves that — and why the answer is also the reason you get time travel for free.

~10 min · grounded in the Iceberg spec, § Optimistic Concurrency and § Sequence Numbers

3 of 5

1.Nobody takes a lock

The whole protocol fits in one sentence from the spec: “An atomic swap of one table metadata file for another provides the basis for serializable isolation.”1

Writers do not coordinate in advance. Each one assumes it will win, does all its work, and finds out at the last instant whether it did.1 Watch what happens when two of them are wrong about that at the same time.

CATALOG WRITER A WRITER B v7 both read v7 builds v8 base = v7 builds v8 base = v7 commit v8 × rejected base v7 is stale refresh → v8 revalidate, retry v9 Both writes land. B never rewrote a data file — only its metadata.
  1. Both writers load the table at v7. Neither knows the other exists, and neither takes a lock.
  2. Both do their real work: write data files, build a candidate metadata file. Both candidates are based on v7.
  3. A commits first. The swap checks that the current version is still v7, and it is, so v8 becomes current.
  4. B attempts the same swap. Its base version is no longer current, so the swap is rejected. Nothing is corrupted; the commit simply did not happen.
  5. B refreshes to v8 and re-checks its assumptions against the new state.
  6. Assumptions hold, so B re-applies its changes on top of v8 and commits v9. Its expensive work — the data files — was never thrown away.
“If the snapshot on which an update is based is no longer current, the writer must retry the update based on the new current version.”1

2.Retry validation: assumptions and actions

Beat 5 is where the design gets interesting. A naive retry would redo everything. Iceberg instead structures a commit as assumptions plus actions: after a conflict, the writer checks whether the assumptions still hold against the new state; if they do, re-applying the actions is safe.2

The canonical example, from the docs

A compaction rewrites file_a.avro and file_b.avro into merged.parquet. That is safe to commit as long as the table still contains both of the original files.2

So a concurrent append does not invalidate the compaction — the append added files, it did not remove the two being merged. But a concurrent delete of file_a.avro does.

This is also where isolation levels come from. The spec puts it plainly: “The conditions required by a write to successfully commit determines the isolation level. Writers can select what to validate.”1 Serializable and snapshot isolation are not two engines — they are two choices about how much to check on retry.

What this means for you in practice

Optimistic concurrency is excellent when writers touch different files and poor when they fight over the same ones. Many small concurrent writers to one hot partition will livelock on retries. The fix is almost never “tune the retry count” — it is to partition the write, or to funnel it through one writer.

3.Sequence numbers: how deletes find their data

Every successful commit is assigned a monotonically increasing sequence number, and it is the mechanism that makes row-level deletes work without rewriting data.1

Assignment

  • A snapshot is optimistically assigned the next sequence number.
  • Commit fails and retries? It is reassigned.
  • All manifests, data files and delete files created for that snapshot inherit it.

Why inheritance

  • New entries are written with null in place of a sequence number.
  • It is filled in from the manifest’s metadata at read time.
  • So a retry only rewrites the manifest list — which was being rewritten anyway.

The payoff shows up in scan planning. A delete applies to a data file only when the data file’s sequence number is less than or equal to the delete’s.1 That single comparison is how a delete committed at sequence 42 knows to affect everything written before it and nothing written after — with no rewrite of either side.

4.History is not a feature, it is a side effect

Every commit writes a new metadata file and keeps the old snapshots listed. Time travel is what you call it when you read one of the ones you didn’t delete.

SNAPSHOT OPERATION LIVE DATA FILES READ IT WITH S1 append a b VERSION AS OF 1 S2 append a b c VERSION AS OF 2 S3 delete a b c VERSION AS OF 3 S4 rewrite merged the current table same rows expire_snapshots(older_than = S3) S1 and S2 leave the history · only now can their unreferenced files be deleted · time travel to them is gone
  1. S1 appends two files. The snapshot records the complete set of live files — not a diff.
  2. S2 appends a third. S1 still exists and still names exactly a and b.
  3. S3 deletes rows in b. Whether b is rewritten or merely masked by a delete file is a table setting, not a format requirement.
  4. S4 compacts. Same rows, fewer files. Because it is just another snapshot, a bad compaction is one rollback away from undone.
  5. Retention is the cost side. Expiring snapshots is what finally lets storage be reclaimed — and it is also what destroys the ability to travel back to them.
Branches and tags are the same idea, named: a tag labels one snapshot, a branch is a mutable named reference updated by committing to it.1

The operational habit this buys

“Roll back the table” becomes a metadata operation instead of a restore from backup. That changes how brave you can be with a pipeline — but only if snapshot retention is long enough that the snapshot you want still exists. Retention policy is your recovery window.

5.Who actually performs the swap

Here is the gap Lesson 1 flagged. The spec defines everything about the metadata — and then declines to define the one operation that makes it work: “The atomic operation used to commit metadata depends on how tables are tracked and is not standardized by this spec.”1

It gives two examples instead.

SchemeHow the swap happensStatus
File system tables Atomic rename, in filesystems that support it — like HDFS. Metadata is v<V>.metadata.json.1 Deprecated. The spec calls it unsafe in object stores and local file systems, and it is being removed in v4.1
Metastore tables A pointer in a metastore or database updated with a check-and-put, which validates that the version the write is based on is still current.1 The real answer. This is the check that rejected Writer B.

Read that deprecation twice

A “catalog” that is just a path on S3 — no service, no compare-and-swap — cannot safely perform this commit, and the spec now says so directly. If you have ever seen a Hadoop-catalog Iceberg table on S3 described as production-ready, this is the sentence to cite.

So the catalog is not an accessory to the format; it is the component that holds the only mutable state and performs the only atomic operation. Which raises the obvious question — if every engine needs to talk to the catalog to read or write a table, what do they talk over? That is the Iceberg REST Catalog, and it turns out to be the pivot the whole 2026 interoperability story turns on. Lesson 5 picks it up.

6.Read this next

Primary source: the Iceberg table spec, § Optimistic Concurrency, § Sequence Numbers, and Appendix § Metastore Tables. Then Reliability › Retry validation for the assumptions-and-actions framing.

Worth doing once for real: open a table’s snapshots and history metadata tables in whatever engine you use, and read the summary map on a few snapshots. Everything in this lesson is visible there.


Sources

  1. Apache Iceberg Table Spec — § Optimistic Concurrency, § Sequence Numbers, § Snapshot References, § Table Metadata, § File System Tables, § Metastore Tables, § Scan Planning.
  2. Apache Iceberg docs › Reliability — concurrent write operations, cost of retries, retry validation.