DEV Community

Cover image for Iceberg Table Maintenance: Compaction, Expire Snapshots, Rewrite Manifests & Orphan Files
Gowtham Potureddi
Gowtham Potureddi

Posted on

Iceberg Table Maintenance: Compaction, Expire Snapshots, Rewrite Manifests & Orphan Files

iceberg table maintenance is the set of housekeeping jobs that keep an Apache Iceberg table fast to read, cheap to store, and correct over time — because Iceberg never edits data in place. Every insert, update, delete, and MERGE writes brand-new files and commits a brand-new snapshot that layers on top of the old one. That immutable, snapshot-per-commit design is exactly what buys you atomic writes, time travel, and safe concurrent readers — and it is also exactly why a table left unattended slowly rots into thousands of tiny files, a mile-high stack of dead snapshots, and a metadata layer so bloated that query planning takes longer than the query.

Maintenance is how you pay that design tax back down. There are four canonical jobs and every serious lakehouse runs them on a schedule: compaction (rewrite_data_files) packs small files into large ones; expiring snapshots (expire_snapshots) drops old snapshots so their data files can actually be deleted; rewriting manifests (rewrite_manifests) consolidates the metadata index so scan planning stays quick; and removing orphan files (remove_orphan_files) sweeps up files that failed jobs left behind. Merge-on-read tables add a fifth, rewrite_position_delete_files, to compact delete files. This guide walks each one the way an interviewer will probe it — what it does, the exact Spark CALL, the knobs that matter — and pairs every section with a Solution-Tail answer: code, a step-by-step trace, an output table, then a concept-by-concept breakdown of why it works.

PipeCode blog header for Iceberg table maintenance — bold white headline 'Iceberg Table Maintenance' with subtitle 'Compaction, Expire Snapshots, Manifests, Orphans' and a stylised many-small-files-into-few-large-files scene on a dark gradient with purple, green, orange, and blue accents and a small pipecode.ai attribution.

When you want hands-on reps immediately after reading, drill file-layout and read-cost problems on the optimization practice library →, rehearse partition-pruning decisions on the partitioning practice set →, and harden your merge and dedupe logic on the deduplication practice set →.


On this page


1. Why Iceberg tables need maintenance

Iceberg never mutates a file in place — that one fact creates every maintenance job you will ever run

The one-sentence invariant: an Iceberg write never edits existing files; it writes new data files and commits a new snapshot that points at the current set of files. Every good thing about Iceberg — atomic commits, snapshot isolation, time travel, safe schema evolution — falls out of that immutability. And so does every maintenance job. If files were edited in place there would be nothing to compact, no dead snapshots to expire, and no orphans to sweep. Maintenance exists precisely because the table grows by accretion.

What a single write actually produces.

  • New data files. An insert or the "new rows" half of a MERGE writes fresh Parquet/ORC/Avro data files. Nothing overwrites the old ones.
  • New delete files (merge-on-read). An update or delete on a MoR table does not rewrite the touched data file; it writes a position or equality delete file that logically masks rows, to be reconciled at read time.
  • New manifests and a manifest list. Each snapshot has a manifest list pointing at manifest files, and each manifest lists data/delete files with their partition and column stats. Every commit adds to this metadata tree.
  • A new metadata.json. The table's root pointer moves to a new metadata file that records the new snapshot; the previous metadata files linger unless configured to be cleaned.

The three ways an unmaintained table degrades.

  • The small-files problem. Streaming ingestion, frequent micro-batches, or many partitions each getting a trickle of rows produce thousands of tiny files. Readers pay a fixed per-file cost (open, seek, footer read, task scheduling), so 10,000 files of 1 MB are dramatically slower and more expensive to scan than 40 files of 256 MB — even though the bytes are identical.
  • Snapshot & metadata growth. Every snapshot you keep pins the data files it references, so storage only ever grows until you expire snapshots. Meanwhile manifests multiply and metadata.json history stacks up, inflating planning time.
  • Delete-file accumulation (MoR). On merge-on-read tables, delete files pile up beside data files. Reads get slower because every scan must merge deletes, and dangling deletes (deletes whose data file is gone) waste I/O.

The four (plus one) canonical jobs.

  • Compaction — rewrite_data_files. Rewrites many small files into fewer target-sized files; optionally sorts or z-orders them.
  • Expire snapshots — expire_snapshots. Removes old snapshots and physically deletes the files only those snapshots referenced.
  • Rewrite manifests — rewrite_manifests. Consolidates and re-clusters the manifest layer so planning stays cheap.
  • Remove orphan files — remove_orphan_files. Deletes files under the table's location that no metadata references (failed-job debris).
  • Rewrite position deletes — rewrite_position_delete_files. Compacts MoR delete files and drops dangling ones.

What interviewers listen for.

  • Do you say "Iceberg is copy-on-write / append-only at the file level, so maintenance is not optional" early? — senior signal.
  • Do you separate logical retention (snapshots) from physical cleanup (orphans) rather than conflating them? — required framing.
  • Do you know that compaction alone does not free storage — the old small files only disappear when their snapshots expire? — the trap most people miss.
  • Do you propose a scheduled order (compact → rewrite manifests → expire → remove orphans) instead of running jobs ad hoc? — senior signal.

Worked example — one table, four commits, a file explosion

