How can a table format change a table's partitioning without rewriting existing data?
answer
- the layout is metadata, not the directory tree
- every file remembers how it was written
- new rule for new writes only
- the reader plans each era separately
basics
~20 sThe table records which partition layout each data file was written under. Changing the layout is a metadata commit that applies to new writes only; existing files keep their old partition values, so nothing is rewritten and no downtime is needed.
solid answer
~50 sPartitioning is a property of the **table metadata**, not of the directory tree, and every data file is tagged with the layout version it was written under. Declaring a new layout is therefore a single metadata commit: subsequent writes produce files tagged with the new version, while every existing file stays exactly where it is, still correctly described by the old one. Readers plan the two groups separately — old files pruned by the old partition fields, new files by the new ones — and union the results, so queries stay correct throughout. The consequence to state out loud is that history does **not** retroactively gain the new pruning: a filter on a newly added field narrows only the files written since the change. Not every format offers this; where partition columns are fixed at table creation, changing the layout means rewriting the table or adopting a clustering mechanism instead.
code
text · 6 lines-- one table, two layout eras after the change
file spec partition values
data/00021.parquet 0 event_month=2024-03
data/00022.parquet 0 event_month=2024-04
data/00311.parquet 1 event_day=2024-05-01, user_bucket=17
data/00312.parquet 1 event_day=2024-05-01, user_bucket=42go deeper
Recall that modern formats keep the partition layout in table metadata, so changing it is a metadata update rather than a move of every file on storage.
Explain that each data file is tagged with the layout it was written under, and that a reader prunes each group of files by the rules that describe it before unioning the results.
An interviewer expects the consequence: history does not gain the new pruning, so state what improves immediately, what improves only as new data lands, and what breaks downstream if jobs read paths instead of the table.
Own the strategy: how many live layouts a table may accumulate before planning complexity outweighs the benefit, and when a layout change should be paired with a bounded backfill rather than left to age in.
## Why this was impossible before In a Hive-style table the directory tree *is* the partitioning. `dt=2024-05-01/` means "partitioned by dt" because there is nowhere else for that fact to live. So changing the partitioning means moving every file into a different tree — in practice: create a second table with the new layout, rewrite all history into it (hours to days of compute on a large table), repoint every downstream job and dashboard, and hope nothing was missed. Teams therefore did not change partitioning. They lived with the wrong key, which is exactly why "your partition key turned out to be wrong, now what?" was such a good interview question for so long. ## The change that makes it possible When the table format tracks data files in its own metadata, two facts become storable that a directory tree cannot hold: 1. The table's partitioning is a **declared specification** — an ordered list of (source column, transform) pairs — held in metadata and versioned. 2. Each data file's entry records the **partition values it holds and which version of the specification produced them**. Once those are true, the layout is no longer implied by the path. It is data, and data can be changed with a commit. ## What actually happens on the change Declaring a new layout writes a new specification version and makes it the table's current one. That is the whole operation: no files are read, none are written, none are moved. The commit is as cheap as any other metadata commit and is atomic like any other, so concurrent readers see either the old layout or the new one, never a half-applied state. From that instant: - **Writers** use the current specification. New files are partitioned the new way and tagged with the new version. - **Existing files** are untouched. They remain tagged with the version that produced them, and their recorded partition values remain accurate descriptions of their contents. ## What a reader has to do Because one table now contains files described by two different specifications, scan planning is no longer a single predicate evaluation. The planner groups files by the specification they were written under and prunes each group with the rules that apply to it: ``` spec 0 (month(event_ts)) -> prune by month spec 1 (day(event_ts), bucket(64, user_id)) -> prune by day and bucket ``` Both results are unioned. Correctness is never in question — every file is pruned by criteria that genuinely describe it — but the *quality* of pruning differs by era. ## The consequence people forget **Old data does not gain the new pruning.** If you add `bucket(64, user_id)` today, a query for one user still reads all of the pre-change history at whatever granularity the old specification provided, and only the post-change files narrow to one bucket. The improvement arrives gradually as new data accumulates, and arrives immediately only for queries that already bound themselves to recent time. The symmetric case matters too: if you *coarsen* the layout (hour to day), old finely-partitioned files still prune by hour and are fine; it is the new files that are coarser. ## Evolution is not free of judgment - **Plans get more complex.** Every additional live specification is another planning group. A table that has been re-specified five times carries five sets of pruning rules; two or three eras is normal, a dozen is a smell. - **Downstream assumptions may break.** Jobs that hardcoded a partition column name in a filter, or that enumerate directory paths directly instead of querying the table, will not follow the change. Query the table, never the paths. - **Sorting and clustering are a separate axis.** Changing the partition specification does not reorder rows within files; if pruning inside a partition matters, that is a different lever. ## Format differences worth stating This capability is not uniform. Apache Iceberg supports changing the partition specification in place with exactly the semantics above. Delta Lake does not evolve partition columns of an existing table — they are fixed when the table is created, and the evolvable alternative it offers is liquid clustering, which is a clustering mechanism rather than directory partitioning. Hudi's partitioning is determined by write configuration and is likewise not something you re-declare over existing data. In an interview, say which behaviour you are describing rather than asserting that "table formats support partition evolution" as a universal. ## The follow-up you should expect "So do you go back and rewrite the old data under the new layout?" That is a cost-benefit decision, not a requirement — and the honest default is no, unless the old data is queried often enough to justify re-reading and re-writing all of it.
- After adding a bucket field to the layout, why does a single-user lookup over last year still read a lot of data?Because only files written after the change carry bucket values. Older files were never organised by that key, so their entries offer nothing to prune on, and the planner must scan whatever the old layout narrowed them to. The benefit accrues going forward, not retroactively.
- What limits how many times you should change a table's partition layout?Every live layout is an extra planning group with its own pruning rules, so plans get wider and harder to reason about, and operators must remember which era behaves how. Two or three eras is unremarkable; a table with many is usually a sign the partition choice was never really settled.
- Which downstream consumers break when the layout changes, even though queries stay correct?Anything that bypasses the table: jobs that glob storage paths, external tools that assume a directory naming convention, and hardcoded filters on a partition column that no longer exists in the current layout. The fix is to make consumers read through the table and its catalog rather than through the filesystem.
- Does changing the partition layout invalidate time travel to older snapshots?No. Older snapshots reference the same untouched data files, and each file still carries the layout it was written under, so reading an old snapshot works exactly as before. The layout change is a metadata commit like any other and takes its place in the table's history.
saying these in an interview costs you the question
- Says existing files are rewritten or re-partitioned in the background
- Claims a filter on the new field prunes historical data too
- Thinks the change requires downtime or a table swap
- Asserts every table format supports partition evolution
- Believes changing the layout also reorders rows in existing files