After changing a table's partition spec, when do you rewrite historical data under the new layout?
answer
- the default answer is no
- measure how much reading touches old data
- rewrite the window queries actually reach
- retention may solve it for free
basics
~20 sRewrite history only when old data is queried often enough that the mixed layout measurably hurts, and the rewrite fits your compute budget. Otherwise let the new layout govern new writes and let retention retire the old data.
solid answer
~50 sTreat it as a cost-benefit decision with a default of *no*. The cost is concrete: re-reading and re-writing every historical file, the compute to do it, the storage held by both old and new files until snapshots expire, and the conflict risk of a long-running rewrite competing with live writers. The benefit is only real for queries that actually reach back into that history — and on most event tables the vast majority of reads bound themselves to recent time, so the mixed layout is invisible. The pragmatic middle path is a **bounded backfill**: rewrite the trailing window that queries genuinely touch, leave colder data alone, and let retention delete it under the old layout. Rewrite everything only when the table is a reference dataset queried across its full history, or when carrying multiple layouts has itself become an operational burden.
code
text · 7 linesquery history, bytes scanned by time range touched
last 7 days 68%
8-90 days 27%
older than 90d 5%
=> backfill the trailing 90 days, leave the 5% tail
under the old layout and let retention expire itgo deeper
Recall that changing a table's partition layout does not oblige anyone to rewrite old data, and that leaving history as it is remains a valid, common choice.
Explain the mechanics of the cost: every historical file is read and written again, both copies occupy storage until old snapshots expire, and the rows themselves are unchanged.
Be ready to justify a bounded backfill with evidence from query history, and to run it as resumable independently-committed units that avoid conflicting with live ingestion.
Own the policy: define when a layout change earns a rewrite budget, how many live layouts a table may accumulate, and how retention is used as the zero-cost path back to a single layout.
## The decision, framed properly Changing a table's partition layout applies to new writes; existing files keep the layout they were written under, and readers plan each era separately. That leaves one open question: do you go back and rewrite history so the whole table shares one layout? The instinct is to say yes for tidiness. The engineering answer is to make it a cost-benefit call, because the costs are large, immediate and certain, while the benefits accrue only to a query population you can actually measure. ## What the rewrite costs - **Compute proportional to the whole table.** Every historical file is read, redistributed and written again. On a multi-terabyte table this is hours to days of cluster time, and it is pure overhead — the rows are unchanged. - **Write amplification and storage.** The new files exist alongside the old ones until the snapshots referencing the old files are expired and the files are cleaned up. Plan for a period of roughly doubled storage, and make sure your retention policy is what releases it. - **Contention with live writers.** A long rewrite of a partition that ingestion is still appending to invites commit conflicts, and a conflict at the end of a multi-hour job is an expensive way to learn this. Rewrite closed periods, or accept retries and checkpoint the work into small units. - **Snapshot and time-travel churn.** A bulk rewrite creates a large volume of metadata change. It does not invalidate time travel, but it does mean the snapshot history is now dominated by an operation that changed no data, and expiring those snapshots is what eventually reclaims the space. ## What the rewrite buys Exactly one thing: queries that read pre-change data get the pruning quality of the new layout. So the question reduces to "how much of our read volume touches pre-change data, and how much would it save?" Answer it with evidence from query history: the distribution of time ranges queried, the bytes scanned by queries reaching past the change point, and how much of that would be eliminated. On a typical append-only event table the answer is a small tail, and the rewrite does not pay. On a slowly-changing reference table that every join reads in full, the answer can be the opposite. ## The three viable policies **1. Forward-only (the default).** Declare the new layout, rewrite nothing. History keeps its old pruning, which is worse but was already what you had. The mixed layout self-heals as retention deletes old data — for a table with a 13-month retention, the problem is gone in 13 months at zero cost. **2. Bounded backfill.** Rewrite the trailing window that real queries touch — the last 90 days, say — and leave the rest. This captures most of the benefit for a small fraction of the cost, and it is the answer that usually reads as the most experienced. Run it as a sequence of small, independently committed units so it is resumable and does not hold a single enormous transaction. **3. Full rewrite.** Justified when the table is queried uniformly across its history, when the number of live layouts has grown enough that plans and operator reasoning suffer, or when the old layout is actively harmful (millions of tiny partitions dragging every query's planning). Schedule it as a project with a compute budget, not as a maintenance job that runs "when there's capacity". ## Guardrails whichever you choose - **Rewrite must be a data-preserving reorganisation**, not a re-derivation of rows from source. Re-running the original pipeline over old input risks changing values — late-arriving data, changed business logic, a fixed bug — and turns a layout change into a silent data change. Read the table, write the table. - **Verify by counting, not by eyeballing.** Row counts and a few aggregate checksums per period before and after, compared automatically. - **Do it in units that commit independently** so a failure halfway leaves a table that is half-rewritten and entirely correct, rather than one long transaction that must be redone. - **Watch downstream incremental consumers.** Anything that follows the table's changes will see a very large volume of rewritten files that contain no new rows; make sure such consumers distinguish a reorganisation from real change, or pause them for the duration. ## The strategic point The reason layout evolution is valuable is precisely that it *decouples* the decision to change the layout from the decision to rewrite the data. Older formats fused them and so made both impossible. A lead who understands that will change the layout today, measure, and only then decide whether history is worth the compute — rather than treating the rewrite as an obligation that comes with the change.
- Why should a backfill read the table itself rather than re-run the pipeline that originally produced it?Because re-running source logic can produce different rows: late-arriving records, changed business rules, or a bug fixed since. That converts a layout reorganisation into a silent data change nobody reviewed. Reading and rewriting the table's own rows guarantees the content is identical and only the file layout differs.
- Storage did not drop after the rewrite completed. What is the likely cause?The old files are still referenced by earlier snapshots. Space is reclaimed only when those snapshots pass their retention and expire, after which orphaned files can be cleaned up. Expect roughly doubled storage from the start of the rewrite until that expiry, and budget for it.
- How do you keep a long rewrite from conflicting with continuous ingestion?Rewrite closed periods that writers no longer touch, commit in small units so any conflict costs one unit rather than the whole job, and make the job retryable. If recent partitions must be included, expect retries and schedule those units during low write volume.
- When would you not change the layout at all, even knowing it is wrong?When the table's retention is short enough that the bad layout expires soon anyway, or when the queries hurt by it are rare and cheap in absolute terms. A layout change adds a planning era and downstream risk; if the payoff is small and time will fix it, doing nothing is a defensible engineering decision.
saying these in an interview costs you the question
- Treats rewriting all history as mandatory after a layout change
- Ignores that old and new files coexist until snapshots expire
- Re-runs the source pipeline instead of reorganising existing rows
- Runs the rewrite as one long transaction with no resume point
- Assumes a rewrite on live partitions will not conflict with writers