Detailed explanation. The fastest way to feel why maintenance matters is to count files after a handful of ordinary writes. A table taking one small append every few minutes does not "fill up" a file — each commit creates its own new files. Four appends of a few rows each leave you with four data files (plus four manifests and four snapshots), and a streaming job doing this every minute leaves you with thousands per day. Nothing went wrong; this is Iceberg working as designed. Maintenance is the counterweight.

Question. A table db.events receives four separate append commits, each writing one small data file. How many data files, snapshots, and (roughly) manifests exist afterward, and what has to happen for the file count to come back down?

Input.

commit rows appended data files written snapshot created
c1 500 1 S1
c2 400 1 S2
c3 620 1 S3
c4 300 1 S4

Code.

-- Four ordinary appends; each is its own snapshot and its own file(s)
INSERT INTO db.events VALUES /* batch 1 */ ... ;   -- -> snapshot S1
INSERT INTO db.events VALUES /* batch 2 */ ... ;   -- -> snapshot S2
INSERT INTO db.events VALUES /* batch 3 */ ... ;   -- -> snapshot S3
INSERT INTO db.events VALUES /* batch 4 */ ... ;   -- -> snapshot S4

-- Inspect what accumulated
SELECT count(*) AS data_files FROM db.events.files;      -- 4
SELECT count(*) AS snapshots  FROM db.events.snapshots;  -- 4
Enter fullscreen mode Exit fullscreen mode

Step-by-step explanation. Each INSERT is an atomic commit that writes at least one new data file and a new snapshot pointing at the full current file set. Iceberg exposes metadata tables — db.events.files, db.events.snapshots, db.events.manifests — so you can query the accumulation directly. After four commits you have four data files (assuming one file each), four snapshots, and roughly four manifests. To shrink the file count you must run compaction; to actually reclaim the storage of the pre-compaction files you must also expire the snapshots that still reference them.

Output.

metric after four commits to reduce it, run
data files 4 (grows unbounded with commits) rewrite_data_files (compaction)
snapshots 4 (each pins its files) expire_snapshots
manifests ~4 rewrite_manifests
untracked files from failed jobs 0 here, but grows in prod remove_orphan_files

Rule of thumb. File count grows with the number of commits, not the volume of data — a chatty streaming pipeline is the classic small-files generator, and it needs compaction scheduled proportional to commit frequency, not data size.


2. Compaction with rewrite_data_files

rewrite_data_files packs small files into target-sized ones — bin-pack for size, sort and z-order for read locality

Compaction is the maintenance job you run most often, and the interviewer's favourite because it has real strategy choices. The Spark procedure system.rewrite_data_files reads a set of data files and rewrites them into fewer, larger files, committing the result as a new snapshot (it does not mutate the old files). Which files it selects and how it lays the rows out inside the new files depends on the strategy and the options you pass.

