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.
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
- Why Iceberg tables need maintenance
- Compaction with rewrite_data_files
- Expiring snapshots to reclaim storage
- Rewriting manifests & metadata
- Orphan files & merge-on-read delete compaction
- Cheat sheet — Iceberg maintenance recipes
- Frequently asked questions
- Practice on PipeCode
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
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 totarget-file-size-byteswithout 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 explicitsort_order). Clustering rows by a frequently-filtered column tightens per-file min/max stats so the reader can skip more files. Costs a shuffle. -
sortwithzorder(...). A space-filling-curve sort that co-locates rows across multiple dimensions at once. Use when queries filter on several columns interchangeably (e.g.countryandevent_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'swrite.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.
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'
)
);
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'
)
);
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 |
- Plain
binpackcannot help here — the files are already the right size; the problem is the arrangement of rows inside them, so you need asort-family strategy. - A single linear
sort_order => 'country, event_date'would clustercountrywell but scatterevent_datewithin each country. Because queries filter on both interchangeably, a z-order curve co-locates rows that are close in both dimensions. - After the rewrite, each file's min/max footer stats for
countryandevent_dateare narrow, so Iceberg's manifest-level filtering prunes files that cannot match the predicate before any data is read. -
partial-progress.enabledlets 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
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_thansets an age cutoff;retain_lastsets 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 lastretain_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.snapshotsand.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 (exceptmain).
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_filesis for.
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
);
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
);
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 |
-
expire_snapshotsdeletes files that are unreachable from surviving snapshots — but a reader that already resolvedS_oldstill expectsS_old's files to exist on disk. - Setting
older_than => now()makesS_oldeligible the instant after compaction, so its exclusive files (the small pre-compaction files) get deleted out from under the in-flight reader. - 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.
-
retain_last => 10adds 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
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 scheduledrewrite_manifestsinstead. -
write.metadata.delete-after-commit.enabledandwrite.metadata.previous-versions-max— automatically prune oldmetadata.jsonfiles so the metadata directory does not grow without bound.
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;
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');
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 |
- 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.
-
write.metadata.delete-after-commit.enabledwithprevious-versions-max = 20makes each commit prune oldmetadata.jsonfiles, so the metadata directory stops growing unbounded. -
rewrite_manifestsfixes the existing manifest bloat once;commit.manifest-merge.enabled = truekeeps future commits from re-accumulating tiny manifests by merging them at write time. - 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-maxbound themetadata.jsonchain so a high-frequency writer does not leave a million root-pointer files behind. -
Inline merge vs scheduled rewrite —
manifest-merge.enabledprevents regrowth at commit time; scheduledrewrite_manifestscleans 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_manifestsis 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
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 adelete-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.
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'
);
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
);
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 |
- The slowness is read-time merge cost: every scan reconciles many small position-delete files against data files, and dangling deletes add pure waste.
-
rewrite_data_fileswithdelete-file-threshold => 3selects 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. -
rewrite_position_delete_filesthen 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. -
expire_snapshotslast 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_filescompacts 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
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')
);
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')
);
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
);
Rewrite manifests (speed up planning).
CALL spark_catalog.system.rewrite_manifests('db.events');
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'
);
Compact merge-on-read position deletes.
CALL spark_catalog.system.rewrite_position_delete_files('db.events');
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'
);
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.





Top comments (0)