How does a catalog make a commit to a lakehouse table atomic when two writers race?
answer
- nobody locks anything for the slow part
- the expensive work happens invisibly first
- the whole commit reduces to one value change
- swap only if the value is still what I read
- exactly one winner, the loser rebases and retries
basics
~20 sEach writer stages new files, then asks the catalog to move the table's metadata pointer only if it still holds the value that writer read. This conditional swap succeeds for exactly one writer; the loser re-reads, revalidates its change and retries.
solid answer
~50 sCommits use **optimistic concurrency** with the catalog supplying the one atomic operation object storage lacks. A writer reads the current metadata pointer, say `M7`, writes its new data files and a new metadata version `M8` that is based on `M7`, then asks the catalog to set the pointer to `M8` **on the condition that it is still `M7`**. Implementations differ — a conditional row update in the metastore's database, a versioned update call, a compare-and-swap in a REST commit request — but the guarantee is identical: at most one writer can win a given transition. The loser gets a conflict, re-reads the now-current pointer, checks whether its change is still valid against it (an append usually is; an overwrite of rows another writer just changed may not be), rebuilds its metadata on the new base and retries. Its already-written data files are normally reusable, so a retry is cheap. Readers see either `M7` or `M8`, never a mixture.
code
text · 10 lineswriter A catalog writer B
read pointer -> M7 read pointer -> M7
write data files (invisible) write data files (invisible)
write metadata M8 (base M7) write metadata M8' (base M7)
swap M7 -> M8 -------------> ok, pointer = M8
swap M7 -> M8' -> CONFLICT
read pointer -> M8
revalidate change against M8
write metadata M9 (base M8)
ok, pointer = M9 <-------- swap M8 -> M9go deeper
Recall that concurrent writers do not overwrite each other: one commit wins, the other is told to retry. Know that files are written before the commit and only become visible when it succeeds.
Walk the read-stage-build-swap loop out loud, name the swap as a conditional compare-and-swap, and explain why the expensive work happens outside any lock.
Show you can diagnose contention: distinguish append conflicts from overwrite conflicts, explain when a retry must recompute, and describe orphan files left by writers that died before committing.
Own the concurrency policy for the platform: who may write which tables, how maintenance is scheduled against user writes, and what isolation guarantee you promise consumers.
## Why the catalog has to be involved A commit to a lakehouse table touches many objects: new data files, possibly delete files, and a new metadata version. Object storage offers no transaction across objects and, classically, no way to say "create this only if nothing else did". So the format reduces the whole commit to a **single-value update**: everything else is written first, in the background, visible to nobody, and the commit becomes one change to one pointer. The catalog is what makes that one change atomic and conditional. ## The optimistic protocol, step by step 1. **Read.** The writer resolves the table and notes the current metadata pointer — call it `M7` — plus enough state (schema, partition spec, the set of files it is modifying) to validate later. 2. **Stage.** It writes its data files under the table's location. These are invisible: no committed metadata version mentions them, so no reader will scan them. 3. **Build.** It writes a new metadata version `M8` describing the table as of `M7` **plus** its changes. 4. **Swap.** It calls the catalog: *set the pointer to `M8`, but only if it is currently `M7`.* 5. **Win or retry.** If the condition holds, the commit is done and instantly visible. If another writer got there first, the call fails, and the writer goes back to step 1 against the new current version. This is a compare-and-swap loop. Nothing is locked while the expensive work happens, which is exactly the point: two jobs writing into different partitions for twenty minutes each do not block one another, and only the final millisecond is serialized. ## How different catalogs implement the swap - **A metastore backed by a relational database** performs the update inside a database transaction, comparing the stored pointer property before writing the new one. Some deployments historically added an explicit table lock around the sequence. - **A managed cloud catalog** exposes an update call that takes the version you believe is current and rejects the write if it has moved. - **A REST catalog protocol** sends the commit as a set of *requirements* (assertions about current state, such as "the current snapshot is still the one I read") together with the *updates* to apply; the server evaluates the requirements and applies the updates atomically or rejects the whole request. Where no catalog is in the path, the same guarantee has to come from somewhere else — a filesystem's atomic rename, a storage conditional put, or an external commit coordinator. A format whose commit ordering lives in its own log needs precisely one of these; "we put the table in a metastore" does not by itself supply atomicity if the metastore is only used for discovery. ## Conflict validation: not every retry is safe Retrying is not just re-swapping the pointer. The writer must re-check that its change still means the same thing against the new base: - **Appends** almost always commute. Adding files to a table someone else also added files to is fine; the retry rebases and succeeds. - **Overwrites, deletes and merges** may not. If your `MERGE` decided to rewrite the files holding January's rows and another commit has just replaced those files, replaying blindly would resurrect deleted rows or lose the other writer's update. The commit therefore asserts something about the files or partitions it read, and if those assertions no longer hold, the operation fails and must be recomputed — sometimes by re-running the whole statement. This is why the isolation level matters: serializable checks that nothing relevant changed at all, while snapshot isolation permits some non-conflicting concurrent changes. Both are implemented by which assertions the commit attaches. ## Failure modes worth naming in an interview - **Live-lock under heavy contention.** Many writers to one table means many losers per second; with retries and backoff exhausted, jobs fail. The fix is usually to reduce writer fan-in (one writer per table or per partition range), batch more per commit, or let maintenance jobs run in windows. - **Orphan files.** A writer that stages files and then dies never commits, leaving unreferenced objects. They cost money and are cleaned by a dedicated orphan-file job, never by deleting anything "that looks unused" by hand. - **A catalog that cannot do a conditional update** is not a valid commit coordinator, however good it is at discovery. - **Two catalogs for one table.** Two independent pointers means two independent compare-and-swap loops, and neither sees the other's commits. ## The reader's view Because the pointer flips in one step and old versions stay readable, a reader either plans against `M7` or against `M8`. A long-running scan started before the swap keeps reading `M7`'s files, which remain until retention removes them — the reason retention windows must exceed your longest query.
- Why does a retrying writer usually not have to rewrite its data files?Because the data files it staged are immutable and were never referenced by any committed version. On retry it rebases only the metadata: take the new current version, add the same files, swap again. Files are rewritten only if conflict validation shows the operation itself is no longer valid, such as an overwrite of rows another commit just replaced.
- What distinguishes an append conflict from a genuine overwrite conflict?Appends commute: adding files alongside someone else's added files changes nothing about their rows, so the retry succeeds. An overwrite, delete or merge asserts something about the files it read; if a concurrent commit replaced those files, replaying would drop or resurrect rows, so the assertion fails and the statement must be recomputed.
- A table gets dozens of failed commits per minute under contention. What do you change?Reduce writer fan-in rather than raising retry counts. Funnel streaming writes through one writer per table, batch more rows per commit so commits are fewer and larger, partition writers so they touch disjoint partitions, and schedule maintenance jobs in windows instead of racing user writes.
It works like changing a shared whiteboard entry with the rule "erase only if it still says what I saw": two people can prepare their updates in parallel, but only one erase lands, and the other has to look again and redo their reasoning.
saying these in an interview costs you the question
- Says last write wins, so the second commit overwrites the first
- Thinks the catalog locks the table for the entire write
- Claims object storage itself makes multi-file commits atomic
- Believes any retry is safe without revalidating the change
- Assumes a discovery-only metastore entry provides atomicity