Apache Iceberg architecture: what actually happens when you commit
A walk through the Iceberg metadata tree from a catalog pointer down to a single data file: what a commit actually swaps, how planning prunes before it reads, and why the catalog is part of the table.
- Architecture: the metadata tree
- Why is the catalog part of the table?
- What does a commit actually do?
- How does a query find its files?
- How does a column survive a rename?
- How is a row deleted without rewriting the file?
- What does a reader do with a delete file?
- Reading it from a running table
- How do branches and tags work?
- Why does the metadata keep growing?
- What happens when two engines commit at once?
- When is a simpler choice better?
- Common misconceptions
- Production tips
- Frequently asked questions
- Conclusion
- References
- Trademarks
TL;DR
- A table is an immutable tree: a catalog pointer to
metadata.json, which names a snapshot, which points at a manifest list, which points at manifests, which list data files. Nothing already written is ever mutated.- The catalog is part of the table, not a lookup service. Nothing on storage says which
metadata.jsonis current, so commit atomicity is a property of the catalog.TableOperations.commitmandates a compare-and-swap and a distinctCommitStateUnknownException, which is the difference between a commit that failed and one whose outcome is unknown. That distinction is why orphan files exist.- Planning is a sequence of prunes over metadata: manifest list by partition range, manifests by partition value, then per-file column bounds. No directory is ever listed.
- Hidden partitioning and partition evolution both work because partition values live in manifests rather than in directory names.
The first Iceberg surprise is operational rather than conceptual. The
table works, queries are fast, and then two writers collide and one fails with
something about stale metadata, or a VACUUM-shaped cleanup deletes files a
reader still wanted. Neither makes sense if you think of Iceberg as a file
layout, because neither is about files.
This is the map. It follows one commit from the statement to the catalog swap, and one query from a filter to the files it actually opens. By the end you should be able to look at a table directory and say what is current, and reason about what two concurrent writers will do to each other.
Assumed knowledge: comfort with SQL and a query engine. No Iceberg internals knowledge is assumed.
Architecture: the metadata tree
Ask Iceberg “which files make up this table right now” and the answer walks a tree, every level of which is an immutable file.
flowchart LR
CAT[("catalog<br/>one pointer")] --> MD["v2.metadata.json<br/>schemas, specs, snapshots"]
MD --> SNAP["current snapshot<br/>7241925443479918015"]
SNAP --> ML["manifest list<br/>snap-...avro"]
ML --> M1["manifest m0<br/>partition ranges, counts"]
ML --> M2["manifest m1"]
M1 --> D1["data files"]
M1 --> DEL["delete files"]
M2 --> D2["data files"]
| Level | File | What it holds |
|---|---|---|
| Catalog | none, it is a service | The single pointer to the current metadata.json |
| Table metadata | metadata/v2.metadata.json |
Schemas, partition specs, sort orders, properties, the snapshot log, the current snapshot id |
| Snapshot | an entry inside metadata.json |
The complete state of the table at one commit |
| Manifest list | metadata/snap-*.avro |
One row per manifest, with partition ranges and added, existing and deleted counts |
| Manifest | metadata/*-m0.avro |
One row per data or delete file, with partition values and per-column bounds |
| Data file | data/city_id=sf/00000-*.parquet |
The rows |
On disk:
s3a://lakehouse-prod/warehouse/trips/
├── metadata/
│ ├── v1.metadata.json
│ ├── v2.metadata.json <- current, per the catalog
│ ├── snap-7241925443479918015-1-a1c2....avro <- manifest list
│ ├── a1c2f3b4-....-m0.avro <- manifest
│ └── a1c2f3b4-....-m1.avro
└── data/
├── city_id=sf/started_at_day=2026-09-15/00000-0-a1c2f3b4-....parquet
└── city_id=nyc/started_at_day=2026-09-15/00000-1-a1c2f3b4-....parquet
Two things are worth pausing on. The data/ directory is laid out by partition
for human convenience, but the engine never relies on it: partition values
come from the manifest. And nothing in that listing says which metadata.json
is current. Open the directory without a catalog and you cannot answer the most
basic question about the table.
Why is the catalog part of the table?
Because the commit is a pointer swap, and something has to make that swap atomic.
TableOperations.commit is an interface with unusually specific
documentation, and it is worth reading because it defines the guarantee:
Implementations must check that the base metadata is current to avoid overwriting updates. Once the atomic commit operation succeeds, implementations must not perform any operations that may fail because failure in this method cannot be distinguished from commit failure.
That first sentence is optimistic concurrency stated as a contract: a writer
prepares a new metadata.json from the base it read, and the commit succeeds
only if the base is still current. If another writer moved the pointer first,
this one fails and retries from the new base.
The second half is subtler and explains a real operational problem:
Implementations must throw a
CommitStateUnknownExceptionin cases where it cannot be determined if the commit succeeded or failed. For example if a network partition causes the confirmation of the commit to be lost. This is important because downstream users of this API need to know whether they can clean up the commit or not, if the state is unknown then it is not safe to remove any files.
A failed commit and an unknown commit are different states. After a clean
failure, the files the writer wrote are garbage and can be deleted. After an
unknown outcome, they may be referenced by a snapshot that did commit, so
deleting them would corrupt the table. Iceberg therefore leaves them, which is
precisely why remove_orphan_files exists and why it takes a conservative
older_than.
| Catalog | Provides the atomic swap via |
|---|---|
| REST | A service implementing the Iceberg REST spec |
| Hive Metastore | A compare-and-swap on the table’s metadata_location property |
| JDBC | A row update in a relational database |
| Glue | A conditional update in the AWS Glue Data Catalog |
| Nessie | A git-like commit against a branch |
| Hadoop | A filesystem rename, which HadoopTableOperations documents as requiring a filesystem that supports atomic rename |
That last row is the caveat to carry. Two Iceberg tables can look byte-identical on storage and have different correctness guarantees under concurrent writes, because the guarantee lives in the catalog rather than in the files.
What does a commit actually do?
flowchart TB
W["writer produces<br/>new data files"] --> M["write a new manifest<br/>listing them"]
M --> ML["write a new manifest list,<br/>reusing unchanged manifests"]
ML --> MJ["write a new metadata.json<br/>with the new snapshot"]
MJ --> CAS{"catalog compare-and-swap:<br/>is the base still current?"}
CAS -->|"yes"| OK["commit succeeds<br/>the new snapshot is live"]
CAS -->|"no"| RETRY["CommitFailedException<br/>re-plan from the new base"]
CAS -.->|"outcome unknown"| UNK["CommitStateUnknownException<br/>leave the files alone"]
RETRY -.-> M
The step that makes snapshot history affordable is manifest reuse. A commit that adds one file writes a new manifest and a new manifest list, but the manifest list points at every unchanged manifest from the previous snapshot rather than copying it. Successive snapshots share most of their metadata, so keeping a thousand snapshots costs metadata proportional to what changed, not to the size of the table.
Retention is therefore not free but is cheap, and it is bounded by
history.expire.max-snapshot-age-ms, which defaults to 5 days.
How does a query find its files?
Planning is a sequence of prunes, and each one reads metadata rather than data.
flowchart LR
Q["query with a filter<br/><i>city_id = 'sf' AND fare > 50</i>"] --> A["manifest list<br/>skip whole manifests<br/>by partition range"]
A --> B["manifests<br/>skip files by<br/>partition value"]
B --> C["column bounds<br/>skip files whose<br/>min/max cannot match"]
C --> D["the files actually opened"]
The manifest list carries partition_spec_id, added_files_count,
existing_files_count, deleted_files_count, row counts and per-partition
lower_bound and upper_bound. That is enough to discard an entire manifest,
and therefore hundreds of files, without opening it.
Notice what is absent: no step lists a directory. This is the structural
difference from Hive-style tables, where planning is a LIST per partition and
cost grows with partition count. Iceberg’s planning cost grows with the number
of manifests, which is what rewrite_manifests exists to control.
Column bounds only prune if the data is sorted enough for them to be narrow.
Unsorted data gives every file a wide min and max, the bounds overlap, and
nothing is skipped. That is why WRITE ORDERED BY and a sorting compaction are
not cosmetic: they are what makes the third prune work at all.
What a commit costs, file by file
The manifest-reuse argument above is testable. Creating a partitioned table and appending three times, counting files in the metadata directory after each step:
after CREATE, before any data metadata_files=2 data_files=0
after append #1 metadata_files=5 data_files=1
after append #2 metadata_files=8 data_files=2
after append #3 metadata_files=11 data_files=3
Three metadata files per commit, every commit: one manifest listing the new
data file, one manifest list for the new snapshot, and one metadata.json
carrying the new table state. The two files present before any data are the
initial metadata.json and the catalog’s version hint.
The snapshot log shows the parent chain those files encode:
snap=3527536884859544448 parent=None op=append added-files=1 total-records=3000
snap=2869732689965030704 parent=3527536884859544448 op=append added-files=1 total-records=6000
snap=4093333610368405196 parent=2869732689965030704 op=append added-files=1 total-records=9000
added-data-files stays at 1 while total-records accumulates, which is the
manifest reuse made visible: each commit describes only its own addition and
inherits the rest.
A caution when comparing this with other write-ups. The metadata filenames
depend on the catalog. A Hadoop catalog produces v1.metadata.json,
v2.metadata.json and a version-hint.text; a REST catalog produces
00000-<uuid>.metadata.json with no version hint, because the pointer lives in
the service rather than in a file. Both were produced in the same session against
the same Iceberg version. Tutorials that describe the metadata directory
differently are usually describing different catalogs.
The pruning evidence
The manifest-list prune is visible in the manifests metadata table, which
carries the partition range of each manifest:
manifest af2b689e-...-m0.avro files=1
partition_summaries=[lower_bound='2026-01-07', upper_bound='2026-01-07']
manifest aef10e73-...-m0.avro files=1
partition_summaries=[lower_bound='2026-01-06', upper_bound='2026-01-06']
manifest 53fea6b1-...-m0.avro files=1
partition_summaries=[lower_bound='2026-01-05', upper_bound='2026-01-05']
A filter on a single day discards two of those three manifests before opening
any of them. The per-file column bounds that drive the third prune are in the
files table, keyed by field id:
lower_bounds={1: ..., 2: ..., 3: 'r0', 4: ...}
upper_bounds={1: ..., 2: ..., 3: 'r2', 4: ...}
Note the keys are field ids, not column names — the same identifiers that make a rename safe.
Why is partitioning hidden?
Because partition values live in manifests rather than in path names, which makes two things possible that are hard to retrofit elsewhere.
Every SQL block below runs in a session with a catalog named prod. On Spark
that is four settings, and prod is the name you chose, not a reserved word:
spark-sql \
--packages org.apache.iceberg:iceberg-spark-runtime-4.1_2.13:1.11.0 \
--conf spark.sql.extensions=org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions \
--conf spark.sql.catalog.prod=org.apache.iceberg.spark.SparkCatalog \
--conf spark.sql.catalog.prod.type=hive \
--conf spark.sql.catalog.prod.uri=thrift://metastore.internal:9083 \
--conf spark.sql.catalog.prod.warehouse=s3a://lakehouse-prod/warehouse
type=hive is the line that picks the commit guarantee discussed above. Swap it
for type=rest plus a uri to use a REST catalog instead.
Hidden partitioning. You declare a transform on a column, and queries filtering that column prune automatically:
CREATE TABLE prod.db.trips (
trip_id STRING,
rider_id STRING,
fare_amount DECIMAL(10,2),
started_at TIMESTAMP,
city_id STRING
) USING iceberg
PARTITIONED BY (city_id, days(started_at));
-- prunes to a day without naming a partition column
SELECT count(*) FROM prod.db.trips
WHERE started_at >= TIMESTAMP '2026-09-15 00:00:00'
AND started_at < TIMESTAMP '2026-09-16 00:00:00';
There is no dt column in the schema, so there is no way for a query to forget
to filter on it, and no correlation for a user to remember.
Partition evolution. Because each manifest records the partition_spec_id
it was written under, the spec can change and old files keep working:
ALTER TABLE prod.db.trips ADD PARTITION FIELD bucket(16, rider_id);
ALTER TABLE prod.db.trips DROP PARTITION FIELD days(started_at);
Measured. A table partitioned by months(ts), then altered with
REPLACE PARTITION FIELD months(ts) WITH days(ts) and appended to again, ends up
with two specs recorded and its files split between them:
partition-specs: [{"spec-id": 0, "fields": [{"name": "ts_month", "transform": "month", "source-id": 2, "field-id": 1000}]},
{"spec-id": 1, "fields": [{"name": "ts_day", "transform": "day", "source-id": 2, "field-id": 1001}]}]
default-spec-id: 1
spec_id per data file: [1, 0, 0]
Two of the three data files still carry spec_id = 0, and nothing rewrote them.
The directory listing shows both layouts side by side — ts_month=2026-01,
ts_month=2026-02 and ts_day=2026-03-09 — and all 900 rows remain queryable
through one table. Both specs reference source-id: 2, the same underlying
column, which is what lets a single predicate on ts prune across both.
Files written before the change keep their original spec; files after use the new one; a scan reads both. In a directory-partitioned table that change is a full rewrite.
Schema evolution works on the same principle one level down: columns carry assigned IDs, so a rename is a metadata edit and data files are never touched.
How does a column survive a rename?
Because columns are identified by number, not by name. Every field in an Iceberg schema carries an ID, the Parquet files record those IDs, and the name is a label stored only in metadata. Renaming a column edits the label and touches no data.
You can watch the IDs rather than trust the claim. metadata.json holds the
schema, so read it before and after:
CREATE TABLE prod.db.t (trip_id STRING, fare DECIMAL(10,2), city STRING) USING iceberg;
-- field ids: trip_id 1, fare 2, city 3
ALTER TABLE prod.db.t RENAME COLUMN fare TO fare_amount;
-- field ids: trip_id 1, fare_amount 2, city 3 <- id 2 unchanged
ALTER TABLE prod.db.t DROP COLUMN city;
ALTER TABLE prod.db.t ADD COLUMN rider STRING;
-- field ids: trip_id 1, fare_amount 2, rider 4 <- id 3 retired, never reused
The retired ID is the important half. A dropped column’s ID is never handed to a new column, so an old data file that still contains column 3 can never be misread as containing the new one. Rows written before any of these statements stay readable with no rewrite.
This is the same idea as hidden partitioning one level down: keep the real identity in metadata, leave the physical files alone, and the operations that are expensive in a directory-and-name-based table become metadata edits.
How is a row deleted without rewriting the file?
Two strategies exist. Copy-on-write rewrites every data file that contains an affected row, so readers stay simple. Merge-on-read writes a small file recording the deletion and leaves the data file alone, so writes stay cheap and the reader does the reconciling.
write.delete.mode, write.update.mode and write.merge.mode all default to
copy-on-write in 1.11.0, so a fresh table rewrites whole data files on
DELETE, UPDATE and MERGE.
ALTER TABLE prod.db.trips SET TBLPROPERTIES (
'write.delete.mode' = 'merge-on-read',
'write.update.mode' = 'merge-on-read',
'write.merge.mode' = 'merge-on-read'
);
Under merge-on-read, a delete becomes a file rather than a rewrite:
| Kind | What it names | Trade |
|---|---|---|
| Position deletes | A data file path and the row positions in it | Precise and cheap to apply, but needs positions |
| Equality deletes | Column values, for example trip_id = 'trip-1001' |
No positions needed, but the reader applies the predicate more widely |
| Deletion vectors | A bitmap per data file, stored in a Puffin file, Iceberg’s format for statistics and index blobs | Format version 3 and later, written by BaseDVFileWriter |
Merge-on-read buys low write latency at the cost of read-side reconciliation
and a compaction job you must operate. The compaction here is
rewrite_position_delete_files, and a growing count in the delete_files
metadata table is the signal it is not running.
What does a reader do with a delete file?
Under merge-on-read the data file keeps the deleted rows, so the reconciliation
happens at read time. The position_deletes metadata table shows exactly what
the reader is given.
Delete three rows from a thousand-row table, then look:
DELETE FROM prod.db.t WHERE id IN (7, 9, 11);
SELECT file_path, pos FROM prod.db.t.position_deletes;
-- three rows, pos = 7, 9, 11
SELECT count(*) FROM prod.db.t.delete_files; -- 1
SELECT count(*) FROM prod.db.t; -- 997
The delete file names a data file and the ordinal positions inside it, not the rows themselves. A scan reads the data file, loads the positions for it, and skips them. That is why position deletes are cheap to apply and why they need a position: the reader has to know where in the file to skip, not what to match.
Manifests distinguish the two kinds. The content column is 0 for a data
manifest and 1 for a delete manifest, so a reader knows which it is holding
before opening it:
SELECT content, path FROM prod.db.t.manifests;
Equality deletes are the alternative when positions are not available, for example a streaming writer that knows the key but not where it landed. They name column values instead, and the reader applies that predicate across every file the delete could affect, which is the cost of not knowing the position.
Either way the work is repeated on every scan until compaction folds the deletes
into the data files. That is what rewrite_position_delete_files is for, and
why delete_files growing faster than it shrinks is the number to watch.
Reading it from a running table
Every claim above is queryable, because Iceberg exposes its own metadata as
tables. MetadataTableType lists sixteen; these three answer most questions.
-- What is the history, and how big is the metadata?
SELECT committed_at, snapshot_id, operation FROM prod.db.trips.snapshots
ORDER BY committed_at DESC LIMIT 10;
-- Small files and skew, per partition
SELECT partition, file_count, record_count
FROM prod.db.trips.partitions ORDER BY file_count DESC LIMIT 20;
-- Is merge-on-read accumulating deletes faster than compaction clears them?
SELECT count(*) AS delete_files FROM prod.db.trips.delete_files;
| Symptom | What it points at |
|---|---|
| Planning slower than the scan | Too many small files, or manifests never rewritten |
| Storage far larger than the data | Snapshots never expired, so nothing is reclaimable |
| Filters not pruning | Unsorted data, so column bounds overlap |
CommitFailedException: Cannot commit: stale table metadata |
Writer contention, thrown by BaseMetastoreTableOperations. Expected under optimistic concurrency |
| Files on storage no snapshot references | A commit whose outcome was unknown, or a failed job |
How do branches and tags work?
Time travel by snapshot ID answers “what did this look like then”. Named
references answer “which state do we treat as meaningful”, and they are ordinary
rows in the refs metadata table.
-- pin a snapshot by name, so retention cannot expire it out from under you
ALTER TABLE prod.db.trips CREATE TAG release_1 AS OF VERSION 7241925443479918015;
-- a branch, which can move independently of main
ALTER TABLE prod.db.trips CREATE BRANCH audit AS OF VERSION 7241925443479918015;
SELECT name, type FROM prod.db.trips.refs;
-- main BRANCH
-- release_1 TAG
-- audit BRANCH
Reading from one is the same VERSION AS OF syntax, taking a name instead of a
number:
SELECT count(*) FROM prod.db.trips VERSION AS OF 'release_1'; -- the pinned state
SELECT count(*) FROM prod.db.trips; -- main, which has moved on
main is itself a branch, which is the detail that makes the model click: the
current state of the table is just the snapshot one named reference points at.
A tag is a reference that does not move, a branch is one that does.
Two practical uses. A tag pins a state that must survive snapshot expiry, which
is the honest way to keep a reproducible dataset rather than disabling
retention. A branch lets a job write and validate without anyone seeing it,
because nothing reading main observes a branch commit.
Why does the metadata keep growing?
Every commit in the previous section wrote a new metadata.json, a new
manifest list and at least one new manifest. None of them overwrote anything,
because immutability is what makes snapshot isolation and time travel work. The
consequence is arithmetic. Twelve single-row inserts into a fresh table:
after create metadata.json= 1 avro (manifests + lists) = 0
after 12 inserts metadata.json= 13 avro (manifests + lists) = 24
snapshots=12 data files=12
One metadata.json per commit, plus one manifest and one manifest list each.
Twelve rows of data are now described by 37 metadata files. Nothing cleans them
up, because both cleanup properties are unset by default:
| Property | Default | What it does |
|---|---|---|
write.metadata.delete-after-commit.enabled |
false |
Delete old metadata.json files on commit |
write.metadata.previous-versions-max |
100 |
How many to keep when deletion is on |
That default is deliberate and it is not what a streaming table wants. Turning it on bounds the growth immediately:
CREATE TABLE c.db.t2 (id BIGINT) USING iceberg
TBLPROPERTIES ('write.metadata.delete-after-commit.enabled'='true',
'write.metadata.previous-versions-max'='3');
after the same 12 inserts -> metadata.json=4
Four instead of thirteen: three previous versions plus the current one.
This only bounds metadata.json. Snapshots, manifests and the data files
they point at are untouched, and they are the larger problem on a table that
commits every few minutes. Three procedures handle the rest, and they do
different jobs:
| Procedure | Removes | Run it when |
|---|---|---|
expire_snapshots |
Old snapshots and the files only they referenced | Always, on a schedule |
remove_orphan_files |
Files on storage no snapshot references | Occasionally; it lists storage, so it is expensive |
rewrite_manifests |
Nothing; it recombines manifests into fewer, larger ones | Planning gets slow and manifest counts are high |
The ordering matters. expire_snapshots is what actually reclaims space, and
it is also what makes time travel to an expired snapshot impossible, so the
retention you choose is a recovery-window decision rather than a storage one.
Run remove_orphan_files with a conservative age threshold, because a file
that looks orphaned may belong to a write that is still in flight.
What happens when two engines commit at once?
Iceberg’s commit is optimistic: a writer builds a new metadata file and then swaps the table’s pointer, checking that the version it started from is still current. If someone else got there first, the swap fails. What happens next depends on the operation, and the difference surprises people.
Six concurrent writers against the same table, varying only what they do:
| What the writers did | Failures | Exception |
|---|---|---|
INSERT, default retries |
0 of 6 | none |
INSERT, commit.retry.num-retries=0 |
5 of 6 | CommitFailedException |
UPDATE on overlapping rows, default retries |
5 of 6 | ValidationException |
Three things follow, and the middle row is the one that explains the other two.
Concurrent appends succeed because Iceberg retries, not because they never
collide. With retries disabled, five of six fail with
CommitFailedException: Version 3 already exists. They were colliding all
along; the default of four retries absorbed it. An append is safe to retry
because adding files to a newer snapshot is still correct.
Concurrent overwrites fail, and retrying does not help. The exception is a different one:
org.apache.iceberg.exceptions.ValidationException:
Found conflicting files that can contain records matching...
That is not a lost race. The writer’s operation was defined against rows that another commit has since changed, so replaying it would produce a different and probably wrong result. Iceberg refuses rather than guess. The fix is application-level: partition the work so writers do not overlap, or serialise the overwrites.
Distinguish the two exceptions when you see them in a log.
CommitFailedException means try again. ValidationException means your
assumption about the data is stale, and trying again unchanged will fail again
or, worse, succeed and lose someone’s write.
None of this is Spark-specific. The commit protocol lives in the Iceberg library and the catalog, so Spark, Trino and Flink writing to one table go through the same mechanism and produce the same two exceptions. What differs is the catalog’s atomic-swap implementation, which is why the catalog choice matters: a Hadoop catalog relies on atomic rename in the filesystem, which object stores historically did not provide, while a REST, Hive or Glue catalog makes the swap a transaction in a real service.
When is a simpler choice better?
Append-only, single-engine, small tables that never need a row corrected. Plain partitioned Parquet with a metastore is less machinery and will not surprise you. Iceberg starts paying for itself when you need atomic commits across many files, row-level mutation, time travel, or planning that does not scale with partition count.
Workloads that are keyed upserts at high frequency. Iceberg has no record
index: a MERGE plans a join whose cost scales with how much of the table it
touches. A format built around record identity will do that work more cheaply.
Anywhere you cannot run a catalog with an atomic swap. Without one you have the file layout but not the guarantee, which is the part worth having.
Common misconceptions
“Snapshot ids are sequential, so the smallest is the oldest.” They are
random 64-bit values. On one test table min(snapshot_id) was the third of five
snapshots in commit order, so a time-travel query built on it silently read the
wrong point in history rather than failing. Order the snapshots metadata table
by committed_at.
“A snapshot is a copy of the table.” It is a pointer to a manifest list, and manifests are shared between snapshots. That is why keeping many snapshots costs metadata rather than data, and why expiring them frees far less space than the snapshot count suggests.
“Hidden partitioning means I do not have to think about partitioning.” It means you do not have to add a derived column or filter on it. The partition spec still decides how files are grouped, a bad spec still produces skewed or tiny files, and evolving it does not rewrite what is already written.
“The catalog is just where the table is registered.” It is where commit
atomicity comes from. Nothing inside the table says which metadata.json is
current, so the swap of that pointer is the commit, and its guarantees are the
catalog’s rather than the format’s.
Production tips
- Choose the catalog as deliberately as the format. It supplies the commit guarantee; a filesystem catalog on storage without atomic rename does not.
- Schedule
rewrite_data_filesbefore planning gets slow, not after. Compaction is cheaper run often. - Declare a write order with
WRITE ORDERED BYif the table is queried with predicates, otherwise column bounds prune nothing. - Expire snapshots deliberately. Long enough for time travel and incremental reads, short enough that storage and erasure obligations stay bounded.
- Run
remove_orphan_filesconservatively, with a generousolder_than. Files from an unknown-outcome commit may still be referenced. - Know that the row-level defaults are copy-on-write. If your workload deletes often, set the three modes explicitly.
- Watch
delete_filescount on merge-on-read tables, not query duration.
Frequently asked questions
Can I point an engine at the directory and skip the catalog?
Some engines allow it for read-only access using the newest metadata.json they
can find, but you lose the commit guarantee and the ability to write safely. The
catalog is not a convenience layer.
What is the difference between a snapshot and a metadata file?
A metadata.json holds the whole table state including the list of snapshots. A
snapshot is one entry inside it, naming a manifest list. One commit writes one
new metadata.json containing one new snapshot.
Why does my table have files that no query reads?
Either an older snapshot still references them and retention has not expired it,
or they came from a commit whose outcome could not be determined. The second
kind is what remove_orphan_files is for.
Is a commit conflict a bug? No. Under optimistic concurrency, a writer that loses the race does its work and then fails, which is the system behaving correctly. The pathology is the opposite: two writers both succeeding and one silently losing rows, which is what a catalog without an atomic swap allows.
How does format version affect this? 1.11.0 creates tables at version 2 and reads and writes up to version 4. Version 3 adds deletion vectors and is the minimum for row lineage. The version only moves forward, so check every reader supports it first.
Conclusion
Back to the two writers colliding, and the cleanup that deleted files someone
still wanted. Neither is about files. The first is optimistic concurrency doing
its job, and the second is the difference between a commit that failed and one
whose outcome was never determined, a distinction the commit contract makes
explicit and most descriptions of Iceberg skip.
The idea that organises the rest is that everything is immutable and addressed through one pointer. Snapshots are cheap because manifests are shared rather than copied. Planning is fast because partition values and column bounds live in metadata that can be read without touching data. Partition evolution works because a manifest records the spec it was written under. Time travel is not a feature that was added; it is what you get when nothing is ever overwritten.
The cost is the catalog. Iceberg asks you to run a service that can do an atomic compare-and-swap, and in exchange it gives you a table that many engines can read safely at once. If your workload is high-frequency keyed upserts, a format with a record index will do that specific job more cheaply. If it is large batches, many readers and layout decisions you expect to revise, this architecture is hard to beat.
References
- Iceberg table specification for snapshots, manifests, delete files and format versions
TableOperations.java, the commit contract quoted aboveTableProperties.javafor every property and its default- Spark procedures for the maintenance surface this post describes
- Apache Iceberg on Spark: the complete cheat sheet for the config surface this post explains
- Copy-on-write or merge-on-read for the delete modes measured against each other
- Running Iceberg in production for the settings and maintenance this architecture demands
- A local Iceberg playground for a three-service environment to try all of this against
Trademarks
Apache Iceberg, Apache Spark, Apache Flink, Apache Hive, Apache Parquet, Apache Avro, Apache Hudi and Apache are either registered trademarks or trademarks of The Apache Software Foundation in the United States and other countries. Delta Lake is a trademark of the Linux Foundation.
Found this useful?
These posts and tools are free. If one saved you an afternoon, you can buy me a coffee.