How would you design a maintenance policy for hundreds of Iceberg tables on one platform?
answer
- one policy per class of table, not per table
- let the table carry its own settings
- the sequence of jobs matters for storage
- retention is a cost decision, not a default
- tags pin what must outlive the window
basics
~20 sClassify tables into a few tiers by write pattern, express retention and target file size as table properties so one generic job serves all of them, run rewrite then expire then occasional orphan removal, and budget the compute compaction consumes.
solid answer
~50 sDo not hand-tune hundreds of tables. Define three or four **tiers** — streaming upsert, hourly batch, daily batch, rarely written — and give each a profile: write mode, `write.target-file-size-bytes`, `history.expire.max-snapshot-age-ms` and `min-snapshots-to-keep`, and a compaction cadence. Encode the profile in table properties so a single generic job can read them and needs no per-table arguments. Fix the order of operations: `rewrite_data_files` first, then `expire_snapshots` to release what the rewrite replaced, and `remove_orphan_files` on a much longer cadence with its conservative default window. Set retention from real requirements — recovery, audit, the slowest consumer — not from a default, and use tags to pin the few snapshots that must outlive it. Budget explicitly: compaction reads and rewrites data, so it is a recurring compute line item, and prioritise the tables where scan cost actually justifies it.
code
sql · 10 lines-- streaming-upsert tier profile, applied at table creation
ALTER TABLE prod.db.events SET TBLPROPERTIES (
'write.delete.mode' = 'merge-on-read',
'write.update.mode' = 'merge-on-read',
'write.target-file-size-bytes' = '268435456',
'history.expire.max-snapshot-age-ms' = '259200000',
'history.expire.min-snapshots-to-keep' = '10',
'write.metadata.delete-after-commit.enabled' = 'true',
'write.metadata.previous-versions-max' = '50'
);go deeper
Recall that these tables need periodic upkeep — combining small files, dropping old versions, clearing leftovers — and that it is scheduled work rather than something automatic.
Explain what each maintenance procedure does and why they run in a particular order, especially that reclaiming storage requires expiry after a rewrite.
Show that you have operated this: retention chosen for real recovery and reader needs, compaction cadence matched to write rate, partial progress on big tables, and metrics that prove the policy is holding.
Own it as platform policy — a small set of tiers encoded as table defaults, an explicit compute budget for compaction, governance around the one destructive procedure, and a review path when a table does not fit its tier.
## Why per-table tuning fails at scale Every Iceberg maintenance knob is a table property or a procedure argument, so the naive approach is a bespoke job per table. At a hundred tables that becomes unowned configuration: nobody knows why one table retains 30 days of snapshots and its neighbour retains one, jobs drift out of sync with the tables they maintain, and a table onboarded last month has no maintenance at all. The design goal is that a **new table is maintained correctly by default**. ## Tier the tables, not the jobs Classify by write pattern, because that is what determines every setting: - **Streaming upsert / CDC.** Merge-on-read for deletes and updates, many small files and delete files per hour. Needs frequent compaction, short snapshot retention (commits are numerous), and a smaller target file size if latency matters. - **Hourly or micro-batch append.** Moderate small-file production. Daily compaction usually suffices. - **Daily batch rebuild.** Few large commits. Copy-on-write, little or no compaction, retention chosen for restatement windows. - **Rarely written reference data.** Almost no maintenance; occasional manifest tidy-up at most. Each tier is a named profile with concrete values for `write.delete.mode` / `write.update.mode` / `write.merge.mode`, `write.target-file-size-bytes`, `write.distribution-mode`, `history.expire.max-snapshot-age-ms`, `history.expire.min-snapshots-to-keep`, and `write.metadata.delete-after-commit.enabled` with `write.metadata.previous-versions-max`. ## Make the table carry its own policy Set the profile as table properties at creation. Then one generic maintenance job iterates the catalog and calls `expire_snapshots` with no arguments — the defaults come from the table's own `history.expire.*` values — and `rewrite_data_files` with tier-appropriate options. Configuration lives next to the table, is visible in `SHOW TBLPROPERTIES`, and survives whoever wrote the job. ## Fix the order of operations 1. **`rewrite_data_files`** — fix layout, and on merge-on-read tables materialise delete files away. Enable `partial-progress.enabled` on large tables so a late failure does not discard hours of work, and use `where` to restrict to recently written partitions rather than rewriting history every night. 2. **`expire_snapshots`** — until the snapshots referencing the pre-compaction files are gone, storage holds both copies. Skipping this is why teams report that compaction increased their bill. 3. **`rewrite_manifests`** — occasionally, on tables with many commits, to keep manifest counts and their partition clustering healthy so planning stays fast. 4. **`remove_orphan_files`** — weekly or monthly, keeping the conservative default age window, always dry-run first, never aimed at a shared location. ## Set retention from requirements Snapshot retention buys three things: rollback after a bad write, time travel for reproducibility or audit, and safety for long-running readers. Ask what each table actually needs. A financial table may need month-end states for years — pin those with tags, which survive expiry, rather than raising the general window for every snapshot. A high-frequency streaming table producing a commit a minute cannot retain 30 days without an enormous snapshot log and file count. Retention is a cost decision, and it should be made once per tier with a documented rationale. ## Budget the compute Compaction reads and rewrites real data; a table compacted daily is read and written daily on top of its query load. Justify it where scan cost is high — the tables driving dashboards and frequent joins — and let cold tables sit fragmented. Track outcomes rather than job success: average file size and file count per partition from the `files` metadata table, delete-file counts on merge-on-read tables, snapshot counts, and total table size before and after expiry. Those numbers say whether the policy is working; a green job status does not. ## Governance and blast radius Maintenance runs with write credentials on every table, and `remove_orphan_files` deletes objects the table has never referenced. Give it a dedicated identity, restrict who can change the safety-relevant arguments, and log what each run deleted. Decide in advance what happens when compaction conflicts with a heavy writer — usually schedule around the write window rather than retry into it. Finally, treat the tier profiles as reviewable defaults: when someone needs an exception, the interesting question is whether the tier is wrong for a whole class of tables rather than just for this one.
- How do you decide which tables get frequent compaction and which get none?By scan cost, not by fragmentation alone. Measure file count and average file size per partition from the `files` metadata table, weight it by how often and how widely the table is queried, and compact where the rewrite compute is cheaper than the repeated scan overhead it removes. A fragmented table nobody reads is not worth the rewrite; a moderately fragmented table behind every dashboard is.
- What single metric would tell you the policy is failing?Delete-file or small-file count per partition trending upward week over week on the tables in the frequent-compaction tier — it means maintenance is not keeping pace with ingest, and read latency will follow. Table size that never drops after expiry is the second signal, usually meaning expiry is not running or retention is too long for the commit rate.
- Why not just run every maintenance procedure nightly on every table?Cost and risk. Compaction rewrites data whether or not the table benefits, so it multiplies compute across cold tables for nothing, and a full-listing orphan sweep on millions of objects is slow and is the one procedure that can delete files a running writer needs. Cadence should follow write rate and read value per tier, with orphan removal deliberately rare.
saying these in an interview costs you the question
- Hand-tunes every table instead of defining tiers
- Expires snapshots before compacting and wonders where storage went
- Copies default retention without asking what recovery needs
- Runs orphan removal nightly on every table
- Measures job success instead of file-size and delete-file trends