The three strategies.

  • binpack (default). Purely a size operation: it bin-packs small files up to target-file-size-bytes without changing row order. Cheapest strategy, no shuffle, fixes the small-files problem directly. Reach for it first.
  • sort. Rewrites files sorted by one or more columns (the table's sort order, or an explicit sort_order). Clustering rows by a frequently-filtered column tightens per-file min/max stats so the reader can skip more files. Costs a shuffle.
  • sort with zorder(...). A space-filling-curve sort that co-locates rows across multiple dimensions at once. Use when queries filter on several columns interchangeably (e.g. country and event_date), where a single linear sort would help one column at the cost of the other.

The options that matter.

  • target-file-size-bytes. The size compaction aims for (defaults to the table's write.target-file-size-bytes, 512 MB). Files near this size are left alone.
  • min-input-files. Minimum number of small files in a group before it is worth rewriting (default 5). Stops compaction churning over already-healthy partitions.
  • min-file-size-bytes / max-file-size-bytes. Define what counts as "too small" (or too big) and thus eligible; by default anything not within 25%–175% of target is a candidate.
  • delete-file-threshold. For MoR, rewrite a data file once it has at least this many delete files attached, even if it is already large — this is how compaction applies deletes and reclaims masked rows.
  • rewrite-all. Force every file to be rewritten regardless of size (e.g. after changing sort order).
  • partial-progress.enabled / partial-progress.max-commits. Commit compaction in several smaller snapshots instead of one huge atomic commit, so a long-running job makes durable progress and is friendlier to concurrent writers.

How file selection works.

  • Compaction plans per partition: within each partition it groups eligible files into file-groups no larger than max-file-group-size-bytes, and each group becomes one Spark job producing target-sized outputs.
  • Because it is a new snapshot, the old small files are still referenced by older snapshots — so compaction on its own does not shrink storage; expiry does that later.

Iconographic Iceberg compaction diagram — a swarm of many tiny data files on the left, a rewrite_data_files engine in the centre with bin-pack, sort, and z-order strategy chips, and a few large target-sized files on the right written as a new snapshot.

Worked example — bin-pack 128 small files into a handful of large ones

Detailed explanation. The canonical compaction is a plain bin-pack: a partition has accumulated 128 files averaging 4 MB from a streaming job, and you want them packed to ~512 MB targets. Bin-pack needs no sort order and no shuffle — it simply concatenates rows from small files into new large files — so it is the cheapest and most common maintenance call.

Question. A partition of db.events holds 128 files at ~4 MB (512 MB of data). Compact them to 512 MB targets with a minimum of 5 input files per group. How many output files result, and how do you invoke it?

Input.

partition input files avg size total bytes
dt=2026-09-14 128 4 MB ~512 MB

Code.

CALL spark_catalog.system.rewrite_data_files(
  table    => 'db.events',
  strategy => 'binpack',
  where    => 'dt = ''2026-09-14''',
  options  => map(
    'target-file-size-bytes', '536870912',   -- 512 MB
    'min-input-files',        '5'
  )
);
Enter fullscreen mode Exit fullscreen mode

Step-by-step explanation. The procedure filters to the dt = '2026-09-14' partition with the where predicate, so only that partition is touched. It finds 128 files well below target, groups them (128 ≥ min-input-files), and bin-packs their rows into new files of ~512 MB. 512 MB of data at a 512 MB target yields roughly one output file (Iceberg leaves a little headroom, so you may see one or two). The procedure returns counts of rewritten and added files and commits them as a new snapshot; the 128 originals remain referenced by prior snapshots until you expire them.

Output.

result field value
rewritten_data_files_count 128
added_data_files_count 1
rewritten_bytes ~512 MB
snapshot new snapshot committed

Rule of thumb. Start with binpack on the partitions your ingestion touches most; only escalate to sort/zorder when profiling shows reads are file-skipping-bound, because those add a shuffle you pay for every compaction.

Iceberg interview question on compaction strategy

Question. A 4 TB events table is queried almost entirely with WHERE country = ? AND event_date = ?, but scans still read far too many files because rows for a country are scattered across every file. The table already has healthy ~512 MB files, so plain bin-pack changes nothing. How do you lay the data out so the reader skips files, and what exactly does the compaction call look like?

Solution Using sort compaction with z-order

Code.

-- Cluster rows across BOTH filter columns at once with a Z-order curve.
CALL spark_catalog.system.rewrite_data_files(
  table       => 'db.events',
  strategy    => 'sort',
  sort_order  => 'zorder(country, event_date)',
  options     => map(
    'target-file-size-bytes', '536870912',
    'rewrite-all',            'true',        -- files are already ~512MB, force rewrite
    'max-concurrent-file-group-rewrites', '8',
    'partial-progress.enabled', 'true',
    'partial-progress.max-commits', '10'
  )
);
Enter fullscreen mode Exit fullscreen mode

Step-by-step trace.

stage what happens effect on file skipping
plan group all files per partition; rewrite-all marks every file eligible full relayout, not just small files
shuffle rows repartitioned along the zorder(country, event_date) curve rows near in both dims land in the same file
write new ~512 MB files, each with tight country/event_date min-max stats most files now exclude a given (country, date)
commit committed in up to 10 partial snapshots long job makes durable progress
  1. Plain binpack cannot help here — the files are already the right size; the problem is the arrangement of rows inside them, so you need a sort-family strategy.
  2. A single linear sort_order => 'country, event_date' would cluster country well but scatter event_date within each country. Because queries filter on both interchangeably, a z-order curve co-locates rows that are close in both dimensions.
  3. After the rewrite, each file's min/max footer stats for country and event_date are narrow, so Iceberg's manifest-level filtering prunes files that cannot match the predicate before any data is read.
  4. partial-progress.enabled lets a multi-terabyte rewrite commit in stages, so a failure late in the job does not throw away all completed work.

Output:

metric before z-order after z-order
files read for country='DE' AND event_date='2026-09-14' ~7,800 (all) ~120 (skipped by stats)
bytes scanned per query full 4 TB region small fraction
file size distribution ~512 MB ~512 MB (unchanged)

Why this works — concept by concept:

  • Bin-pack vs sort — bin-pack fixes file size; sort/z-order fix file content ordering. When files are already sized right, only a content relayout improves skipping.
  • Z-order locality — a space-filling curve interleaves multiple columns so a range filter on any of them still touches a compact set of files, which a single-column sort cannot deliver for more than one column.
  • Min-max pruning — Iceberg keeps per-file column bounds in manifests; tight bounds after clustering let planning eliminate files without opening them, turning a full scan into a targeted one.
  • Partial progress — committing in several snapshots bounds the blast radius of a failure and reduces contention with concurrent writers on a long rewrite.
  • Cost — sort/z-order compaction costs a full shuffle, roughly O(bytes rewritten · log) sort work, paid once, in exchange for O(files skipped) savings on every subsequent query.

Optimization
Topic — optimization
File-layout and read-cost optimization problems

Practice →

Partitioning Topic — partitioning Partitioning and clustering-for-pruning problems

Practice →


3. Expiring snapshots to reclaim storage

expire_snapshots is the only job that frees disk — it deletes files no kept snapshot still references

Compaction makes reads fast but, on its own, makes storage worse: it writes new large files while the old small files stay referenced by older snapshots. expire_snapshots is the job that actually reclaims space — it removes old snapshots from the table's history and physically deletes the data files, delete files, and manifests that no remaining snapshot references. This is the maintenance job people most often forget, and the reason a "compacted" table can still be paying for three copies of its data.

The mechanism.

  • Pick which snapshots to drop. older_than sets an age cutoff; retain_last sets a minimum number of recent snapshots to always keep. A snapshot is removed only if it is both older than the cutoff and not within the last retain_last.
  • Reference-count the files. Iceberg computes the set of files reachable from the surviving snapshots. Any file referenced only by expired snapshots is now unreachable.
  • Delete the unreachable files. Those data/delete/manifest files are physically deleted from object storage. Files still shared with a kept snapshot are never touched.
  • Update history. The expired snapshots vanish from db.tbl.snapshots and .history, so you can no longer time-travel to them.

The knobs that matter.

  • older_than — a timestamp; expire snapshots created before it. Omitting it uses the table's default max-snapshot-age.
  • retain_last — a hard floor on how many recent snapshots survive regardless of age; guards against expiring everything.
  • snapshot_ids — expire specific snapshots by id (surgical cleanup).
  • max_concurrent_deletes — parallelism for the physical file deletes, which dominate wall-clock time on large expirations.
  • stream_results — stream the file list to the driver instead of collecting it, for very large expirations.

The table properties that automate it.

  • history.expire.max-snapshot-age-ms — default 5 days; snapshots older than this are eligible when engines auto-expire.
  • history.expire.min-snapshots-to-keep — default 1; the always-retained floor.
  • history.expire.max-ref-age-ms — controls expiry of named branches/tags (except main).

Failure modes interviewers probe.

  • Retention vs time travel. Expiring aggressively saves money but shrinks your rollback window — you can only revert to, or read AS OF, snapshots you kept. Set retention to your recovery-objective, not to zero.
  • Order with orphan removal. Expiry only deletes files it can prove are unreachable through metadata. Files that were never in any snapshot (failed writes) are invisible to it — that is what remove_orphan_files is for.

Iconographic Iceberg expire-snapshots diagram — a vertical stack of snapshot layers with a cut-off line at older_than, kept recent snapshots highlighted, expired old snapshots and their now-unreferenced data files swept to a trash glyph, and a retain_last guard chip.

Worked example — expire everything older than seven days, keep the last five

Detailed explanation. The everyday expiry policy pairs an age cutoff with a retention floor: "drop anything older than a week, but never leave fewer than five snapshots so recent rollback still works." The two conditions are ANDed, so a snapshot survives if it is either recent enough or inside the retained window.

Question. A table has 40 snapshots spanning 30 days. Expire snapshots older than 7 days while always keeping at least the 5 most recent. What is deleted, and what is the call?

Input. Snapshots S1 (oldest, 30 days) … S40 (newest, minutes ago); the last 5 are S36–S40.

Code.

CALL spark_catalog.system.expire_snapshots(
  table       => 'db.events',
  older_than  => TIMESTAMP '2026-09-08 00:00:00',   -- 7 days before now
  retain_last => 5,
  max_concurrent_deletes => 8
);
Enter fullscreen mode Exit fullscreen mode

Step-by-step explanation. Iceberg lists all snapshots older than 2026-09-08 — say S1 through S31. It then applies retain_last => 5, which protects S36–S40 (already newer than the cutoff, so no conflict here). The survivors are S32–S40. Iceberg reference-counts files from those nine survivors; any data/manifest file referenced only by S1–S31 is unreachable and physically deleted, using 8 parallel delete threads. Files that S32+ still share (for example, large compacted files that predate the cutoff but are still current) are kept.

Output.

result field value
deleted_data_files_count files unique to S1–S31
deleted_manifest_files_count manifests unique to S1–S31
snapshots remaining S32–S40 (9)
earliest time-travel target S32

Rule of thumb. Set older_than/retain_last from your recovery objective, then let a scheduled expiry enforce it — never run expiry with no retain_last, or a bad recent write plus a fast expiry can leave you nothing to roll back to.

Iceberg interview question on retention correctness

Question. You compacted a table at 02:00 (new snapshot S_new referencing large files), then immediately ran expire_snapshots with older_than = now and retain_last = 1 to "save money." A reader running a long query that started at 01:59 against snapshot S_old fails or reads garbage. What went wrong, and how should retention be set so compaction + expiry is safe?

Solution Using an age-based retention window

Code.

-- WRONG: expires the snapshot a long-running reader/writer may still use
-- CALL ... expire_snapshots(table => 'db.events', older_than => now(), retain_last => 1);

-- RIGHT: keep a window wider than your longest reader / rollback objective
CALL spark_catalog.system.expire_snapshots(
  table       => 'db.events',
  older_than  => TIMESTAMP '2026-09-12 00:00:00',   -- 3 days ago, not "now"
  retain_last => 10
);
Enter fullscreen mode Exit fullscreen mode

Step-by-step trace.

policy S_old (pre-compaction) long reader on S_old rollback window
older_than => now, retain_last => 1 expired instantly, files deleted reads deleted files → fails none
older_than => 3 days, retain_last => 10 retained completes safely ~3 days / 10 snapshots
  1. expire_snapshots deletes files that are unreachable from surviving snapshots — but a reader that already resolved S_old still expects S_old's files to exist on disk.
  2. Setting older_than => now() makes S_old eligible the instant after compaction, so its exclusive files (the small pre-compaction files) get deleted out from under the in-flight reader.
  3. The fix is to keep a retention window wider than your longest-running query and your rollback objective — expiry should target snapshots from hours or days ago, never the immediate past.
  4. retain_last => 10 adds a count-based floor so even a burst of commits cannot shrink history below a safe number of rollback points.

Output:

metric reckless policy safe policy
in-flight reader crashes on missing files completes
storage reclaimed maximal, unsafe slightly less, safe
rollback available 0 snapshots 10 snapshots / ~3 days

Why this works — concept by concept:

  • Reference-counted deletion — expiry deletes a file only when no surviving snapshot references it, so retaining more snapshots directly protects the files in-flight readers need.
  • older_than vs retain_last — one is an age gate, the other a count floor; they are ANDed, and using both prevents either "too old but only copy" or "recent but too many" edge cases.
  • Reader/writer safety window — snapshots must outlive the longest consumer that could still resolve them; expiry against "now" races live readers.
  • Rollback objective — retention is your recovery policy; the storage you keep is the price of being able to undo a bad write.
  • Cost — physical deletes dominate: O(unreachable files) delete calls to object storage, parallelised by max_concurrent_deletes; metadata work is negligible by comparison.

Optimization
Topic — optimization
Storage-reclaim and retention-policy problems

Practice →

Deduplication Topic — deduplication Snapshot, versioning and dedupe-history problems

Practice →


4. Rewriting manifests & metadata

rewrite_manifests keeps the metadata index small and partition-clustered — so scan planning stays fast

Compaction fixes the data layer; the metadata layer needs its own housekeeping. Every commit adds manifest files, and a table with tens of thousands of manifests spends real time in query planning — before a single data byte is read — just walking manifests to find candidate files. system.rewrite_manifests consolidates many small manifests into fewer large ones and, crucially, re-clusters manifest entries by partition so that planning can skip whole manifests whose partition range cannot match a predicate.

Why manifest layout matters.

  • Planning reads manifests, not data. To build a scan, Iceberg reads the manifest list, then the manifests, filtering files by partition and column stats. Thousands of tiny, partition-mixed manifests make this step slow.
  • Clustered manifests prune in bulk. If each manifest holds entries for one partition (or a tight partition range), a predicate like dt = '2026-09-15' lets the planner discard entire manifests without inspecting their file entries.
  • Small manifests come from small commits. The same chatty ingestion that creates small data files creates small manifests, so tables that need frequent compaction usually need periodic manifest rewrites too.

The procedure and its knobs.

  • rewrite_manifests('db.tbl'). The basic call; rewrites manifests for the current snapshot into target-sized, partition-clustered manifests.
  • use_caching. Cache manifest entries in Spark during the rewrite (default true); disable if you hit memory pressure on huge manifest sets.
  • spec_id. Restrict the rewrite to manifests for a specific partition spec, useful after partition evolution left mixed specs.

The table properties that govern manifests and metadata.

  • commit.manifest.target-size-bytes — target size for rewritten/merged manifests (default 8 MB).
  • commit.manifest.min-count-to-merge — how many new manifests trigger an automatic merge at commit time.
  • commit.manifest-merge.enabled — whether Iceberg merges manifests automatically on write (default true); some high-throughput writers disable it and rely on scheduled rewrite_manifests instead.
  • write.metadata.delete-after-commit.enabled and write.metadata.previous-versions-max — automatically prune old metadata.json files so the metadata directory does not grow without bound.

Iconographic Iceberg rewrite-manifests diagram — many small scattered manifest files on the left, a rewrite_manifests engine clustering entries by partition in the centre, and a few large manifests grouped by partition on the right enabling better partition pruning.

Worked example — collapse 900 tiny manifests into a few partition-clustered ones

Detailed explanation. A streaming table committing every minute for two weeks can accumulate hundreds of manifests, each listing a handful of files spread across many partition-days. Query planning slows because the planner opens nearly every manifest. rewrite_manifests rebuilds them into a small number of large manifests, each holding entries for a contiguous partition range.

Question. A table db.events has 900 manifests averaging 20 KB, entries mixed across partitions. Rewrite them and describe the planning improvement.

Input.

metric before
manifest count 900
avg manifest size ~20 KB
partition clustering mixed (each manifest spans many days)

Code.

CALL spark_catalog.system.rewrite_manifests('db.events');

-- verify
SELECT count(*) AS manifests FROM db.events.manifests;
Enter fullscreen mode Exit fullscreen mode

Step-by-step explanation. The procedure reads all current-snapshot manifest entries, sorts them by partition, and writes new manifests up to commit.manifest.target-size-bytes (8 MB). 900 tiny manifests holding, say, 90,000 file entries collapse into a handful of large manifests, each covering a tight partition range. A subsequent query with WHERE dt = '2026-09-15' reads the manifest list, sees each manifest's partition bounds, and opens only the one or two manifests whose range includes that day — instead of all 900.

Output.

metric before after
manifest count 900 ~5
manifests opened for one-day query ~900 ~1–2
planning time seconds milliseconds

Rule of thumb. If EXPLAIN shows planning time dominating short queries, or the manifests metadata table has thousands of tiny entries, schedule rewrite_manifests — it is cheap and touches no data files.

Iceberg interview question on metadata growth

Question. A Flink job writes to an Iceberg table every 10 seconds. After a month, both query planning and the metadata.json directory have ballooned — planning takes seconds and object storage lists thousands of metadata files. Reads are slow before they even start. How do you bring both the manifest layer and the metadata-file history back under control?

Solution Using rewrite_manifests plus metadata table properties

Code.

-- 1) Automatically cap the metadata.json history so it stops growing
ALTER TABLE db.events SET TBLPROPERTIES (
  'write.metadata.delete-after-commit.enabled' = 'true',
  'write.metadata.previous-versions-max'       = '20',
  'commit.manifest.target-size-bytes'          = '8388608',   -- 8 MB
  'commit.manifest-merge.enabled'              = 'true'
);

-- 2) Consolidate the accumulated manifests now
CALL spark_catalog.system.rewrite_manifests('db.events');
Enter fullscreen mode Exit fullscreen mode

