skip to content

Two Spark jobs write to the same Delta Lake table at once — how does the transaction log settle it?

level: seniorimportance: should knowfreq 60%

answer

  1. nobody takes a lock before writing
  2. the file name is the contention point
  3. only one writer can create version n+1
  4. the loser rebases if reads did not overlap
  5. the exception name tells you what overlapped

basics

~20 s

Delta uses optimistic concurrency: each writer tries to create the next numbered commit file in _delta_log, and only one can. The loser re-reads the commits it missed, checks whether they logically conflict with what it read, then retries or throws a concurrency exception.

solid answer

~50 s

Every Delta write follows read–write–commit. The writer records the snapshot version it read and which files or partitions it read, writes its new Parquet files (outside the log, so they are invisible until committed), then tries to create `<n+1>.json`. Creation succeeds only if that file does not already exist, so exactly one writer wins version n+1. The loser does not simply overwrite: it reads the commits that landed in between and applies conflict rules. If the winning commit only appended files that the loser never read, the loser silently retries at the next version reusing its already-written data files. If the two logically overlap, it fails with a specific exception — `ConcurrentAppendException` when files were added to a partition the transaction read, `ConcurrentDeleteReadException` when a file it read was deleted, `ConcurrentDeleteDeleteException`, `MetadataChangedException`, `ProtocolChangedException`. The fix is usually narrowing read sets with partition predicates so writers do not overlap.

code

sql · 9 lines
sql
-- conflicts: the merge reads the whole table, so any concurrent writer overlaps
MERGE INTO events t USING updates s
  ON t.event_id = s.event_id
WHEN MATCHED THEN UPDATE SET *;

-- disjoint: the partition predicate proves this writer read only its own partition
MERGE INTO events t USING updates s
  ON t.event_id = s.event_id AND t.dt = '2024-05-01'
WHEN MATCHED THEN UPDATE SET *;

go deeper

for a junior

Know that Delta writes are optimistic: writers do not lock, they race to create the next numbered commit file, and only one can win a given version.

for a middle

Explain the three phases — read snapshot, write data files, attempt commit — and why a loser can often retry without rewriting data. Name at least one conflict exception and what it indicates.

for a senior

Diagnose a real conflict: read the exception name, tie it to what the transaction read, and fix it by narrowing read sets with partition predicates, adjusting layout, or restructuring which job owns which data.

for a principal

Own the write topology. Decide who may write to which table, whether multi-cluster writers on object storage are supported by the configured log store, and where isolation level and retry budgets sit for the platform.

