Iceberg storage keeps growing: expire snapshots and orphan files
IcebergThe object store bill for an Iceberg table fed from Kafka keeps climbing. The table’s metadata directory holds thousands of JSON and Avro files. The snapshots metadata table returns one row for every commit since the table was created, even when the data the table currently holds has stayed about the same size.
Iceberg storage growth is a table that keeps old snapshots, and the data and metadata files those snapshots reference, long after any query needs them. A streaming writer commits every few minutes, each commit adds a snapshot and a new metadata file, and nothing is removed until a maintenance job removes it.
The Iceberg maintenance documentation lists expiring snapshots, removing old metadata files and deleting orphan files as recommended maintenance, and says tables with frequent commits, like those written by streaming jobs, may need to clean metadata files regularly. For how tables, snapshots and metadata fit together, see the complete Iceberg guide.
This page is about storage and metadata that pile up over time. If the table holds many tiny data files and queries are slow, start with too many small files and how to compact them, because compaction adds its own snapshots that this page then expires. If the Flink job writing the table restarts because its checkpoints fail, start with Flink checkpoints that fail or time out.
Why the table keeps growing
A write to an Iceberg table adds new files and a new version of the table, and the earlier version stays readable until it is expired.
Every commit adds a snapshot and a metadata file
The maintenance documentation says each write to an Iceberg table creates a new snapshot, or version, of the table, and that snapshots accumulate until they are expired. It also says each change to a table produces a new metadata file to provide atomicity, and that old metadata files are kept for history by default. Each snapshot points to a manifest list, and the manifest list points to manifests that track the data files, so a commit adds files at several levels of the tree.
The writer sets how often this happens. The Flink writes documentation says the Iceberg commit happens after a successful Flink checkpoint, so the checkpoint interval is the commit interval. The Flink key below sets that interval, and a value of 5min commits about every five minutes.
execution.checkpointing.interval: 5min
The Apache Iceberg Kafka Connect sink is the open source, self-managed way to write Kafka topics into Iceberg, and it commits on a timer that defaults to 300,000 ms, five minutes.
iceberg.control.commit.interval-ms=300000
At one commit every five minutes, a table gains 288 snapshots and 288 metadata files a day before any compaction or delete job adds its own commits.
Removed data stays referenced
Deleting rows, overwriting a partition or compacting small files does not free storage at commit time. The Spark procedures documentation says each write, update, delete, upsert or compaction produces a new snapshot while keeping the old data and metadata around for snapshot isolation and time travel. The files a compaction replaced stay on storage as long as an older snapshot still lists them, so a table that is compacted but never expired grows faster, not slower.
Manifests are merged only after they pile up
On write, Iceberg merges manifests automatically once enough of them have accumulated. The table configuration documentation sets the minimum number of manifests to accumulate before merging at 100 and the target size of a merged manifest at 8 MB.
commit.manifest-merge.enabled=true
commit.manifest.min-count-to-merge=100
commit.manifest.target-size-bytes=8388608
Measure what the table is holding
Pick the metric first, then change one thing at a time and measure again. Iceberg already exposes its own metrics, metadata and statistics as metadata tables, documented in the Spark queries documentation, and they can be queried like any table, so growth can be measured before anything is changed.
Rule out real data growth
Storage that climbs because the data itself grew is a capacity question, not a retention problem, so rule it out first. The files table lists only the current snapshot’s files, so its total is the size of the live table.
SELECT sum(file_size_in_bytes) AS live_bytes
FROM prod.db.events.files
The summary map of each snapshot can hold the same total as total-files-size, which the table specification lists as an optional summary field, so the history of live bytes comes from the snapshots table.
SELECT date_trunc('day', committed_at) AS day,
max(CAST(summary['total-files-size'] AS BIGINT)) AS live_bytes
FROM prod.db.events.snapshots
GROUP BY 1
ORDER BY 1
If live bytes climb as fast as storage does, the table holds more data and expiry will not shrink it. If live bytes stay flat while storage climbs, old snapshots are the cause, so continue with the steps below.
Count snapshots and commits per day
The snapshots table shows the valid snapshots for a table, with committed_at, operation, manifest_list and a summary map that holds counts such as added-data-files and total-data-files.
SELECT date_trunc('day', committed_at) AS day,
count(*) AS snapshots,
min(committed_at) AS oldest
FROM prod.db.events.snapshots
GROUP BY 1
ORDER BY 1
The oldest committed_at shows how much history the table keeps. If it reaches back to the day the pipeline started, snapshots have never been expired.
Compare live files with all referenced files
The data_files table shows the current snapshot’s data files, while all_data_files returns data files across all snapshots. The gap between the two byte totals is the data file storage that only older snapshots keep alive.
Data files in the current snapshot:
SELECT count(*) AS files, sum(file_size_in_bytes) AS bytes
FROM prod.db.events.data_files
Data files across all snapshots, each counted once:
SELECT count(*) AS files, sum(file_size_in_bytes) AS bytes
FROM (
SELECT DISTINCT file_path, file_size_in_bytes
FROM prod.db.events.all_data_files
)
The documentation warns that the “all” metadata tables may produce more than one row per data file or manifest file, because a file can belong to more than one snapshot, so the second query selects distinct paths rather than counting rows. Run both again after expiry to see the gap close.
Count manifests and metadata files
The manifests table lists the current snapshot’s manifests with their length and added_data_files_count, all_manifests lists manifests across snapshots, and metadata_log_entries lists the metadata JSON files the table still tracks.
Manifests in the current snapshot:
SELECT count(*) AS manifests, sum(length) AS bytes
FROM prod.db.events.manifests
Metadata files the table still tracks:
SELECT count(*) AS tracked_metadata_files
FROM prod.db.events.metadata_log_entries
Metadata files and snapshots are counted separately. The table specification adds an entry to the metadata log each time a new metadata file is created, and an entry to the snapshot log each time the current snapshot changes, so the newest committed_at in snapshots shows the last write while metadata_log_entries also grows with other changes to the table.
A large manifest count for the current snapshot points at manifest rewriting. A metadata directory holding far more JSON files than metadata_log_entries tracks points at orphaned metadata files, covered below.
Check what must be kept
The refs table lists branches and tags with the snapshot each one points to and its retention settings. Snapshots that a branch or tag still references are not removed by expiry, so read this table before choosing a cutoff.
SELECT name, type, snapshot_id, max_reference_age_in_ms
FROM prod.db.events.refs
Which cleanup removes what
The fix has four parts: expire snapshots with a retention rule that keeps the history readers and the writer still need, let the writer delete old metadata files, clear orphan files with a safe cutoff, and rewrite manifests only if planning is slow. Expiry and orphan removal are different procedures with different risks.
Expire snapshots, then clean metadata and orphan files
Run the steps in this order and measure after each one. Expiry comes first because it removes the snapshots and the files only they referenced.
Expire old snapshots
The expire_snapshots procedure removes old snapshots and the data files uniquely required by them, and the documentation states it will never remove files still required by a non-expired snapshot. older_than defaults to 5 days ago and retain_last, the number of ancestor snapshots kept regardless of age, defaults to 1.
-- keep the same window as history.expire.max-snapshot-age-ms
CALL prod.system.expire_snapshots(
table => 'db.events',
older_than => current_timestamp() - INTERVAL 7 DAYS,
retain_last => 100,
stream_results => true
)
The procedure returns deleted_data_files_count, deleted_manifest_files_count and deleted_manifest_lists_count, which is the metric that proves the run removed something. The documentation recommends stream_results => true to keep the Spark driver from running out of memory when many files are deleted. The maintenance documentation also describes a Spark action, expireSnapshots, that runs expiry in parallel for large tables.
Set the retention once, on the table
When older_than and retain_last are omitted, the procedure uses the table’s expiration properties. Setting them on the table means every scheduled run, from any engine, applies the same rule. The configuration documentation gives the defaults as 432000000 ms, five days, and one snapshot.
ALTER TABLE prod.db.events SET TBLPROPERTIES (
'history.expire.max-snapshot-age-ms' = '604800000',
'history.expire.min-snapshots-to-keep' = '100'
)
The value above keeps seven days. Choose the age from the longest time travel query or rollback the team needs, not from storage cost alone. Iceberg lets engines such as Spark, Trino and Flink work with the same table at the same time, as the Apache Iceberg home page describes, so this one rule sets how far back every engine reading the table can time travel. Snowflake documents Iceberg tables whose files sit in external cloud storage the owner manages, and says that owner is responsible for data protection and recovery. The gc.enabled property, true by default, must stay true, because it allows garbage collection operations such as expiring snapshots and removing orphan files.
Let the writer delete old metadata files
Expiry removes snapshots, not old metadata JSON files. The maintenance documentation describes two table properties for those. With write.metadata.delete-after-commit.enabled set to true, the writer deletes the oldest tracked metadata file each time it creates a new one, keeping up to write.metadata.previous-versions-max, which defaults to 100.
write.metadata.delete-after-commit.enabled=true
write.metadata.previous-versions-max=100
The setting only applies to files still tracked in the metadata log. The documentation gives the example of a table with the setting off and previous-versions-max at 10, which after 100 commits has 10 tracked and 90 orphaned metadata files, and those 90 can only be cleaned by orphan file deletion.
Remove orphan files with a safe cutoff
Failed tasks and jobs can leave files that no metadata references. The remove_orphan_files procedure removes them and defaults to files created more than 3 days ago. Start with a dry run, which lists the candidates without deleting them.
CALL prod.system.remove_orphan_files(
table => 'db.events',
dry_run => true
)
The maintenance documentation calls it dangerous to remove orphan files with a retention interval shorter than the time expected for any write to complete, because in-progress files can be treated as orphans and deleted. It also warns that Iceberg compares paths as strings, so a path that changed over time, for example after an HDFS authority change, can lead to data loss. Check that the paths in the metadata tables match the storage listing before the real run. The documentation adds that the action can take a long time on large directories and may not need to run often.
Rewrite manifests when planning is slow
The rewrite_manifests procedure rewrites manifests to optimize scan planning, sorted by the fields in the partition spec, and returns rewritten_manifests_count and added_manifests_count. The maintenance documentation says this helps when the table’s write pattern does not line up with the query pattern.
CALL prod.system.rewrite_manifests('db.events')
A manifest rewrite is a commit, so it adds a snapshot of its own that the next expiry run removes in turn.
Keep time travel and running readers working
Expiry is safe for history as long as the cutoff is chosen deliberately. The maintenance documentation says expired snapshots are no longer available for time travel queries, and that data files are not deleted until no snapshot that may be used for time travel or rollback references them.
Pin the history that matters with tags
Spark time travel uses TIMESTAMP AS OF or VERSION AS OF, and VERSION AS OF accepts a snapshot ID or a branch or tag name. The branching and tagging documentation describes tags as the way to retain important historical snapshots for auditing, with their own retention, for example one snapshot per week kept for a month.
ALTER TABLE prod.db.events CREATE TAG `EOW-01` AS OF VERSION 7 RETAIN 30 DAYS
The 7 stands for a snapshot_id taken from the snapshots table.
A tagged snapshot survives expiry for as long as the tag’s retention, while the untagged commits around it expire on the table’s normal rule.
Leave room for long queries
A query that started against an older snapshot reads the files that snapshot lists. If expiry removes that snapshot and its unique files while the query is still running, the query can fail to find them. Set the age cutoff well beyond the longest query, batch job or downstream incremental read, and run expiry outside the windows when long reports run.
Keep the Flink writer’s last snapshot
The Flink writes documentation says Flink streaming write jobs keep the last committed checkpoint ID in the snapshot summary and store uncommitted data as temporary files, so expiring snapshots and deleting orphan files could corrupt the Flink job’s state. The complete Flink guide covers the job itself. The documentation’s guidance is to keep the last snapshot created by the Flink job, identified by the flink.job-id property in the summary, and to delete only orphan files that are old enough.
SELECT snapshot_id, committed_at, summary['flink.job-id'] AS job
FROM prod.db.events.snapshots
WHERE summary['flink.job-id'] IS NOT NULL
ORDER BY committed_at DESC
LIMIT 1
Because retain_last counts the newest snapshots of any kind, a compaction or manifest rewrite committed after the writer’s last snapshot pushes that snapshot further down the list. This is arithmetic on retain_last, not wording from the Iceberg documentation. Count the snapshots committed after the writer’s last one:
SELECT count(*) AS newer_snapshots
FROM prod.db.events.snapshots
WHERE committed_at > (
SELECT max(committed_at)
FROM prod.db.events.snapshots
WHERE summary['flink.job-id'] IS NOT NULL
)
Set retain_last above that count plus a margin, then run the first query again after expiry to confirm the snapshot is still listed. A tag on that snapshot is the other option, because the Spark procedures documentation says snapshots still referenced by a branch or tag are not removed.
Run maintenance from Flink instead of Spark
The Flink TableMaintenance documentation describes an API that runs snapshot expiry, small file compaction and orphan file cleanup from Flink, either inside the streaming pipeline or as a separate Flink job, without a Spark cluster. Its ExpireSnapshots task defaults to a maximum snapshot age of 5 days, keeps at least 1 snapshot and deletes in batches of 1000 files. Its DeleteOrphanFiles task defaults to a minimum age of 3 days. The documentation describes lock management so that only one maintenance operation runs per table at a time, and marks lock configuration as no longer required for a single Flink job.
The Flink writes documentation also describes post-commit maintenance on the IcebergSink, where expireSnapshots() and deleteOrphanFiles() run after each commit. The same page labels the SinkV2 based IcebergSink as an experimental feature to use with caution.
Watch the pipeline while maintenance runs
Maintenance runs beside the writer, not instead of it, and the pipeline should keep committing throughout. Two signals on the Kafka side show whether ingestion into the table keeps up: connector task errors, for example a task failing after a backward-incompatible schema change, and consumer lag. The Kafka Connect user guide describes a REST endpoint that returns the status of a connector, including the state of every task and error information.
GET /connectors/{name}/status
Lag becomes data loss when it reaches the topic’s retention. The Kafka broker configuration defines log.retention.hours, 168 by default, as the number of hours to keep a log file before deleting it. Retention is time or size based and does not check consumer offsets, so lag that reaches retention means unread records are deleted. If a sink is stopped during a maintenance window or after a failed commit, alert on lag approaching retention rather than on a fixed record count. The guide to monitoring Kafka consumer lag covers how to measure it, and the guide to diagnosing a failed Kafka Connect connector covers task failures.
On the Flink side, a table that keeps growing while the job’s own state also grows is two problems. The guide to Flink state that keeps growing covers the job side.
Where Kpow and Flex fit
This page is published by Factor House, which makes Flex, its Flink job management and monitoring product, and Kpow, its Kafka management and monitoring product, so weigh this section with that in mind. Neither product’s documentation describes Iceberg snapshots, metadata tables, snapshot expiry or orphan file removal, so those stay in the Iceberg metadata tables and the Spark or Flink maintenance jobs above.
Kafka side. Kpow’s Kafka Connect management documentation describes pausing, restarting and deleting connectors, editing their configuration, and restarting individual tasks or viewing their stack traces when a task is in an ERROR state. Its Prometheus metrics glossary lists group_offset_lag, the total lag of all assignments of a consumer group.
Flink side. Flex’s job inspection documentation describes a Checkpoints tab with counts by status, average checkpoint size and duration, and a history table of every checkpoint attempt, plus an Events tab that logs the job’s lifecycle. Because the Iceberg sink commits after each successful checkpoint, the checkpoint history shows the commit cadence that sets how fast snapshots accumulate. It does not show whether an Iceberg commit succeeded or how many snapshots the table holds.
Flex is one of several ways to watch a Flink job, and the comparison of tools for self-managed Apache Flink sets it beside the Flink web UI, the REST API and metrics stacks.
FAQ
Does expire_snapshots break time travel?
It removes time travel to the snapshots it expires, and only those. The maintenance documentation says expired snapshots are no longer available for time travel queries, while snapshots inside the retention window and snapshots referenced by a branch or tag stay readable.
What is the difference between expire_snapshots and remove_orphan_files?
expire_snapshots removes snapshots from the table metadata and deletes the files only those snapshots needed. remove_orphan_files lists the table location and deletes files that no metadata references at all, such as files left by failed writes and untracked metadata files.
Why does the metadata directory keep growing after expiry?
Expiry does not delete old metadata JSON files. Set write.metadata.delete-after-commit.enabled=true so the writer deletes the oldest tracked file on each commit, and run orphan file removal once to clear files that were already untracked.
How long should snapshots be kept on a streaming table?
The Iceberg defaults are 5 days and at least one snapshot. The right value is the longest time travel, rollback or incremental read the team needs, plus enough recent snapshots to keep the Flink writer’s last one, with tags for any snapshot that must be kept longer.