skip to content

Apache Iceberg

You will learn Apache Iceberg, the open table format that has become the neutral default for multi-engine lakehouses: its layered metadata spec, hidden partitioning with in-place partition and schema evolution, and the maintenance procedures that keep a table healthy. Interviewers ask about Iceberg because it is now the format most vendors agreed on.

on this pageshow

explore

questions

18

In Apache Iceberg, which files does the expire_snapshots procedure actually delete?

level: middleimportance: must knowfreq 72%

answer

  1. old commits still hold on to files
  2. age of a file is not the test
  3. reachable from a kept snapshot or not
  4. two phases: drop snapshots, then delete files
  5. retain_last and history.expire.max-snapshot-age-ms

basics

~20 s

It removes snapshot entries older than the retention window from table metadata, then deletes the data files, delete files, manifests and manifest lists that no remaining snapshot references. Anything still reachable from a kept snapshot survives, whatever its age.

solid answer

~40 s

Every Iceberg commit writes a new snapshot, and old snapshots keep pointing at the files they saw, so nothing is reclaimed until snapshots are expired. `expire_snapshots` works in two phases: it drops the selected snapshots from table metadata, then computes the files reachable from the snapshots that remain and physically deletes the data files, delete files, manifest files and manifest lists that only the removed snapshots referenced. **Reachability decides, not file age** — a five-year-old data file the current snapshot still lists is untouched. Arguments are `older_than` and `retain_last`, defaulting to the table properties `history.expire.max-snapshot-age-ms` (5 days) and `history.expire.min-snapshots-to-keep` (1). Snapshots held by a branch or tag are protected by that ref's own retention. It does not remove old `metadata.json` files or uncommitted orphan files — those are separate settings and a separate procedure.

code

sql · 12 lines
sql
-- keep 5 recent snapshots, expire anything older than the cutoff
CALL prod.system.expire_snapshots(
  table => 'db.events',
  older_than => TIMESTAMP '2026-08-14 00:00:00',
  retain_last => 5
);

-- make retention declarative instead of per-call
ALTER TABLE prod.db.events SET TBLPROPERTIES (
  'history.expire.max-snapshot-age-ms' = '604800000',
  'history.expire.min-snapshots-to-keep' = '5'
);

go deeper

for a junior

Recall that Iceberg keeps every past version of a table as a snapshot, and that a cleanup procedure is what eventually frees the old files. Know the name expire_snapshots and that it trades history for storage.

for a middle

Explain the two phases — drop snapshot entries, then delete only files no surviving snapshot references — and name the older_than and retain_last arguments plus the history.expire.* table properties that supply their defaults.

for a senior

Show judgment about the retention window: long enough for the slowest reader and any audit or rollback requirement, short enough to control object-storage cost, with tags pinning snapshots that must outlive it. Be ready to explain why storage only drops after expiry follows compaction.

for a principal

Own retention as policy rather than a job argument: retention tiers set through table properties, tagging conventions for compliance points, and the cost-versus-recoverability tradeoff you are signing the platform up for across every table.

