Lakehouse Table Format — Understanding Apache Iceberg Through Its Metadata
Two ways to change rows in files you can't modify — copy-on-write and merge-on-read
In one line
Since a Parquet file cannot be edited, changing one row means either writing that file anew (copy-on-write) or writing a separate delete file that says "row N of file X is deleted" and filtering it out when reading (merge-on-read). The former makes writes expensive and reads cheap, and the latter is the opposite. Three table properties decide this.
Why it is a problem — one row in a million-row file
One order turned into a refund. The Parquet file containing that order holds another million orders. Once a file is written, it does not change (the requirement of the spec is "no in-place writes"). Then do we have to rewrite a million rows to fix one?
In format version 1, yes. Version 2 added delete files to solve this problem. It became possible to delete or change individual rows inside immutable data files without rewriting them.
How it works — three shapes of deletion
Row-level deletes come in two branches.
| Kind | What deletes with | Mainly used by |
|---|---|---|
| Position delete file (version 2) | (data file path, row position) | Spark's merge-on-read |
| Deletion vector (version 3 and up) | One bitmap per data file | Position deletes of version 3 tables |
| Equality delete file | A column value (for example order_id = 'O123') | Streaming upsert (Flink, etc.) |
A position delete file is a file with two columns, file_path and pos (the row number counted from 0), written sorted by file_path and then pos. While reading a data file, the reader skips the deleted positions of that file. A deletion vector holds the same information more compactly as a roaring bitmap inside a Puffin file, with at most one per data file. In version 3, position delete files are deprecated.
An equality delete file does not know positions, only values. A streaming upsert cannot afford to look up which file and which position the old row is in, so it writes only "delete the old row of this key". In exchange, the reader has to match the rows of all related data files against the values, so it is the most expensive.
Which deletes apply to which data is decided by the sequence number. According to the scan planning rules, a position delete applies to data files with an equal or smaller sequence number, and an equality delete only to data files with a strictly smaller one. That is why a new row inserted again with the same key after an equality delete is not deleted.
Write mode — the table remembers
Which way Spark does DELETE, UPDATE, and MERGE is decided by the table properties write.delete.mode, write.update.mode, and write.merge.mode, and all three default to copy-on-write. Merge-on-read is available only in version 2 and up.
- copy-on-write: Writes the data files that contain changed rows anew as a whole and removes the old files from the list. There are no delete files, so reads are fastest.
- merge-on-read: Leaves the data files as they are and writes only delete files (and, for changed rows, the data files holding the new rows). Writes are small and fast, but deletes have to be applied on every read.
MERGE INTO does WHEN MATCHED (delete, update) and WHEN NOT MATCHED (insert) in one commit. If two changes hit one source row, it fails because it cannot tell which to apply, so keep the change set down to one per key.
What it looks like in the field
A CDC load table gets slower and slower. A merge-on-read table that MERGEs every few minutes keeps accumulating delete files and has to match them on every read. You have to melt the deletes into the data with periodic compaction (the delete ratio and count criteria of rewrite_data_files, rewrite_position_delete_files).
Another engine returns rows that were deleted. An engine that does not understand delete files reads even the deleted rows of a merge-on-read table. If several engines read the same table, first check which delete formats each engine supports. In the lab environment, DuckDB 1.5.5 and pyiceberg 0.12 read with version 2 position deletes applied.
A table modified heavily once a day. If a nightly batch rewrites yesterday's partition as a whole and dashboards read it hundreds of times during the day, copy-on-write is better. The write cost is paid once and the read benefit comes hundreds of times.
What really matters in practice
- copy-on-write is write amplification, and merge-on-read is read amplification. Choose by how often you write and how often you read.
- The mode is a table property. All three default to copy-on-write, and merge-on-read needs version 2 or up.
- For a merge-on-read table, compaction is part of operations. When delete files pile up, reads get slow.
- Check that the reading engine supports the delete format. An engine that does not know it returns deleted rows.
What you will do in the next lab
With the same March orders, you create a copy-on-write table and a merge-on-read table, run the same DELETE, and confirm from the commit summaries that one side rewrites data files while the other only adds position delete files. After running the same MERGE on both tables with a change set, you open one position delete file directly with pyarrow to see which row of which file it points to, compare the bytes the two MERGEs wrote, and check whether DuckDB and pyiceberg read with the deletes applied.