What happens to existing Iceberg data files when you ALTER TABLE ADD PARTITION FIELD?
answer
- the DDL finishes instantly on a huge table
- only new writes see the new layout
- every manifest remembers which spec it used
- planning projects predicates per spec
basics
~20 sNothing is rewritten. Iceberg appends a new partition spec to table metadata and points new writes at it; files already written keep the spec id they were written under, and scan planning handles each spec separately.
solid answer
~50 sPartition spec evolution in Iceberg is a metadata-only operation. `ALTER TABLE ... ADD PARTITION FIELD region` adds a new spec to the `partition-specs` list and moves `default-spec-id` to it. Every manifest records the spec id its entries belong to, so files written yesterday keep their old partition tuple and files written after the change carry the new one — no data is rewritten and no query changes. At plan time Iceberg projects each predicate through the spec that the manifest was written with, so old files still prune on the fields they do have, and prune on the new field only through per-file column bounds. The same applies to `DROP PARTITION FIELD` and `REPLACE PARTITION FIELD days(ts) WITH hours(ts)`. If you want the new layout applied to historical data, that is a separate rewrite of those files, not something the DDL does.
code
sql · 9 lines-- metadata-only: no data file is read or written
ALTER TABLE prod.db.events ADD PARTITION FIELD region;
-- change granularity going forward
ALTER TABLE prod.db.events
REPLACE PARTITION FIELD days(event_ts) WITH hours(event_ts);
-- stop partitioning new writes by day
ALTER TABLE prod.db.events DROP PARTITION FIELD days(event_ts);go deeper
Recall that changing an Iceberg table's partitioning is a metadata change: the statement returns immediately and no existing data is rewritten. New data uses the new layout.
Explain the mechanics: a new spec is appended, the default spec id moves, manifests record which spec their entries belong to, and planning projects predicates through each spec separately.
Bring the operational consequences: mixed-spec planning, why old data does not prune on a newly added field, when a rewrite of the hot range is worth its IO cost, and the file-count blowup a new high-cardinality field can cause.
The strategic point is that layout decisions stop being irreversible. Argue for shipping a reasonable spec early, measuring, and evolving, rather than long design debates about partitioning before any data exists.
## The DDL With the Iceberg Spark SQL extensions the three operations are: ```sql ALTER TABLE prod.db.events ADD PARTITION FIELD region; ALTER TABLE prod.db.events ADD PARTITION FIELD bucket(16, user_id); ALTER TABLE prod.db.events DROP PARTITION FIELD days(event_ts); ALTER TABLE prod.db.events REPLACE PARTITION FIELD days(event_ts) WITH hours(event_ts); ``` Each one completes in the time it takes to write a new table metadata file. None of them reads or writes a single data file. ## What changes in metadata Table metadata holds a list of partition specs, each with a `spec-id`, and a `default-spec-id` naming the one writers must use. Evolving the spec appends a new entry and moves the default; the old entries stay forever, because files written under them still exist. ```json { "default-spec-id": 1, "partition-specs": [ {"spec-id": 0, "fields": [ {"name": "event_ts_day", "transform": "day", "source-id": 3, "field-id": 1000}]}, {"spec-id": 1, "fields": [ {"name": "event_ts_day", "transform": "day", "source-id": 3, "field-id": 1000}, {"name": "region", "transform": "identity", "source-id": 4, "field-id": 1001}]} ] } ``` Two details are worth noticing. The transform names in metadata are singular (`day`, `hour`, `bucket[16]`, `truncate[10]`), while the Spark DDL spells them as functions (`days(...)`, `hours(...)`). And each partition field references its source column by id, not by name, so renaming the source column later does not break the spec. ## What does not change Existing data files are untouched: same paths, same contents, same partition tuple. Existing manifests are untouched too, and every manifest carries the spec id of the entries it holds. Existing snapshots are untouched, so time travel to a snapshot from before the change reads exactly the files and partition values it always did. Queries are untouched, because with hidden partitioning no query ever referenced the partition value in the first place — that is precisely what makes evolution possible. ## How planning copes with several specs A scan may span manifests written under spec 0 and spec 1. Iceberg projects the query's predicates through each spec separately: manifests bound to spec 0 are filtered with spec 0's projection, spec 1's with spec 1's. The practical consequences are worth stating explicitly. - After adding `region`, a query filtering on `region` prunes new files by partition and old files only by their recorded column bounds. Old data is not suddenly organised by region, so it may be read in full. - After `REPLACE PARTITION FIELD days(event_ts) WITH hours(event_ts)`, a one-hour predicate prunes new files to that hour and old files to the containing day, which is the correct and expected outcome. - After `DROP PARTITION FIELD`, files written before the change are still partitioned by the dropped field and still prune on it; files written after are not. ## The v1 wrinkle In v2 tables partition fields carry explicit, unique field ids, so a field can be dropped cleanly. Older v1 tables keep a dropped field in the spec with the `void` transform, which always produces null, so the shape of the partition tuple is preserved. Seeing `void` in a spec is therefore a sign of history, not a mistake. ## When you do want old data reorganised Spec evolution changes the future, not the past. If historical data must be laid out the new way — because the old granularity makes a common query read too much — that is a rewrite of those files under the new spec, done by the table's maintenance procedures rather than by DDL, and it costs full read-and-write IO for the range you rewrite. Deciding whether to rewrite is a cost question: often the honest answer is to evolve the spec, let old partitions age out under the retention policy, and rewrite only the range that is actually hot. ## Failure modes and interview traps The common wrong answer is that Iceberg rewrites data to match the new spec, or that old data becomes unreadable. Neither is true, and the second confusion usually comes from Hive experience where the layout is the contract. A subtler mistake is evolving the spec repeatedly: a table with many specs is legal but harder to reason about, and every scan must consider each one. Another is expecting an immediate query improvement — the win only appears for data written after the change. Finally, remember that adding a partition field increases the number of output files per write, since each commit now produces at least one file per new partition combination; on a small table that can turn one healthy file into dozens of small ones.
- After replacing days with hours, why does an hour-long predicate still read a whole day of old files?Files written before the change carry a day-level partition value, and their manifests are bound to the old spec, so the predicate can only be projected to day granularity for them. They prune to the containing day and are filtered further only by per-file column bounds. New files prune to the hour. Rewriting the old range under the new spec is the only way to change that.
- Does adding a partition field break time travel to older snapshots?No. Snapshots reference manifests that already exist, and each manifest records the spec it was written with, so reading an old snapshot reproduces the layout of that moment exactly. Spec evolution only appends to the spec list and moves the default for new writes.
- What is the operational risk of adding a partition field to a small, frequently written table?File count. Each commit now writes at least one file per distinct combination of partition values, so a high-cardinality identity field can turn one well-sized file per commit into many tiny ones. That is a small-files problem created by DDL, and it usually shows up within hours of the change.
saying these in an interview costs you the question
- Says Iceberg rewrites existing data to match the new spec
- Claims old files become unreadable after evolving the spec
- Expects historical queries to speed up immediately
- Thinks queries must be rewritten to name the new partition field
- Believes only one partition spec can exist per table