TT Lab
Get started
Learn Learning paths Courses

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

Continue in TT Lab

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.

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

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.