Delta Concurrent Append Conflicts — Retry Smart
Append rejected? Delta's optimistic concurrency blocks overlapping writes.
20+ years shipping production backend systems. Written from production experience, not tutorials.
- ✓Delta Lake basics: tables, versions, and time travel
- ✓Batch and streaming writes (MERGE, append, foreachBatch)
- ✓Partitioning concepts and basic SQL predicate logic
- Delta Lake guards tables with optimistic concurrency: overlapping commits on the same data get rejected with concurrent-append or commit-conflict errors, not silent corruption
- Retry idempotently: deterministic job logic plus merge/upsert keys (or streaming transaction IDs) make re-runs safe instead of duplicating rows
- Cut contention by partitioning so writers touch different partitions, and by separating OPTIMIZE/VACUUM windows from write traffic
- Pick isolation deliberately: Serializable (default) fails on any overlapping change, WriteSerializable allows compatible appends — then alert on retry rates per table
Think of a Delta table as a shared ledger where writers check 'did anyone edit since I looked?' at commit time. If someone touched the same pages, your commit is rejected and you retry on the fresh version. That rejection is the concurrent-append error — protection, not corruption. Real ledgers fix it the same way: split the book so clerks rarely collide (partitioning), agree pure additions never conflict (isolation levels), and make entries re-writable without duplication (idempotent retries).
It starts as an occasional flake: a ConcurrentAppendException at 3 AM, a retried Airflow task, green by morning. Then traffic doubles, the streaming sink and the hourly backfill start colliding daily, OPTIMIZE jobs trip over writers, and 'transient Delta errors' become the top on-call category — each retry re-reading the table, each collision wasting a full commit's work. Teams discover that Delta's guarantees are strict in exactly the ways concurrent pipelines violate.
Delta Lake uses optimistic concurrency control: each commit reads the latest table version, does its work, then commits only if no conflicting transaction landed meanwhile. Serializable isolation (the default) rejects any overlapping data change; WriteSerializable permits concurrent appends that don't touch the same files. The errors look scary but they're the system working — the alternative (last-writer-wins on parquet directories) is silent data loss, which is strictly worse.
This article turns conflicts from pages into plumbing. You'll learn what each conflict error actually guards, how to write idempotent jobs that retry safely, how partitioning and scheduling dissolve contention, when to relax isolation (and when never to), and how to monitor conflict rates so growth never surprises you again.
Optimistic Concurrency: The Ledger That Checks at Commit
Delta commits are atomic log appends: each transaction reads the latest version N, computes new files, then appends version N+1 — but only if the table is still at N. If another writer committed N+1 first with conflicting changes, your commit is rejected and you retry against the fresh snapshot. No locks, no waiting, no silent overwrites: contention surfaces as explicit errors instead of lost data. This is optimistic concurrency — assume no conflict, verify at commit, retry on rejection.
Serializable (default) treats any overlapping data change as conflict: concurrent appends to the same partition, MERGE vs overwrite on shared files, OPTIMIZE rewriting files a writer read. WriteSerializable relaxes this for blind appends — writers that add files without reading existing ones (pure inserts, streaming appends) commit alongside each other freely. The distinction is semantic, not just performance: Serializable preserves full snapshot equivalence for mixed workloads; WriteSerializable trades snapshot strictness for append throughput where business logic tolerates it.
Read the error as the guard working: ConcurrentAppendException names the overlapping versions, commit-conflict messages name the files. Both tell you who (operation type), when (versions), and where (partitions) — the full collision report. Log these fields per table and you get contention analytics for free; swallow them into generic retries and you get mystery lag.
Idempotent Writers: Retries That Converge
A retry is only safe if re-running can't duplicate: that's idempotency, and non-idempotent retries caused our incident's $410K double-apply. The failure mode is 'committed but unacknowledged' — the commit lands, the acknowledgment dies in the conflict, the orchestrator reports failure, and the retry applies everything again. Blind appends and full overwrites re-applied are duplication machines; keyed MERGEs re-applied converge to the same rows.
Build idempotency per writer shape. Batch corrections: MERGE on business keys (order_id + updated_at) so re-runs update the same rows to the same values. Streaming sinks: foreachBatch with transaction tracking (streaming query ID + batch ID recorded in the target, or Delta's built-in exactly-once append handling) so replays skip committed batches. Partition overwrites: dynamic overwrite of the exact same partition predicate — re-running writes identical files, and Delta's log makes the second commit a no-op equivalent.
Prove it, don't assert it: the idempotency test runs every writer twice against fixed input and diffs row counts plus checksums. Two runs, identical table — merge that test and ambiguous failures become boring. Without it, every Airflow retry bump is a loaded duplication gun pointed at the money table.
Partition Fencing: Writers That Never Meet
The best conflict is the one structurally impossible: partition fencing gives each writer exclusive territory so snapshots never overlap. Streaming owns the fresh tail (last 6 hours), backfills own settled history (older than 6 hours), enforced by write predicates both jobs share from one config — not by schedule luck. Overlapping predicates are the collision; disjoint predicates are the cure, regardless of timing.
Granularity sets fence strength. Date-partitioned tables fence by day (coarse but simple); hourly sub-partitioning fences tighter for hot tails. When two writers must touch one partition (late-arriving corrections to today's data), fence by key range within it instead — corrections MERGE on order_id ranges the streaming append never rewrites. Fences compose with isolation: fenced Serializable writers never meet, so they pay no conflict cost while keeping full guarantees.
Maintenance gets fenced too: OPTIMIZE and VACUUM rewrite and expire files, making them contenders like any writer. Run them in writer-free windows (4 AM, not midnight alongside backfills), scope OPTIMIZE to settled partitions (WHERE dt < today), and never let VACUUM's retention window (default 7 days) approach streaming lateness — a long stall plus aggressive retention deletes files a restarting reader still needs, converting conflicts into data loss.
Isolation Levels: Serializable vs WriteSerializable
Choose isolation per table semantics, never cluster-wide. Serializable (default): any overlapping committed change aborts yours — required wherever updates, deletes, or MERGEs mix with other writers (money tables, dimension rebuilds), because snapshot equivalence is the correctness property. WriteSerializable: concurrent blind appends commit freely, conflicts only on file-level modification overlap — right for append-only event sidecars and high-throughput ingest where writers never update each other's rows.
The cost of wrong choice is asymmetric. Serializable on pure-append ingest buys conflicts you didn't need (throughput collapses under fan-in). WriteSerializable on mixed workloads buys anomalies you can't see (lost updates that pass every row-count check). When unsure, keep Serializable and fence harder — its errors are loud and retryable, while relaxed anomalies are silent and forensic.
Snapshot isolation (Databricks: the Serializable implementation detail) adds one more subtlety: readers see consistent snapshots regardless, so long-running readers never block writers — but writers with stale snapshots fail more as table churn rises. Shorten batch windows and commit granularly (per-partition commits over whole-table rewrites) to keep snapshots fresh: the smaller your commit's footprint, the less surface conflicts can strike.
Scheduling and Maintenance Windows That Don't Collide
Conflicts are timetable failures as often as semantic ones: the midnight stack (streaming + backfill + OPTIMIZE + VACUUM) collides because everything runs when engineers sleep, not when tables are free. Draw the contention calendar per table — every writer and maintenance job with its partitions and duration — and the overlaps glow. Stagger heavy commits (backfill at :15, never :00 with the streamer), move OPTIMIZE to 4 AM writer-free windows, and run VACUUM weekly against 7-day retention with readers' max lateness modeled in.
Streaming deserves special scheduling respect because it never stops: its micro-batch commits (every 2 minutes in our incident) make any overlapping batch writer a continuous collision source. Either fence the streamer off the batch's partitions entirely or pause-and-catch-up around bounded backfills (stop stream, run backfill, resume — lag spikes once instead of conflicting for hours). The math favors the pause: one 20-minute catch-up beats 47 conflicted retries by orders of magnitude.
Orchestrator config finishes the job: retries with exponential backoff plus jitter (thundering herds re-collide deterministically), conflict-aware alerting (page on sustained conflict rate, not single flakes), and SLA-budgeted backfill windows that shrink as volume grows — the 6-hour window that 'worked' at half volume is the midnight page at full volume.
Monitoring Conflict Health as You Scale
Conflict rate per table is the leading indicator volume growth shows up in first — weeks before lag pages. Derive it from the transaction log: commits landed (DESCRIBE HISTORY) versus attempts made (writer-side attempt logging), per operation type, per day. A table drifting 0 → 5 → 20 daily conflicts is 2 quarters from a storm; the fix (re-fence, split tables, graduate to WriteSerializable append paths) is cheap at 5 and structural at 47.
Dashboard the full contention surface, not just errors: streaming lag p99, backfill duration trend, OPTIMIZE overlap minutes with writers, and retry-count histograms per job. Retries hide conflicts from task-green dashboards while lag and cost grow — the incident's Airflow showed green for weeks while streaming lag climbed 2 → 40 minutes. Alert on the lag and the rate, never on single flakes (one conflict is the system working).
Review on onboarding, not on fire: every new writer gets a contention review (partitions touched, isolation of target, idempotency proof, schedule vs existing calendar) before its first production commit. Growth re-collides fenced writers within 2-3 quarters as volumes shift predicates and schedules drift — the review is the fence maintenance, and the conflict-rate graph proves it holds.
The Midnight Collision: Streaming Sink vs Backfill, 47 Conflicts a Night
- Retries without idempotency convert conflicts into duplication: the $410K double-apply happened because 'failed' can mean 'committed but unacknowledged'. Key every writer (merge keys, transaction IDs) so re-runs converge instead of multiply.
- Bumping retries treats the dashboard while feeding the disease — each retry re-does full rewrites and raises collision odds. Dissolve contention structurally (partition fencing, keyed merges, separated maintenance windows) instead of out-retrying it.
- Serializable is a default, not a law: audit-style append-only tables run fine (and faster) on WriteSerializable, while money tables keep Serializable plus fencing. Choose per table, document why, and alert on conflict rates to catch drift.
| File | Command / Code | Purpose |
|---|---|---|
| collision_forensics.sql | DESCRIBE HISTORY lake.orders LIMIT 25; | Optimistic Concurrency |
| idempotent_merge.py | from delta.tables import DeltaTable | Idempotent Writers |
| fenced_writers.py | FRESH_HOURS = 6 # streaming owns the tail; backfill owns settled history | Partition Fencing |
| IsolationChoice.scala | spark.sql("""ALTER TABLE lake.orders SET TBLPROPERTIES (""" + | Isolation Levels |
Key takeaways
Common mistakes to avoid
5 patternsBumping orchestrator retries instead of fencing writers
Retrying non-idempotent writers on ambiguous failures
Running OPTIMIZE/VACUUM alongside writers at midnight
Relaxing isolation to fix overlapping writers
Alerting on task failure instead of conflict rate and lag
Interview Questions on This Topic
What does a ConcurrentAppendException actually protect against?
Frequently Asked Questions
20+ years shipping production backend systems. Written from production experience, not tutorials.
That's Databricks. Mark it forged?
5 min read · try the examples if you haven't