Step-by-step trace.

layer symptom fix applied result
metadata.json thousands of old versions listed delete-after-commit.enabled + previous-versions-max=20 keeps only last 20 per commit
manifests 900+ tiny, partition-mixed rewrite_manifests few large, partition-clustered
new commits still create small manifests manifest-merge.enabled=true merged inline going forward
  1. Two different problems hide behind "slow planning": too many manifests (metadata index) and too many metadata.json files (root-pointer history). They need different fixes.
  2. write.metadata.delete-after-commit.enabled with previous-versions-max = 20 makes each commit prune old metadata.json files, so the metadata directory stops growing unbounded.
  3. rewrite_manifests fixes the existing manifest bloat once; commit.manifest-merge.enabled = true keeps future commits from re-accumulating tiny manifests by merging them at write time.
  4. None of this touches data files, so it is safe to run on a live table and composes with a separate compaction schedule.

Output:

metric before after
metadata.json files retained thousands 20
manifest count 900+ ~5
planning time (short query) ~2 s ~30 ms

Why this works — concept by concept:

  • Manifest consolidation — fewer, larger, partition-clustered manifests mean planning opens a handful of files instead of thousands, cutting fixed pre-scan latency.
  • Metadata history capping — delete-after-commit + previous-versions-max bound the metadata.json chain so a high-frequency writer does not leave a million root-pointer files behind.
  • Inline merge vs scheduled rewrite — manifest-merge.enabled prevents regrowth at commit time; scheduled rewrite_manifests cleans up whatever still accumulates, and high-throughput writers often turn inline merge off to keep commits cheap and lean on the schedule.
  • Planning is pre-scan cost — this work pays off before a single data byte is read, which is why it dominates short, selective queries specifically.
  • Cost — rewrite_manifests is O(file entries) metadata I/O with no data shuffle; the property changes are O(1) and apply going forward.