## Why snapshots accumulate Every write to an Iceberg table — append, overwrite, `DELETE`, `MERGE`, or a compaction job — is a commit that produces a new `metadata.json` containing a new snapshot. A snapshot points at one manifest list (`snap-<id>-<seq>-<uuid>.avro`), which points at manifest files, which list the data files and delete files that make up the table at that instant. Nothing is ever mutated in place. That immutability is what gives Iceberg snapshot isolation, time travel and rollback — and it is also why a table that never expires anything keeps every data file it has ever written, plus a metadata file whose snapshot log grows without bound. ## What the procedure does `expire_snapshots` runs in two phases. 1. **Metadata phase.** It commits a new `metadata.json` in which the selected snapshots no longer appear in `snapshots` or in the snapshot log. 2. **Cleanup phase.** It computes the set of files reachable from the snapshots that remain, compares it with the files reachable from the removed snapshots, and physically deletes the difference: data files, delete files (position and equality), manifest files, and manifest lists. The decisive rule is **reachability, not age**. A data file written years ago that the current snapshot still lists is never deleted by expiry. Conversely a data file written five minutes ago can be deleted if the only snapshot that referenced it was just expired — which is exactly what happens to the pre-compaction files after `rewrite_data_files`. ## Arguments and defaults ```sql CALL prod.system.expire_snapshots( table => 'db.events', older_than => TIMESTAMP '2026-08-14 00:00:00', retain_last => 5 ) ``` `older_than` defaults to now minus `history.expire.max-snapshot-age-ms`, whose table default is 5 days. `retain_last` defaults to `history.expire.min-snapshots-to-keep`, whose default is 1, and acts as a **floor**: that many recent ancestors of the current snapshot are kept even if they are older than the cutoff. The procedure also accepts `snapshot_ids` to expire specific snapshots, `max_concurrent_deletes` to parallelise the file deletions, and `stream_results` so the driver does not collect the whole delete list. The Java equivalent is `table.expireSnapshots().expireOlderThan(ts).retainLast(n).commit()`. Setting the two `history.expire.*` properties on the table makes retention declarative, so a generic maintenance job needs no per-table arguments. Branches and tags protect what they reference. A snapshot held by a tag is retained under that ref's retention (`max-ref-age-ms`, and per-ref minimum snapshots and age), so tagging a month-end snapshot is the supported way to keep one specific point in time far beyond the general window. ## What it does not clean up - **Old `metadata.json` files.** They are governed by `write.metadata.delete-after-commit.enabled` (default `false`) and `write.metadata.previous-versions-max` (default 100). A high-frequency streaming table with the default off accumulates a metadata file per commit. - **Orphan files** — files under the table location that no metadata ever referenced, typically left by a failed or killed writer. Expiry only walks metadata, so it cannot see them; `remove_orphan_files` lists the storage location and diffs it against metadata. ## What expiry costs you Once a snapshot is expired, everything that depended on it is gone: time travel `AS OF` that snapshot id or a timestamp inside the expired range fails, `rollback_to_snapshot` to it fails, and an incremental read whose start snapshot has been expired fails. A reader that resolved an older snapshot before the job ran and is still scanning can hit missing files, which is why the retention window should comfortably exceed your longest-running query and your slowest downstream consumer's lag. Choose the window from real requirements — audit and recovery needs — rather than copying the 5-day default blindly. ## Operating it Inspect state through metadata tables: `SELECT * FROM db.events.snapshots` for what exists and when, `db.events.refs` for branches and tags, `db.events.metadata_log_entries` for metadata file history. The common surprise is "we compacted and storage went **up**" — compaction writes new files while the pre-compaction files stay reachable from the previous snapshots, so space is only reclaimed at expiry. That fixes the order of operations: rewrite first, then expire.

  • We ran compaction and storage grew instead of shrinking. Why?
    Compaction writes new, larger data files and commits a new snapshot, but the previous snapshots still reference the small files they replaced, so both sets are on disk. Iceberg cannot delete the originals while any retained snapshot points at them. Space is reclaimed only when `expire_snapshots` removes those snapshots. That is why maintenance runs rewrite first, expiry second.
  • How do you keep one specific snapshot far longer than the retention window?
    Tag it. `ALTER TABLE ... CREATE TAG` pins a snapshot id under a named ref, and expiry will not remove a snapshot that a branch or tag still references, subject to that ref's own retention settings such as `max-ref-age-ms`. This is the supported way to hold a month-end or pre-migration state for audit while everything else expires on the normal schedule.
  • Does expire_snapshots remove old metadata.json files?
    No. Metadata file retention is separate: `write.metadata.delete-after-commit.enabled` is `false` by default, so previous `metadata.json` versions are kept indefinitely, and `write.metadata.previous-versions-max` (default 100) caps how many are retained once you turn deletion on. On a frequently committing table this is a real source of metadata bloat that snapshot expiry will never touch.

saying these in an interview costs you the question

  • Says it deletes every data file older than the cutoff
  • Confuses it with removing uncommitted orphan files
  • Believes time travel still works after a snapshot expires
  • Assumes compaction alone reclaims storage without expiry
  • Thinks it also prunes old metadata.json versions

context

open as a page

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

level: middleimportance: must knowfreq 66%

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.

open as a page

What happens to existing Iceberg data files when you ALTER TABLE ADD PARTITION FIELD?

level: middleimportance: must knowfreq 72%

basics

~20 s

Nothing is rewritten. Iceberg appends a new partition spec to table metadata and points new writes at it; files already written keep the spec id they were written under, and scan planning handles each spec separately.

open as a page

In Apache Iceberg, what does hidden partitioning change about how a partitioned table is queried?

level: middleimportance: must knowfreq 82%

basics

~10 s

Iceberg derives partition values from a transform on a real column, so queries filter that column and still prune. Hive-style tables need a separate derived column that every reader must remember to filter on.

open as a page

In Apache Iceberg, how does a reader reach the data files a query must scan?

level: middleimportance: must knowfreq 80%

basics

~20 s

The catalog names the current metadata file. That JSON gives the current snapshot, whose manifest-list Avro file names the manifests. Each manifest lists data files with partition values and column bounds, so the reader prunes manifests then files and scans only survivors.

open as a page

Two Spark jobs commit to one Iceberg table at the same moment — what happens?

