skip to content

How does an open lakehouse table format handle two writers committing at the same time?

level: middleimportance: should knowfreq 62%

answer

  1. nobody takes a lock while writing
  2. the collision is detected at the very end
  3. the swap is conditional on the base you read
  4. exactly one wins, the other rebases
  5. two appends rarely conflict; two rewrites do

basics

~20 s

Both writers work optimistically against the snapshot they read, then race to swap the current pointer. The swap is conditional, so exactly one wins; the loser re-reads the new snapshot, checks whether the winner's changes conflict with its own, and retries or fails.

solid answer

~60 s

Writers take no locks while they work. Each reads the current snapshot as its **base**, writes its data files invisibly, then attempts the commit as a **conditional swap**: make my metadata current, but only if the pointer is still the base I read. Because that swap is atomic, exactly one concurrent writer succeeds. The loser's swap fails and it enters a retry: it re-reads the now-current snapshot and asks whether the winning commit is **compatible** with its own operation. Two blind appends are compatible — the loser can rebase onto the new snapshot and re-attempt without recomputing anything, since its files are still valid. A conflict is real when the winner changed something the loser's operation depended on: both deleted or rewrote the same files, or the loser's `MERGE`/conditional update read rows the winner has since changed. Then the loser must recompute against the new state or abort. Retries are bounded, so under sustained heavy write concurrency a slow writer can lose repeatedly and fail — which is why the practical fix is to reduce contention rather than raise the retry count.

code

text · 9 lines
text
writer A            writer B
read base = v41     read base = v41
write files         write files
CAS 41 -> 42  OK
                    CAS 41 -> 42  FAIL (current is 42)
                    re-read v42
                    overlap with A's changes?
                      no  -> rebase, CAS 42 -> 43  OK
                      yes -> recompute against v42, or abort

go deeper

for a junior

Know that writers do not lock the table, that one commit wins the race, and that the other is retried rather than silently discarded.

for a middle

Explain the conditional swap on the base version and the difference between a rebase-and-retry and a genuine conflict. Interviewers expect the append versus MERGE contrast.

for a senior

Show operational judgment: recognise starvation and compaction-versus-ingest conflicts, and prescribe contention reduction — writer ownership, batching, partition scoping — over bigger retry budgets.

for a principal

Own the platform-level write topology: who is allowed to write which tables and partitions, how maintenance windows are scheduled against ingest, and what commit frequency the catalog and metadata layer can sustain.

## Optimistic, not pessimistic Acquiring a lock for the duration of a lakehouse write would be disastrous: writes commonly run for minutes to hours, and locking a table that long serializes the whole platform. So table formats use **optimistic concurrency control**: assume no one else is writing, do all the work unsynchronised, and detect the collision only at the moment of commit. The reason this is safe is that all the work really is unsynchronised — data files are written where nothing references them, so two writers producing files simultaneously cannot interfere at all. ## The race ```text writer A: read base = v41 ──▶ write files ──▶ CAS v41 -> v42 SUCCEEDS writer B: read base = v41 ──▶ write files ──▶ CAS v41 -> v42 FAILS (current is v42) └▶ rebase on v42, re-check, retry ``` Both writers read version 41 and plan their work against it. Both attempt a compare-and-swap conditioned on 41 still being current. The swap is a single atomic conditional operation, so exactly one succeeds. This is the entire mutual-exclusion story — there is no lock, only a losing conditional write. ## What the loser does next Failing the swap is not an error the user should ever see; it is the start of a retry loop with a genuine correctness check in the middle. The loser re-reads the current snapshot and compares the winner's changes with its own operation's assumptions: - **Compatible (rebase and retry).** Two appends touch nothing in common: the loser's data files are still perfectly valid, so it simply rebuilds its commit on top of version 42 and swaps again. Nothing is recomputed and no data is re-read. Most real concurrency is this case. - **Conflicting (recompute or abort).** The winner removed or rewrote files the loser also intends to remove or rewrite — two concurrent compactions over the same files, two updates hitting the same rows, an overwrite of the same partition. Or the loser's operation had a **read set**: a `MERGE` or a conditional `UPDATE ... WHERE` decided which rows to change by reading data the winner has since changed. Re-committing the loser's precomputed result would silently apply it to a state it never saw, producing a lost update. The loser must redo the operation against version 42, or fail loudly. The important nuance for an interview is that conflict detection is about **overlap of what each commit read and wrote**, not about wall-clock timing. Two writers to different partitions of the same table usually both succeed, one after a harmless rebase. ## What readers see meanwhile Nothing surprising: a reader resolves the current pointer once and then reads an immutable file set. A commit that lands mid-query does not change that query's results, and a query started after the commit sees all of it. That is snapshot isolation, and it falls out of the same design. ## Where it goes wrong in production - **Starvation.** A writer whose work takes twenty minutes competing with a job that commits every thirty seconds may lose every race until it exhausts its retries. Retry counts do not fix this; reducing commit frequency, batching the fast writer, or splitting the slow one does. - **Retry storms on conflicting work.** If two pipelines genuinely update the same rows, retries burn compute and still conflict. Route both through one writer, or partition ownership so each partition has exactly one writer. - **Compaction fighting ingestion.** Maintenance that rewrites files an ingest job is also touching is a classic recurring conflict; schedule it against quiet partitions or accept and re-run it. - **False confidence in "idempotent retry".** A retry that blindly re-swaps a precomputed result without re-checking the read set is how a lost update happens. If you implement anything like this yourself, the check is the point, not the loop. ## The design rules that follow One writer per table (or per partition) is the simplest correct topology and the one most platforms converge on. Commit in reasonably sized batches rather than continuously, so the pointer is not being swapped constantly. Keep long-running rewrites narrow so their conflict surface is small. And treat repeated commit failures as a contention signal to be designed away, not an error to be retried harder.

  • Why does an append usually retry successfully while a MERGE may have to be recomputed?
    An append asserts nothing about existing data: its files stay valid no matter what committed in between, so it can rebase onto the new snapshot and swap again. A MERGE decided which rows to insert, update or delete by reading the table. If the winning commit changed rows in that read set, re-applying the precomputed result would overwrite work it never saw, so it must be recomputed.
  • A long compaction job keeps losing to a streaming writer and eventually fails. What do you change?
    Not the retry count — the contention. Narrow the compaction to partitions the stream is not writing, run it during a quieter window, or have the stream commit in larger, less frequent batches. Optimistic concurrency starves the slow writer by design, so the fix is to stop the two from overlapping rather than to race harder.
  • What does a reader see if a commit lands in the middle of its query?
    Nothing changes for it. The reader resolved the current pointer at planning time and reads an immutable file set, so the new commit is simply not part of its scan. A query started after the commit sees all of it. This is snapshot isolation, and it comes free from immutable files plus a single pointer.
  • Do two writers appending to different partitions of the same table conflict?
    Usually not. They still race for the same pointer, so one will lose the swap, but the check that follows finds no overlap between what each commit read and wrote, so the loser rebases onto the winner's snapshot and commits without redoing work. The user sees one commit after the other.

Two people editing the same shared document offline: whoever saves first wins outright, and the second is told the base changed — a purely additive edit merges cleanly, but two rewrites of the same paragraph must be redone by hand.

saying these in an interview costs you the question

  • Says the table is locked for the duration of the write
  • Claims both concurrent commits are merged automatically
  • Thinks any retry is safe without re-checking the read set
  • Says raising the retry limit fixes a starving writer
  • Believes a mid-query commit changes what the running query reads

context