Optimization
Topic — optimization
Metadata and query-planning optimization problems

Practice →

Partitioning Topic — partitioning Partition-pruning and manifest-clustering problems

Practice →


5. Orphan files & merge-on-read delete compaction

remove_orphan_files sweeps untracked debris; rewrite_position_delete_files compacts MoR deletes

Two final jobs close the loop. Orphan files are files sitting under the table's storage location that no table metadata references — the debris of jobs that wrote files then failed before committing, or of retried tasks. Expiry cannot touch them (expiry only follows metadata, and orphans are outside the metadata graph). system.remove_orphan_files is the only job that lists the physical storage location, diffs it against every file reachable from metadata, and deletes the difference — behind a mandatory time guard so it can never delete a file a concurrent commit is still writing.

Separately, merge-on-read tables accumulate position delete files that mask rows in data files, and those must be periodically applied and compacted with system.rewrite_position_delete_files.

Orphan removal — the mechanism.

  • List vs reference. Iceberg enumerates all files under the table location and computes the set reachable from all snapshots' manifests/metadata. Files in the listing but not in the reachable set are orphans.
  • The time guard (older_than, default 3 days). Only files older than this are eligible. This is a safety mechanism, not a retention policy: an in-flight write may have created files seconds ago that are not yet in a committed snapshot, and deleting them would corrupt the commit. The default 3-day window makes accidental deletion of live writes essentially impossible.
  • dry_run => true. List what would be deleted without deleting. Always run this first on a table whose write patterns you do not fully control.
  • location. Scan a non-default location (e.g. a specific data path).