## The protocol Delta Lake has no lock manager and no coordinator between writers. Concurrency control is optimistic and lives entirely in the commit file names under `_delta_log`. A write goes through three phases: 1. **Read.** The writer resolves the table's current snapshot — say version *n* — and remembers both that number and its **read set**: the predicates it filtered on and the files it actually consulted. 2. **Write.** It writes new Parquet data files into the table directory. These are invisible to every reader, because no `add` action references them yet. A crash here leaves orphan files, not a corrupt table. 3. **Commit.** It attempts to create `<n+1>.json` containing its actions. The commit attempt is the only synchronisation point, and it hinges on one primitive: **create-if-not-exists**. Whoever creates version n+1 first wins; every other writer's attempt fails because the object already exists. ## What the loser does A failed commit is not automatically an error. The loser enters conflict resolution: it reads the commits that appeared between its read version and now, and asks whether any of them invalidates the work it did. - If the winner only **added** files, and none of them fall inside the partitions or predicate range the loser read, there is no logical conflict. The loser rebases — retries the commit at version n+2 — **reusing the Parquet files it already wrote**. No data is rewritten. This is why a table can absorb many concurrent appenders cheaply. - If the commits do overlap, the transaction fails with a named exception, and the application must re-run the operation against the newer snapshot. A pure append is flagged `isBlindAppend: true` in `commitInfo`: it read nothing, so it can never conflict with another append. Two streaming jobs appending to the same table will interleave happily forever. ## The conflict exceptions and what each means - **`ConcurrentAppendException`** — another transaction added files into a partition or predicate range that this transaction read. The classic case is two `MERGE` jobs against the same unpartitioned table, or against the same partition. The transaction's result might have been different had it seen those rows, so Delta refuses. - **`ConcurrentDeleteReadException`** — a file this transaction read was removed by a concurrent commit; typically a `MERGE` racing an `OPTIMIZE`, a `DELETE`, or a retention cleanup. - **`ConcurrentDeleteDeleteException`** — both transactions tried to delete the same file. - **`MetadataChangedException`** — the table's `metaData` action changed underneath, e.g. a concurrent schema evolution. - **`ProtocolChangedException`** — the `protocol` action changed, e.g. someone upgraded the table to enable a feature mid-flight. - **`ConcurrentTransactionException`** — two runs of the same streaming query (same `txn` appId) committed concurrently, which usually means two instances of a job are running against one checkpoint. ## Making writers not conflict Conflicts are decided by **read sets**, so the practical fix is to make read sets disjoint. If two jobs each own a date partition, put the partition value in the `MERGE` condition itself — `ON t.id = s.id AND t.dt = '2024-05-01'` — so Delta can prove each transaction only read its own partition. Without that predicate the merge reads the whole table and every other writer conflicts with it. Other levers: partition or cluster the table along the natural writer boundary; funnel mutations through a single job and let concurrency be appends only; wrap operations in a bounded retry with backoff, since a retry re-plans against the newer snapshot and often succeeds. ## Isolation level Delta exposes a table property `delta.isolationLevel` with values `Serializable` and `WriteSerializable` (the default). `WriteSerializable` permits blind appends to be ordered differently from how a reader might infer, in exchange for far fewer conflicts; `Serializable` is stricter and rejects more. Readers are unaffected either way — a reader always sees one consistent snapshot and never blocks a writer or is blocked by one. ## The storage requirement underneath The whole scheme depends on the storage layer providing mutually exclusive create-if-not-exists for the commit object. HDFS and Azure storage give atomic rename/create. Amazon S3 historically did not, which is why Delta ships alternative log stores — notably a DynamoDB-backed one in its storage module — to coordinate multi-cluster writes to S3 tables; S3 has since gained conditional writes. This matters in interviews: **a single Spark cluster writing to an S3 Delta table is safe because the driver coordinates in-process, but two independent clusters writing to the same S3 table without a coordinating log store can both believe they created the same version.** Check the configured log store before you promise multi-writer safety. ## What readers see throughout Nothing partial, ever. A reader pins a snapshot version at query start and reads exactly those files. A concurrent commit does not disturb it, and a reader is never blocked by a writer nor a writer by a reader — the cost of concurrency is paid entirely by writers losing races.

  • Why does adding a partition predicate to a MERGE condition reduce ConcurrentAppendException?
    Conflicts are evaluated against what the transaction read. Without a partition predicate the merge reads the whole table, so any concurrent append anywhere overlaps. With `AND t.dt = '2024-05-01'` in the condition, Delta knows the read set is one partition, and a writer touching a different partition no longer conflicts.
  • Two independent clusters write to the same Delta table on S3. What must you check?
    Which log store is configured. The commit protocol needs mutually exclusive create-if-not-exists on the commit object; a single cluster is safe because its driver coordinates in-process, but two clusters need a coordinating store — historically a DynamoDB-backed log store — or conditional-write support. Without it, two writers can both believe they created the same version.
  • Does a concurrent writer ever block or corrupt a running Delta reader?
    No. A reader pins a snapshot version at query start and reads exactly the files that version lists. New commits are new files and new versions; they do not modify anything the reader is using. Writers never block readers, and readers never block writers — only writers race each other.
  • Is retrying a failed Delta commit always safe?
    For an idempotent operation planned against a fresh snapshot, yes — that is what Delta's internal rebase does. Blindly retrying application logic that computed values from a stale read is not: re-run the operation against the new snapshot rather than replaying the old result, and bound retries so a permanently conflicting job fails loudly instead of spinning.

Two people editing the same shared document offline: whoever saves first wins, and the second must rebase. If their edits touched different sections the rebase is automatic; if they edited the same paragraph, the tool refuses and asks a human.

saying these in an interview costs you the question

  • Saying Delta takes a table lock before writing
  • Assuming the second writer's data is silently overwritten or lost
  • Believing every concurrent write fails and must be serialized manually
  • Thinking a reader can observe a half-finished commit
  • Claiming multi-cluster S3 writers are safe with no coordinating log store

context