level: seniorimportance: must knowfreq 60%

basics

~20 s

Both write their new metadata files, then each asks the catalog to swap the table pointer from the base metadata they read to their own. The swap is an atomic compare-and-swap, so exactly one wins; the loser refreshes, re-applies its changes on the new base, and retries.

open as a page

In an Iceberg table PARTITIONED BY (days(event_ts)), which WHERE clause actually prunes partitions?

level: juniorimportance: should knowfreq 60%

basics

~20 s

A predicate on event_ts itself, such as a timestamp range or equality. Iceberg projects it through the days transform onto each file's stored partition value. There is no event_ts_day column to filter, and wrapping event_ts in a function can block the projection.

open as a page

What does an Apache Iceberg table store under its metadata/ and data/ directories?

level: juniorimportance: should knowfreq 62%

basics

~20 s

An Iceberg table's data/ directory holds the immutable data files (Parquet by default, optionally ORC or Avro). Its metadata/ directory holds the JSON metadata files, the Avro manifest lists named snap-*.avro, and the Avro manifest files that describe those data files.

open as a page

Why can a column in an Iceberg table be renamed without rewriting any data files?

level: middleimportance: should knowfreq 55%

basics

~20 s

Iceberg tracks every column by a unique integer field id that is recorded in the table schema and written into the data files. Readers resolve columns by id, so a name lives only in metadata and a rename touches nothing on disk.

open as a page

When an Apache Iceberg append adds one file, which metadata files are written?

level: middleimportance: should knowfreq 58%

basics

~20 s

An append writes the data file, a new manifest listing it, a new manifest list that references that manifest plus the parent snapshot's still-valid manifests, and a new JSON metadata file carrying the new snapshot. Existing manifests and data files are reused untouched.

open as a page

Which statistics does an Iceberg manifest store, and how do they prune a scan?

level: middleimportance: should knowfreq 55%

basics

~20 s

Each Iceberg manifest entry stores per-file record_count, file_size_in_bytes, column_sizes, value_counts, null_value_counts, nan_value_counts and per-column lower_bounds and upper_bounds. Planning compares the predicate against those bounds and drops files that cannot match, before opening any data file.

open as a page

Why does Iceberg's remove_orphan_files procedure default to a three-day age threshold?

level: seniorimportance: should knowfreq 46%

basics

~20 s

Because it finds candidates by listing the storage location and subtracting everything the metadata references, files written by a still-running job look identical to orphans. The conservative default window keeps it from deleting data an in-flight commit is about to reference.

open as a page

After rewrite_data_files, an Iceberg table still has thousands of small files. Why?

level: seniorimportance: should knowfreq 56%

basics

~20 s

Usually the job skipped most groups: file groups are built per partition, and a partition with fewer than min-input-files candidates or with files already near the target size is left alone. A where filter, an unchanged small target size, or fresh ingest after the run explain the rest.

open as a page

In an Iceberg table partitioned by bucket(16, user_id), why does a range filter on user_id scan every bucket?

level: seniorimportance: should knowfreq 45%

basics

~20 s

Bucketing hashes the value, so ordering is destroyed and neighbouring ids land in unrelated buckets. Iceberg can project only equality and IN predicates through the bucket transform; a range predicate matches every bucket, so nothing is pruned.

open as a page

What do Iceberg format versions 1, 2 and 3 each change about table metadata?

level: seniorimportance: should knowfreq 45%

basics

~20 s

Version 1 tracks data files only. Version 2 adds row-level delete files, marks manifests as data or delete content, and adds sequence numbers that order commits. Version 3 adds deletion vectors stored in Puffin blobs, row lineage, default column values and new data types.

open as a page

How would you design a maintenance policy for hundreds of Iceberg tables on one platform?

level: principalimportance: should knowfreq 38%

basics

~20 s

Classify tables into a few tiers by write pattern, express retention and target file size as table properties so one generic job serves all of them, run rewrite then expire then occasional orphan removal, and budget the compute compaction consumes.

open as a page

In Iceberg v2, how do equality delete files differ from position delete files at read time?

level: seniorimportance: nice to knowfreq 36%

basics

~20 s

A position delete names a data file and the row offsets to drop, so a reader skips exactly those rows. An equality delete names column values, so the reader must evaluate them against every candidate data file in the partition — far more expensive.

open as a page

In Iceberg, what does ALTER TABLE ... WRITE ORDERED BY change about later writes?

level: seniorimportance: nice to knowfreq 30%

basics

~10 s

It records a sort order in table metadata and makes it the table's default, so engines writing afterwards sort rows to match it. Existing files are untouched and readers do not enforce it.

open as a page