Why orphans are dangerous to over-clean.

  • Orphan removal is the one maintenance job that can delete a live file if misconfigured, because it works by physical listing rather than by following metadata. Respect the time guard, and never point it at a location other jobs write to concurrently.

Merge-on-read delete compaction.

  • rewrite_position_delete_files('db.tbl'). Rewrites and compacts position delete files: it merges many small delete files, and — during data compaction with a delete-file-threshold — deletes are applied so masked rows physically disappear.
  • Dangling deletes. A position delete whose target data file has already been rewritten away is "dangling"; compaction drops these so reads stop paying to evaluate deletes that match nothing.
  • CoW vs MoR. Copy-on-write tables rewrite data files on update (no delete files, heavier writes); merge-on-read defers the cost to read time via delete files, which is why MoR specifically needs delete compaction.

The safe scheduling order.

  • Run compaction → rewrite manifests → expire snapshots → remove orphans. Compaction and manifest rewrites create new files and new snapshots; expiry then drops the now-superseded old snapshots and their files; orphan removal last sweeps anything physical that no snapshot ever claimed. Running orphan removal before expiry wastes work and running it against active writers is the classic foot-gun.

Iconographic Iceberg orphan-files and delete-compaction diagram — a table location containing tracked files plus untracked orphan files from failed jobs on the left, a remove_orphan_files sweep with a 3-day guard in the centre, and merge-on-read position delete files being compacted by rewrite_position_delete_files on the right.

Worked example — dry-run then remove orphans behind a time guard

Detailed explanation. The safe orphan-removal pattern is always two steps: dry-run to see the candidate list, then a real run with a sensible older_than. A failed Spark job left partial files in the table location; they were never committed, so no snapshot references them, and expiry will never remove them.

