Lakehouse Table Format — Understanding Apache Iceberg Through Its Metadata
Two writers commit to one table at once — who wins, what's left behind, what's lost
Goal
With pyiceberg, deterministically reproduce a race in which two writers read the same metadata and then commit in turn. An append gets both in through automatic retry, and if you turn retries off, one side fails and the half-written files are left behind as orphans. A conditional overwrite is caught by validation even with retries, and with a single counter you confirm that a lost update is avoided only if you read anew and recompute afterward.
Why it matters
Iceberg has no locks. Each writer finishes writing its files and builds a new metadata, and then tries a conditional swap in the catalog, "change main only if the one I read is still the same". If two arrive at the same time, only one wins, and the loser reads the new metadata again and layers its change on top again. This is optimistic concurrency — you believe that most of the time there is no collision, and if there is, you do it again. The catch is that there are "changes that may be layered on again" and "changes that may not". Adding new files can be layered on again no matter what anyone did in between. But the change "set the rows matching this condition to this value" produces a wrong result if someone changed rows of the same condition in the meantime. The library validates that at the file level and rejects it, but if your code writes again with a value it read earlier without recomputing, nobody can stop it.
Steps
- Put a helper that reads one day's CSV with pyarrow in /root/ice/conc/common.py (
order_tsin the UTC time zone), and with /root/ice/conc/setup.py createlake.conc.ordersand load 2026-03-01. - With /root/ice/conc/race.py, read two Table objects
aandbfirst, then haveaappend 03-02 andbappend 03-03, and write the result to /root/ice/conc/out/race.json. - With /root/ice/conc/noretry.py, set the table property
commit.retry.num-retriesto0, run the same race again (03-04 · 03-05), and write the exception name of the losing side to /root/ice/conc/out/noretry.txt. - With /root/ice/conc/orphans.py, find the Parquet files that are in the data directory of the table location but that no snapshot points to, and write them to /root/ice/conc/out/orphans.json.
- With /root/ice/conc/conflict.py, create a counter table
lake.conc.counters(hits = 10), have two workers read the same value and each overwrite with +5, and write the exception name of the losing side to /root/ice/conc/out/conflict.txt. - With /root/ice/conc/retry.py, redo the losing side's +5 correctly so that the counter becomes 20.
- In /root/ice/conc/report.md, write three sections,
## 자동 재시도## 실패한 커밋의 흔적## 잃어버린 갱신(keep these headings as written; they stand for automatic retry, traces of a failed commit, and lost updates).
Notes
- Run the scripts like
cd /root/ice/conc && python3 race.py(from common import …). - Every commit in this lab is pyiceberg. When pyiceberg loses the conditional swap, it reads anew and tries again as many times as the retry setting of the table properties, and then prints 'Commit failed due to a concurrent update, retrying'.
- You get the exception name with
type(exc).__name__. - A common mistake: in step 6, after catching the exception, writing 15, the value read earlier (10) plus 5, again — the commit succeeds but the counter stops at 15. To redo steps 1–4 from the start, run
spark-sql -e "DROP TABLE lake.conc.orders PURGE"and begin at step 1. - Official docs: Reliability — Concurrent write operations · Configuration — Table behavior properties · Spec — Optimistic Concurrency · pyiceberg — Write support
A table made with Python
In /root/ice/conc/common.py, make day("YYYY-MM-DD") return that day's CSV as a pyarrow table (amount int32, order_ts a timestamp in the UTC time zone), and with /root/ice/conc/setup.py create lake.conc.orders with format-version 2 and load 2026-03-01.
pyiceberg's create_table(이름, schema=pyarrow_스키마) assigns a field ID to each column (the placeholders stand for the table name and the pyarrow schema). If you attach a time zone to the timestamp, it becomes the same timestamptz as a table Spark creates. The grader checks that the first snapshot added the row count of March 1.
The race — an append can just be layered on again
In /root/ice/conc/race.py, create a = load_table(…) and b = load_table(…) both first, then commit in the order a.append(03-02), b.append(03-03), and write the number of snapshots and the number of rows to /root/ice/conc/out/race.json as {"snapshots", "rows"}.
b holds the metadata from before a's commit, so its first commit attempt hits the condition. An append may be layered on top of whatever came in meanwhile, so pyiceberg reads anew and tries again and succeeds. The grader checks that the three snapshots line up in one chain (did not fork) and the row count.
If you turn retries off, the losing side fails
With /root/ice/conc/noretry.py, change the property commit.retry.num-retries of lake.conc.orders to "0", then, as in step 2, create a and b first and have a append 03-04 and b append 03-05. Write the exception name of b on the first line of /root/ice/conc/out/noretry.txt.
If retries are 0, an exception is raised the moment it loses the conditional swap. But b had already written all its data files before it tried to commit. The grader checks the exception name, the property value, and that March 4 is in while March 5 is not in the table.
The files a failed commit left behind
With /root/ice/conc/orphans.py, compare the list of files that all snapshots point to with the actual Parquet files in the data directory under the table location (t.location()), and write the paths of the files not in the list (orphans) to /root/ice/conc/out/orphans.json as {"orphans": [경로, …]} (the placeholder in the code stands for the paths).
An orphan file is not part of the table, so it is not read, but it takes up space. You must compare with the files that all snapshots point to, not just the current snapshot — the files of old snapshots are files that are alive for time travel. Clearing such files is remove_orphan_files in the next module.
An overwrite is caught by validation
With /root/ice/conc/conflict.py, create lake.conc.counters(name STRING, value BIGINT) anew (dropping it if it exists) with one row, hits = 10, and have a and b both read the value, then have a write 읽은 값 + 5 (the value read plus 5) with overwrite(…, overwrite_filter=EqualTo("name", "hits")) and b write in the same way. Write the exception name of b on the first line of /root/ice/conc/out/conflict.txt.
b's first attempt loses the conditional swap, and in the retry it is stopped by the validation "a file matching my condition (name = hits) came in meanwhile". Unlike an append, an overwrite layered on again as it is would erase a's result. The grader checks the exception name and whether there is a snapshot in which the counter was once 15.
Read anew and recompute
With /root/ice/conc/retry.py, redo b's +5. On every attempt, read the table anew, add 5 to the value read then, try the conditional overwrite, and start over from the beginning if it fails. When it is done, hits must be 20.
The library's retry only layers the commit on again; it does not re-read the value you read earlier for you. If you write the 15 you computed from the 10 read earlier as it is, the commit succeeds and a's +5 disappears — a lost update. The grader checks that the counter is exactly 20.
Concurrent write rules as team rules
In /root/ice/conc/report.md, write three sections, ## 자동 재시도 ## 실패한 커밋의 흔적 ## 잃어버린 갱신 (automatic retry, traces of a failed commit, and lost updates). In the second section, put the number of orphan files found in step 4, and in the third, the final counter value, as numbers.
If two or more jobs write to the same table, write down which jobs can be left to retries and which must start over from reading, and who clears the files of a failed job and when.