When two jobs write to the same Iceberg or Delta table simultaneously, both operate under optimistic concurrency control by reading the latest table snapshot and writing candidate Parquet files in isolation. The job that finishes first commits its metadata changes successfully, winning the race and creating a new snapshot version. The second job then encounters a commit collision and must evaluate whether its changes conflict with the newly committed snapshot before attempting an automatic retry.
Optimistic concurrency protocol
Under optimistic concurrency, neither writer places a distributed lock on the table before starting work:
- If both jobs are pure append operations that do not modify existing data files, their changes do not conflict. The second writer simply re-bases its commit on top of the newly committed snapshot version and appends its files without reprocessing data.
- If one job executes a DELETE or MERGE that touches data files that the first job just rewrote or removed, a write conflict occurs.
- The second job cannot blindly retry its commit because the data files it used as input no longer exist in the active table state. To recover safely, it must restart its query, scan the updated snapshot, re-evaluate its predicate logic, and generate fresh output files.
Identifying commit conflicts
Conflict resolution depends on transaction isolation levels:
- Under snapshot isolation, operations check whether concurrent commits modified the exact files or partitions they intended to touch.
- Under serializable isolation, stricter checks detect phantom reads and predicate overlaps, causing more transactions to fail.
- Retrying blindly on concurrent updates without validating input state causes data corruption or stale overwrites.
Pipeline design to minimize retries
To prevent continuous retry storms and wasted compute in production, follow practical scheduling practices:
- Partition write jobs by time or business entity, ensuring separate jobs write to distinct partitions or cluster keys.
- Consolidate multiple streaming merge streams into a single structured streaming job or stage updates through a message queue.
- Keep transaction runtimes short by compacting files during separate off-peak maintenance hours rather than bundling heavy compaction into critical streaming writes.