skip to content

In Iceberg, what changes when write.delete.mode is set to merge-on-read?

level: middleimportance: must knowfreq 66%

answer

  1. immutable files cannot be edited in place
  2. either rewrite the file or record the change
  3. the write gets cheap, someone else pays
  4. per-operation table property, three of them
  5. delete files merged at scan time until compaction

basics

~20 s

Instead of rewriting every data file that contains a matching row, the engine writes small delete files that mark rows as removed and commits those. Writes get much faster; reads must merge deletes at scan time until compaction materialises them.

solid answer

~40 s

Iceberg exposes three properties — `write.delete.mode`, `write.update.mode` and `write.merge.mode` — each set to `copy-on-write` or `merge-on-read`, with copy-on-write as the documented default. Under **copy-on-write**, a `DELETE` reads every data file holding a matching row and rewrites it without those rows, so the commit is expensive but the result is clean, immediately-fast-to-read data files. Under **merge-on-read**, the engine leaves the data files alone and commits small delete files that identify the removed rows; the write touches far less data, but every reader must now apply those deletes while scanning, and the delete files accumulate until a compaction job merges them away. Merge-on-read requires table `format-version` 2 or higher. The choice is per operation type, so frequent small deletes can be merge-on-read while a rare bulk overwrite stays copy-on-write.

code

sql · 7 lines
sql
ALTER TABLE prod.db.events SET TBLPROPERTIES (
  'write.delete.mode' = 'merge-on-read',
  'write.update.mode' = 'merge-on-read',
  'write.merge.mode'  = 'copy-on-write'
);

DELETE FROM prod.db.events WHERE user_id = 42;

go deeper

for a junior

Recall that data files in this format cannot be edited, so a delete either rewrites whole files or records the removed rows somewhere separate. Know that a table property picks between those two behaviours.

for a middle

Explain both paths concretely: which files are rewritten under copy-on-write, what a delete file contains under merge-on-read, and that the properties are set per operation type with copy-on-write as the default.

for a senior

Demonstrate that merge-on-read is deferred work: quantify the read-side merge cost, tie the mode to a compaction cadence you would actually run, and be ready to diagnose a table whose delete files have outrun maintenance.

for a principal

Own the mode as a platform default per workload class — streaming upserts versus batch restatement — including engine compatibility across every reader, the format-version floor, and who is accountable for the compaction that merge-on-read presumes.

## The problem both modes solve Iceberg data files are immutable. Deleting or updating a row therefore cannot edit a file in place; the table format has to encode the change some other way. Iceberg offers two encodings, selected per operation by table property: ```sql ALTER TABLE prod.db.events SET TBLPROPERTIES ( 'write.delete.mode' = 'merge-on-read', 'write.update.mode' = 'merge-on-read', 'write.merge.mode' = 'copy-on-write' ); ``` The accepted values are `copy-on-write` and `merge-on-read`, and the documented default for all three is `copy-on-write`. Note that these are Iceberg *write modes*; they are a per-operation setting on one table, not two different table types. ## Copy-on-write The engine finds every data file containing at least one matching row, reads it, and writes a replacement file without those rows. The commit removes the old files and adds the new ones. Deleting one row from a 512 MB file rewrites the whole 512 MB — classic write amplification. In exchange, the table after the commit consists only of plain data files: readers scan them directly with no merge step, no extra file opens, and no accumulating debt. Copy-on-write suits tables that are written in large batches and read constantly: a nightly restatement, a GDPR erasure run, a slowly changing dimension rebuilt daily. ## Merge-on-read The engine leaves the data files untouched and writes **delete files** that mark rows as removed, committing them alongside the unchanged data files. Iceberg format v2 has two kinds: *position deletes*, which name a data file path and the row positions inside it, and *equality deletes*, which name column values whose rows are deleted. Format v3 replaces position delete files with **deletion vectors** stored in Puffin files — one compact bitmap of deleted positions per data file — which readers apply far more cheaply than a scattered set of small position delete files. The write is now proportional to the number of rows changed, not to the size of the files containing them, so a streaming CDC pipeline applying thousands of small upserts a minute becomes viable. The cost moves to the reader: for each data file the scan must locate the applicable delete files, load them, and filter out the marked rows. Iceberg decides applicability by sequence number within a partition — a position delete applies to data files whose sequence number is less than or equal to its own, an equality delete to data files with a strictly smaller sequence number. ## Where the cost accumulates Merge-on-read is a deferral, not a discount. Every un-compacted delete file is a file every relevant scan must open and merge. A table under continuous MoR writes degrades steadily: query latency climbs, and eventually a read fans out over thousands of tiny delete files. Two maintenance jobs pay the debt down. `rewrite_data_files` reads data files together with the delete files that apply to them and writes clean replacements, dropping the deletes; its `delete-file-threshold` option targets exactly the files carrying many deletes. `rewrite_position_delete_files` is a cheaper intermediate step that compacts the position delete files themselves without rewriting the data. Choosing merge-on-read therefore commits you to running compaction on a cadence matched to the write rate — that is the tradeoff an interviewer is probing. ## Choosing per operation Because the three properties are independent you can mix modes. A common shape: `write.delete.mode = merge-on-read` for frequent narrow deletes arriving from CDC, `write.merge.mode = copy-on-write` for a nightly `MERGE INTO` that touches a large fraction of the table anyway — when a rewrite would happen regardless, copy-on-write avoids the delete-file debt for free. Two more inputs matter. First, engine support: not every reader implements every delete type, so verify that all your engines can read the mode you enable. Second, format version: merge-on-read needs `format-version` 2 or later, so a v1 table must be upgraded first. ## Diagnosing which mode is in play The metadata tables answer it directly. `SELECT * FROM db.events.files` lists both data and delete files with their content type; a growing count of delete-content rows means merge-on-read writes are outrunning compaction. `SHOW TBLPROPERTIES prod.db.events` shows the configured modes, and the absence of the properties means the copy-on-write default is in force.

  • Reads on a merge-on-read table have got steadily slower. What is the fix?
    Delete files are accumulating faster than they are being merged away. Run `rewrite_data_files`, which reads each data file with the delete files that apply to it and writes clean replacements, using the `delete-file-threshold` option to target the worst offenders first. `rewrite_position_delete_files` compacts the delete files themselves as a cheaper stopgap. Then expire snapshots so the replaced files are actually released.
  • When is copy-on-write still the better choice for deletes?
    When the operation rewrites most of the affected files anyway — a bulk restatement or a `MERGE` touching a large share of rows — since you pay the rewrite once and end with clean files and no delete-file debt. It is also right when read latency is the hard requirement, when a consuming engine cannot read delete files, or when nobody will reliably run compaction.
  • Does merge-on-read work on a format v1 Iceberg table?
    No. Delete files were introduced in format version 2, so a v1 table must be upgraded — `ALTER TABLE ... SET TBLPROPERTIES ('format-version' = '2')` — before merge-on-read modes have any effect. Check first that every engine reading the table supports v2 delete files, because upgrading the format version is not something readers can opt out of afterwards.

saying these in an interview costs you the question

  • Calls it a table type rather than a per-operation property
  • Claims merge-on-read makes both writes and reads faster
  • Thinks delete files disappear on their own without compaction
  • Confuses it with Hudi's Merge-on-Read table type
  • Assumes it works on a format v1 table

context