Question. A crashed ingestion job left ~40 partial files under db.events's location, none referenced by any snapshot and all written over a week ago. Safely delete them without risking any live write.

Input.

file class count referenced by a snapshot? age
tracked data files 1,200 yes mixed
orphan partials 40 no 8 days

Code.

-- Step 1: see what would be deleted (deletes nothing)
CALL spark_catalog.system.remove_orphan_files(
  table   => 'db.events',
  dry_run => true
);

-- Step 2: delete orphans older than 3 days (the safety guard)
CALL spark_catalog.system.remove_orphan_files(
  table      => 'db.events',
  older_than => TIMESTAMP '2026-09-12 00:00:00'
);
Enter fullscreen mode Exit fullscreen mode

Step-by-step explanation. The dry run lists all files under the table location, subtracts every file reachable from metadata, and returns the 40 orphan partials without touching them — you eyeball the list to confirm none are live. The real run applies older_than = 3 days ago, so only files older than the guard (all 40 orphans qualify; any file a current job just wrote would be excluded) are physically deleted. Tracked data files are never candidates because they are reachable from a snapshot.

Output.

result field value
orphan_file_locations (dry run) 40 paths
deleted files (real run) 40
tracked files touched 0

Rule of thumb. Always dry_run first, always keep older_than at days-not-minutes, and never run orphan removal against a location that concurrent jobs are actively writing — it is the one maintenance job that can delete a live file.

Iceberg interview question on merge-on-read cleanup

Question. A merge-on-read table takes frequent UPDATE/DELETE statements. Reads have gotten slow, and profiling shows scans spend most of their time merging thousands of small position-delete files, many of which point at data files that were already compacted away. What maintenance restores read performance, and in what order do you run it?

Solution Using delete-aware compaction plus rewrite_position_delete_files

Code.

-- 1) Compact data files that have accumulated many deletes, APPLYING the deletes
CALL spark_catalog.system.rewrite_data_files(
  table    => 'db.events',
  strategy => 'binpack',
  options  => map(
    'delete-file-threshold', '3',          -- rewrite any data file with >=3 delete files
    'target-file-size-bytes', '536870912'
  )
);

-- 2) Compact and de-dangle the remaining position delete files
CALL spark_catalog.system.rewrite_position_delete_files('db.events');

-- 3) Expire old snapshots so the pre-compaction files + applied deletes free storage
CALL spark_catalog.system.expire_snapshots(
  table => 'db.events', older_than => TIMESTAMP '2026-09-12 00:00:00', retain_last => 10
);
Enter fullscreen mode Exit fullscreen mode

Step-by-step trace.

step procedure effect on delete files
1 rewrite_data_files (delete-file-threshold=3) rewrites hot data files, applies their deletes, masked rows physically gone
2 rewrite_position_delete_files merges remaining small delete files, drops dangling deletes
3 expire_snapshots frees the superseded data + delete files from storage
  1. The slowness is read-time merge cost: every scan reconciles many small position-delete files against data files, and dangling deletes add pure waste.
  2. rewrite_data_files with delete-file-threshold => 3 selects data files carrying three or more delete files and rewrites them with the deletes applied, so those rows disappear and their delete files are no longer needed.
  3. rewrite_position_delete_files then compacts whatever position deletes remain into a few files and removes dangling deletes (those whose data file was just rewritten away), so future scans evaluate far fewer, larger delete files.
  4. expire_snapshots last reclaims the storage of the now-superseded data and delete files; order matters — compaction and delete-rewrite must commit before expiry can free their predecessors.

Output:

metric before after
position delete files scanned per read ~3,000 small ~20 large (or 0 on rewritten files)
dangling deletes many 0
read latency high (merge-bound) back to baseline

Why this works — concept by concept:

  • delete-file-threshold — ties compaction to delete pressure rather than file size, so hot MoR data files get rewritten and their deletes applied exactly when reads would otherwise pay for them.
  • Applying vs compacting deletes — rewriting a data file applies its deletes (rows physically removed); rewrite_position_delete_files compacts deletes that remain — you usually need both.
  • Dangling-delete cleanup — dropping deletes whose data file is gone removes pure read-time overhead that neither masks anything nor can ever match a row.
  • Ordering guarantee — compaction and delete-rewrite create new snapshots; expiry must run after so it can prove the old files are unreachable and safely delete them.
  • Cost — delete-aware compaction is O(rows in hot files) rewrite work paid once; it converts an every-read O(delete files) merge penalty into a one-time maintenance cost.

Optimization
Topic — optimization
Merge-on-read and read-path optimization problems

Practice →

Deduplication Topic — deduplication Delete-file, upsert and dedupe-reconciliation problems

Practice →


Cheat sheet — Iceberg maintenance recipes

Compaction — bin-pack (fix small files).

CALL spark_catalog.system.rewrite_data_files(
  table => 'db.events', strategy => 'binpack',
  options => map('target-file-size-bytes','536870912','min-input-files','5')
);
Enter fullscreen mode Exit fullscreen mode

Compaction — sort / z-order (fix read locality).

CALL spark_catalog.system.rewrite_data_files(
  table => 'db.events', strategy => 'sort',
  sort_order => 'zorder(country, event_date)',
  options => map('rewrite-all','true')
);
Enter fullscreen mode Exit fullscreen mode

Expire snapshots (reclaim storage).

CALL spark_catalog.system.expire_snapshots(
  table => 'db.events',
  older_than => TIMESTAMP '2026-09-08 00:00:00', retain_last => 5
);
Enter fullscreen mode Exit fullscreen mode

Rewrite manifests (speed up planning).

CALL spark_catalog.system.rewrite_manifests('db.events');
Enter fullscreen mode Exit fullscreen mode

Remove orphan files (dry-run first, always).

CALL spark_catalog.system.remove_orphan_files(table => 'db.events', dry_run => true);
CALL spark_catalog.system.remove_orphan_files(
  table => 'db.events', older_than => TIMESTAMP '2026-09-12 00:00:00'
);
Enter fullscreen mode Exit fullscreen mode

Compact merge-on-read position deletes.

CALL spark_catalog.system.rewrite_position_delete_files('db.events');
Enter fullscreen mode Exit fullscreen mode

Automate cleanup with table properties.

ALTER TABLE db.events SET TBLPROPERTIES (
  'write.target-file-size-bytes'                 = '536870912',
  'history.expire.max-snapshot-age-ms'           = '432000000',   -- 5 days
  'history.expire.min-snapshots-to-keep'         = '5',
  'write.metadata.delete-after-commit.enabled'   = 'true',
  'write.metadata.previous-versions-max'         = '20'
);
Enter fullscreen mode Exit fullscreen mode

Which job for which symptom.

Symptom Maintenance job
Too many tiny files, slow scans rewrite_data_files (binpack)
Reads touch too many files despite good sizes rewrite_data_files (sort / zorder)
Storage never shrinks after compaction expire_snapshots
Slow query planning, thousands of manifests rewrite_manifests
Files on disk not tracked by metadata remove_orphan_files
MoR reads slow, many small delete files rewrite_position_delete_files

Frequently asked questions

What is Iceberg table maintenance?

Iceberg table maintenance is the set of scheduled housekeeping jobs that keep an Apache Iceberg table fast, cheap, and correct as it accumulates writes. Because every commit writes new immutable files and a new snapshot, tables grow small files, dead snapshots, and bloated metadata over time. The core jobs are compaction (rewrite_data_files), snapshot expiry (expire_snapshots), manifest rewriting (rewrite_manifests), orphan-file removal (remove_orphan_files), and — for merge-on-read tables — position-delete compaction (rewrite_position_delete_files).

What is the small-files problem and how does compaction fix it?

The small-files problem is that frequent commits (streaming, micro-batches, many thin partitions) produce thousands of tiny data files, and readers pay a fixed per-file cost — opening, seeking, reading footers, scheduling tasks — so many small files are far slower to scan than a few large ones holding the same bytes. Compaction with system.rewrite_data_files and the binpack strategy packs those small files up to target-file-size-bytes (512 MB by default) as a new snapshot, so scans read fewer, larger files. It fixes read speed but not storage, because the old files remain until you expire their snapshots.

What does expire_snapshots actually delete?

expire_snapshots removes old snapshots from the table's history and then physically deletes the data files, delete files, and manifests that no remaining snapshot references. It never deletes files that a kept snapshot still shares, so a compacted table only reclaims the storage of its old small files once the snapshots referencing them expire. You control it with older_than (an age cutoff) and retain_last (a floor on recent snapshots to keep), and it is the only maintenance job that frees storage held by tracked files.

What is the difference between expire_snapshots and remove_orphan_files?

expire_snapshots works through metadata: it deletes files that were once referenced by a snapshot but no longer are. remove_orphan_files works by physical listing: it enumerates every file under the table location, diffs it against all files reachable from metadata, and deletes files that were never part of any snapshot — typically debris from failed or retried jobs. Expiry can never see orphans (they are outside the metadata graph), and orphan removal is guarded by a default 3-day older_than window because it is the one job that could otherwise delete a file an in-flight commit is still writing.

When should I use sort vs z-order compaction?

Use a single-column sort when your queries filter predominantly on one column — clustering by it tightens that column's per-file min/max stats so the planner skips files. Use zorder(colA, colB, ...) when queries filter on several columns interchangeably; a Z-order space-filling curve co-locates rows that are close across all those dimensions at once, which a linear sort cannot do beyond its leading column. Both cost a shuffle every compaction, so reach for plain binpack first and only escalate to sort/z-order when profiling shows reads are file-skipping-bound rather than file-size-bound.

How do I schedule Iceberg maintenance safely?

Run the jobs in dependency order: compaction (rewrite_data_files) and rewrite_manifests first to produce new files and snapshots, then expire_snapshots to drop the now-superseded old snapshots and free their storage, and remove_orphan_files last to sweep untracked debris. Keep snapshot retention wider than your longest-running query and your rollback objective so expiry never races a live reader, always dry_run orphan removal before a real run, and keep the older_than guard at days, not minutes. Many teams schedule frequent light compaction proportional to commit rate and run expiry/orphan removal daily.

Practice on PipeCode

Pipecode.ai is Leetcode for Data Engineering — every Iceberg maintenance idea above, from bin-pack and z-order compaction to snapshot retention, manifest clustering, and merge-on-read delete compaction, maps to a hands-on practice room where you tune real file layouts against graded inputs. PipeCode pairs each reading with 450+ DE-focused problems and a real-time scoring engine, so your answer to "how would you keep this lakehouse table fast and cheap?" holds up under a senior interviewer's depth probes.

Practice optimization problems now →
Partitioning drills →

Top comments (0)