Technologies that implement this pattern: PostgreSQL · Distributed SQL · DynamoDB · Cassandra · Apache Kafka · Elasticsearch
Why This Matters#
Every schema change on a live system is a distributed deploy that happens to touch a database. For some window — minutes for a column, months for a datastore move — old code and new code run side by side, against old data and new data, and every combination has to work. Most migration outages are not caused by the migration's logic. They are caused by an ALTER TABLE that waited 40 seconds for a lock and blocked every query queued behind it, a backfill that pushed replicas 15 minutes behind, or a cutover after which 0.3% of rows turned out to have been written by a code path nobody dual-wrote.
Most candidates treat migrations as a DDL question: "run the ALTER, then deploy the code." Staff engineers treat them as a compatibility-window question. The design question is not what is the target schema but what sequence of individually safe, individually reversible steps gets there, what must be true about old and new code at each step, how do I prove the data matches before I trust it, and at which step does rollback stop being free? The target schema takes five minutes to design. The sequence is the interview.
The second reframe: the risky step is never the copy; it is the switch, and then the delete. Copying data is slow but harmless — you can stop, throttle and resume it. Switching reads is the first moment users see new data, so it needs verification before it happens. Dropping the old column, table or datastore is the only truly one-way step, and it should happen weeks after everyone has stopped worrying, not the afternoon the cutover looked fine.
If you can walk an interviewer from "expand, migrate, contract" to "how each step is made non-blocking" to "how I verify before switching reads" to "where rollback ends and who signs off on crossing that line," you are answering at Staff level.
The 60-Second Version#
- Never change a thing in place — add the new thing, move to it, then remove the old thing. Expand → migrate → contract. A column rename is three deploys and a backfill, not one
ALTER. Each step must work with both the previous and next version of the code. - Locks are the outage, not the DDL. In PostgreSQL, many
ALTER TABLEforms need anACCESS EXCLUSIVElock; while one waits behind a 2-minute analytics query, every new read on that table queues behind it. Setlock_timeoutto 2–5s and retry, rather than letting one migration stall all traffic. - Backfills are throttled by replica lag, not by speed. Batches of 1–10K rows, keyed by primary key, with a pause whenever replica lag exceeds ~1s. A 1B-row table at 5K rows/s is ~55 hours — plan in days, make it resumable, and checkpoint.
- Verify before you switch reads. Shadow-read from the new path, compare to the old, and require a mismatch rate of 0 (or explained) over at least 24 hours and 1M+ comparisons before serving users from new data.
- Writes move last, and the old path stays writable until you are sure. Order: dual-write → backfill → verify → move reads → move writes → soak → contract. Rollback is free until you stop writing the old path; plan a soak of 1–4 weeks before dropping anything.
- Online schema tools exist because triggers and table copies hurt. Trigger-based copy tools add write overhead and lock contention; log-based tools (gh-ost for MySQL) copy in chunks and apply changes from the binlog, and can pause on replica lag. Use the store's native online DDL first; reach for a copy tool when the change rewrites the table.
The Problem#
A payments service needs to move currency and amount from a single amount_text column ("12.50 USD") into two typed columns, on a 2-billion-row table taking 6,000 writes/s, with 40 services reading it, three of which are batch jobs owned by other teams. The naive plan — ALTER TABLE to add the columns with a type conversion, deploy code that reads the new columns, drop the old one — rewrites the table under an exclusive lock for hours, breaks every reader still on the old code during the deploy, and leaves no way back once the old column is gone. Even the careful plan has traps: the backfill saturates I/O and replicas fall 20 minutes behind; a reporting job nobody listed writes amount_text directly and never populates the new columns; the read switch happens on a Friday and a 0.04% rounding mismatch surfaces on Monday's settlement report. The job is to break the change into steps that each hold no long lock, work with both adjacent code versions, can be paused and reversed, and are verified with production traffic before users depend on them — and to know exactly which step is the one-way door.
Case Studies That Use This Pattern#
- Sharded Database — Resharding is the largest online data migration most teams ever run: copy, catch up, cut over per shard
- Choosing a Database — Picking a store is a two-way door only if you've planned how to migrate off it
- Payments — Ledger and payment schemas change under strict correctness; dual-run and reconciliation before any switch
- Multi-Region — Schema changes must roll through regions with mixed versions live for hours
- Deployment System — Migrations ride the same staged rollout; schema and code compatibility gate each wave
- Feature Flags — Read and write path switches are flags with percentage rollouts and instant rollback
- Search Engine — Reindexing to a new mapping behind an alias: build new, dual-write, verify, swap
- Key-Value Store — Rebalancing partitions online is a migration the system runs on itself, continuously
Which Problem Are We Solving?#
"Migrate the data" hides four goals that lead to different designs. Name them and commit before drawing a box.
| Intent | Constraint | Strategy | Failure Mode | Correctness Bar |
|---|---|---|---|---|
| Schema change on a live table (add column, index, constraint, type change) | Minutes to hours; same database; no downtime | Native online DDL where it exists; lock timeouts; concurrent index builds; NOT VALID then VALIDATE; copy tools for rewrites | Lock queue blocks all traffic; table rewrite saturates I/O | Zero blocked-query incidents; change reversible until contract |
| Data model refactor (split a column, move a field to a new table, change keys) | Days to weeks; application semantics change | Expand → dual-write → backfill → shadow-read → switch reads → switch writes → contract | Missed writer path; semantic mismatch found after switch | Mismatch rate 0 on shadow reads before switch; old path writable until soak ends |
| Datastore migration (monolith to sharded, one engine to another) | Weeks to months; two systems run in parallel | Dual-write or CDC replication, bulk backfill, dark reads, per-tenant or per-shard cutover | Divergence during the long parallel run; cutover lag; cost of two systems | Per-entity checksums match; cutover reversible per slice |
| Rebalancing / resharding (move ranges or tenants between nodes) | Continuous or periodic; same store | Copy + change-stream catch-up + brief write pause + directory flip | Writes lost in the flip window; hot destination | No lost writes; write pause bounded (e.g. < 30s per slice) |
🎯 Staff Move: "This is intent two: the semantics of the field change, so I'm not looking for a clever ALTER. I'll add the new columns empty, dual-write from the service, backfill in throttled batches, shadow-read and compare for a week, then move reads behind a flag, then writes. The old column stays until a month after the switch. The only one-way step is the drop, and that one needs the reading teams' sign-off."
The Core Tradeoff#
| Strategy | What Works | What Breaks | Who Pays |
|---|---|---|---|
| Maintenance window (stop writes, migrate, restart) | Simplest; no compatibility window; easy to verify | Downtime proportional to data size; a 2B-row rewrite doesn't fit in any window; rollback means another window | Users and the business during the outage; on-call doing it at 3am |
In-place DDL (plain ALTER TABLE) | One statement; instant for metadata-only changes | Rewrites and exclusive locks on large tables; lock queue blocks reads; not pausable | Every request to that table while the lock is held or queued |
| Online schema tools (gh-ost, pt-online-schema-change, pg_repack-style copies) | Shadow table copied in chunks; throttle and pause; cut-over is a quick rename | Doubles storage temporarily; trigger-based variants add write overhead; cut-over still needs a brief lock | DBA or platform team; disk headroom |
| Expand / contract in the application | Every step small, independent and reversible; works on any store | Many deploys over weeks; code carries both paths; easy to forget the contract step | The owning team's velocity; the codebase until cleanup |
| Dual write + backfill + shadow read | Verifies correctness on real traffic before users depend on it | Dual write is itself a consistency risk; every writer must be found; doubles write load | The owning team; the new store's capacity |
| CDC replication + cutover | Captures every writer, including ones you didn't know about; no app dual-write | Lag at cutover; transformation logic lives in the pipeline; needs log access | Platform team running CDC; cutover coordination |
| Big-bang copy and switch | Fast to execute; one cutover | No verification on live traffic; rollback loses writes made after the switch | Users if the new path is wrong; the team doing the 2am rollback |
Staff Default Position#
Expand, migrate, contract — every step independently deployable, compatible with the code on either side of it, throttled, verified on production traffic before users depend on it, and reversible until the final contract step, which waits weeks and needs sign-off.
The default sequence for a model change: (1) expand — add new columns or tables, nullable, with no constraints that existing writes would violate; (2) dual-write — the service writes both shapes in the same transaction (or CDC populates the new shape); (3) backfill — a resumable job copies historical rows in primary-key order, 1–10K rows per batch, paused whenever replica lag exceeds ~1s; (4) verify — shadow-read the new shape on a sample or all of production traffic and compare, plus a full offline checksum; (5) move reads behind a percentage flag; (6) move writes, keeping the old shape written for rollback; (7) soak for 1–4 weeks; (8) contract — stop writing, then drop. Schema DDL along the way uses lock_timeout, concurrent index builds and NOT VALID constraints; any change that rewrites a large table goes through an online copy tool. Every migration has an owner, a written step list with the rollback for each step, and a named point of no return.
When to Deviate#
- Small tables and quiet systems — below ~1M rows or in a system with real off-hours and an SLA that permits it, a plain
ALTERor a 5-minute maintenance window is cheaper than three deploys. Measure the table and the lock first; don't ceremony a 40K-row config table. - Metadata-only changes — adding a nullable column, or one with a constant default in PostgreSQL 11+, or an
INSTANTcolumn add in MySQL 8.0, doesn't rewrite the table. Still setlock_timeout, because even an instant change needs the lock briefly. - The old system can't be dual-written — a vendor system or a store you can only read. Use CDC or periodic snapshots plus an incremental catch-up, and accept a short write freeze at cutover.
- Regulatory or financial records — ledgers where "both paths for a month" creates two sources of truth auditors will ask about. Run the new path as a pure shadow with reconciliation, and switch in a single declared cutover with a signed reconciliation report.
One Question, Three Levels#
| Behavior | Senior (L5) | Staff (L6) | Principal (L7) |
|---|---|---|---|
| First move | "Write the migration script, run it, deploy the new code" | "What's the step sequence where each step is compatible with the code on either side and reversible?" | "How many teams run migrations like this a quarter, and do they all reinvent the backfill, the verifier and the lock handling?" |
| Locks | Runs ALTER TABLE in a deploy hook | Knows which DDL forms rewrite or lock; sets lock_timeout; uses concurrent and NOT VALID forms; copy tools for rewrites | Lints migrations in CI for dangerous forms; schema changes go through a platform with automatic throttling |
| Backfill | One UPDATE ... SET new = f(old) | Resumable, primary-key-ordered batches throttled on replica lag, checkpointed, idempotent | Shared backfill framework with lag-aware throttling and progress reporting every team uses |
| Verification | Spot-checks rows after the migration | Shadow reads and full checksums before switching reads; mismatch budget of zero, investigated | Makes verification a gate in the migration platform: no read switch without a verifier report |
| Rollback | "We have backups" | Rollback plan per step; old path writable until soak ends; names the point of no return | Sets org policy on soak periods and who signs off on one-way steps |
| Ownership | The team that wrote the script | Owning team runs it; downstream readers listed and signed off; DBA reviews locks | Migration platform team owns tooling; owning teams own semantics; contract step has a cross-team sign-off |
Why "First move" separates levels
The L5 plan — script, run, deploy — works on a small table in a single-service world, and it gets downleveled because it has no answer for the deploy window. During a rolling deploy of 200 pods over 15 minutes, half the fleet runs old code against the new schema. If the migration renamed or retyped a column, half the fleet is broken. The Staff candidate designs the sequence so that every adjacent pair of (code version, schema version) works. The Principal candidate notices that the org runs dozens of these a quarter and that the backfill, verifier and lock-safety logic should be built once.
Why "Locks" separates levels
The most common migration outage isn't a slow migration; it's a fast one that waits. PostgreSQL grants locks in queue order. An ALTER TABLE needing ACCESS EXCLUSIVE waits behind a long-running SELECT, and every subsequent SELECT on that table queues behind the ALTER. The table is effectively offline until the first query finishes. A Staff engineer sets lock_timeout = '3s' so the migration fails fast and retries, rather than taking the table down for the length of someone's report.
Why "Verification" separates levels
A migration without verification has only one way to find a mismatch: a user or a downstream report. Shadow reads — run both paths, return the old, compare and log — find semantic differences (rounding, time zones, null handling, a writer you missed) while the old path is still the source of truth and nothing is broken. The Staff bar is "I won't switch reads until the comparison has been clean on real traffic for a defined period," with the period and sample size stated.
Where the Design Splits#
| # | Fault Line | The Tension |
|---|---|---|
| 1 | Database-Level Online DDL vs Application-Level Expand/Contract | One tool-driven operation vs many small deploys that change semantics safely |
| 2 | Application Dual-Write vs Log-Based Replication | Explicit control and same-transaction writes vs catching every writer without touching code |
| 3 | Big Switch vs Gradual Cutover | One coordinated moment vs percentage or per-tenant ramps with two paths live |
| 4 | Sampled Verification vs Full Reconciliation | Cheap continuous confidence vs exhaustive proof that costs a full scan |
| 5 | Keep the Old Path Writable vs Cut Over One-Way | Free rollback for weeks vs paying double writes and two sources of truth |
Fault Line 1: Database-Level Online DDL vs Application-Level Expand/Contract#
Some changes are purely physical — add an index, widen a column, add a constraint — and the right tool is the database's own online path or a copy tool. Others change meaning — split a field, change a key, move data between tables — and no database tool knows the semantics. Who pays: forcing a semantic change through a single DDL makes users pay during a lock or rewrite and makes the team pay with no verification; forcing a physical change through expand/contract makes the team pay three deploys for something CREATE INDEX CONCURRENTLY does in one. Staff default: physical change → native online DDL or a copy tool; semantic change → expand/contract with dual-write and verification. Deviate when: the store has no safe online DDL for the change (e.g. a type change that rewrites a 2B-row table) — then the physical change also becomes expand/contract: new column, backfill, swap.
Fault Line 2: Application Dual-Write vs Log-Based Replication#
Dual-writing from the application gives same-transaction consistency when both shapes live in one database, and lets you write transformed data with full business context. CDC replication from the old store's log catches every writer — batch jobs, admin scripts, other services — without code changes. Who pays: application dual-write pays when a writer is missed (the batch job nobody listed) and, across two databases, it is a dual write with all the divergence that implies; CDC pays in pipeline lag at cutover, transformation logic living outside the service, and a new piece of infrastructure. Staff default: same-database model changes use in-transaction dual-write; cross-store migrations use CDC (or an outbox) from the old store as the source of truth, with the application cutting over reads and then writes. Deviate when: the old store has no accessible log — then dual-write plus a periodic reconciliation scan, and plan a short write freeze at cutover.
Fault Line 3: Big Switch vs Gradual Cutover#
A big switch flips all reads (or writes) at once. A gradual cutover moves 1% → 10% → 50% → 100% of traffic, or one tenant or shard at a time. Who pays: big switch makes every user pay if the new path is wrong, but keeps the period of two live paths short; gradual cutover makes the team pay in complexity — both paths serve production for days, caches and read-your-writes must work across them — but caps the blast radius of a mistake at the current percentage. Staff default: gradual for reads (they are safe to mix if both paths have the same data), per-slice for writes in datastore migrations (one tenant or shard at a time, with a brief pause per slice), single switch for writes within one database. Deviate when: the two paths can't serve the same entity consistently at the same time — then cut over by entity (tenant, shard), never by request percentage.
Fault Line 4: Sampled Verification vs Full Reconciliation#
Shadow reads compare what users actually request, continuously. Full reconciliation scans both sides and compares checksums per row or per chunk. Who pays: sampling pays in blind spots — cold rows nobody reads aren't checked; full reconciliation pays in load (a full scan of 2B rows on both sides, best run on replicas or snapshots) and in time (hours to days). Staff default: both — shadow reads at 1–100% of traffic from the moment dual-write starts, plus a chunked checksum reconciliation (e.g. hash per 10K-row PK range, drill into mismatching ranges) before the read switch and again before the contract step. Deviate when: the data is derived and can be rebuilt (a search index, a cache) — then rebuild and spot-check rather than reconcile.
Fault Line 5: Keep the Old Path Writable vs Cut Over One-Way#
After writes move to the new path, you can keep writing the old one (reverse dual-write) so rollback is a read switch, or stop immediately. Who pays: keeping it writable pays in double write load and in a period where two places claim to be the truth; cutting over one-way pays the day a bug appears and rollback means losing every write since the switch or replaying them by hand. Staff default: keep the old path written for a soak period of 1–4 weeks, scaled to how much is at stake and how long the slowest consumer (month-end reports, quarterly jobs) takes to exercise the new path. Deviate when: keeping the old store costs real money at scale (a second database cluster at $50K+/month) — then shorten the soak, but run full reconciliation before decommissioning and keep a snapshot.
Common Interview Mistakes#
| What Candidates Say | What Interviewers Hear | What Staff Engineers Say |
|---|---|---|
| "Run an ALTER TABLE to rename the column" | "Hasn't seen a rolling deploy break on a rename" | "Add the new column, dual-write, backfill, switch reads, switch writes, then drop the old one weeks later." |
| "We'll run the migration during low traffic" | "Lock queue risk not understood" | "Low traffic doesn't help if one long query holds the table. I'll set lock_timeout to 3 seconds and retry." |
| "Backfill with one UPDATE statement" | "One giant transaction and a replica lag spike" | "Batches of 5K rows in primary-key order, checkpointed, paused whenever replica lag passes one second." |
| "We'll verify after the migration" | "Users will find the bugs" | "Shadow reads compare old and new on live traffic for a week before any user is served from the new path." |
| "We can always restore from backup" | "Rollback loses every write since the switch" | "The old path stays written for a month. Rollback is a flag flip until we drop it, and the drop needs sign-off." |
| "Dual-write to both databases" | "Two commits, one will fail alone" | "Across stores I'll replicate from the old store's log, so every writer is captured and there's no second commit in the request path." |
Quick Reference#
Staff Sentence Templates#
"This change alters [meaning / layout], so it goes through [expand-contract / online DDL]. Each step works with the code on either side, and the only one-way step is [dropping X], which waits [N weeks] and needs [reading teams'] sign-off."
"Every DDL statement runs with lock_timeout [3s] and retries, because one long query holding [table] would otherwise queue every read behind the migration."
"The backfill is [rows] rows at about [rate] rows per second, so roughly [hours]. It runs in primary-key batches of [N], checkpointed, and pauses whenever replica lag exceeds [1s]."
"I won't switch reads until shadow comparison has been clean for [duration] and [count] reads, and a full chunked checksum matches. Rollback until the contract step is [a flag flip]."
Implementation Deep Dive#
1. Lock-Safe DDL — PostgreSQL#
Most PostgreSQL DDL is fast; the danger is the lock it needs and the queue it creates while waiting. Every statement gets a short lock_timeout, and the slow forms get their non-blocking variants.
-- session settings for every migration
SET lock_timeout = '3s'; -- fail fast instead of queueing all traffic behind us
SET statement_timeout = '0'; -- long-running concurrent builds are fine; locks are not
-- add a column: metadata-only in PG 11+ when the default is a constant
ALTER TABLE payments ADD COLUMN amount_minor bigint; -- nullable, instant
ALTER TABLE payments ADD COLUMN currency char(3) DEFAULT 'USD'; -- constant default, no rewrite
-- index: never block writes on a big table
CREATE INDEX CONCURRENTLY idx_payments_currency ON payments (currency);
-- cannot run inside a transaction; on failure leaves an INVALID index -> DROP INDEX CONCURRENTLY and retry
-- constraints: add without scanning under a heavy lock, validate separately
ALTER TABLE payments ADD CONSTRAINT amount_minor_nonneg
CHECK (amount_minor >= 0) NOT VALID; -- brief lock, no scan
ALTER TABLE payments VALIDATE CONSTRAINT amount_minor_nonneg; -- scans with a lighter lock; writes continue
-- NOT NULL without a full-table scan under ACCESS EXCLUSIVE (PG 12+):
ALTER TABLE payments ADD CONSTRAINT amount_minor_nn CHECK (amount_minor IS NOT NULL) NOT VALID;
ALTER TABLE payments VALIDATE CONSTRAINT amount_minor_nn;
ALTER TABLE payments ALTER COLUMN amount_minor SET NOT NULL; -- uses the valid CHECK, skips the scan
ALTER TABLE payments DROP CONSTRAINT amount_minor_nn;
-- forms that rewrite the whole table under ACCESS EXCLUSIVE: avoid on large tables
-- ALTER COLUMN ... TYPE (unless binary-coercible), ADD COLUMN ... DEFAULT volatile_fn()
-- -> do these as expand/contract: new column, backfill, swap
retry wrapper (migration runner):
for attempt in 1..20:
try: run(statement); return
except LockNotAvailable:
metrics.incr("migration.lock_timeout", tags=[table])
sleep(jitter(5s × attempt)) # back off; maybe the long query finishes
page(owner, "could not acquire lock on {table} after 20 attempts") # find the blocker, don't force it
🎯 Staff Insight: A CI lint that flags dangerous forms —
ALTER COLUMN TYPE,CREATE INDEXwithoutCONCURRENTLY,ADD CONSTRAINTwithoutNOT VALIDon tables over 1M rows, any DDL withoutlock_timeout— prevents more migration outages than any amount of review, because reviewers don't remember which forms rewrite.
2. Expand / Contract — A Rename in Six Deploys#
The canonical example: rename amount_text to typed amount_minor + currency. Each row is one deploy; each must work with the code from the row above and below.
| Step | Schema | Writes | Reads | Rollback |
|---|---|---|---|---|
| 1. Expand | Add amount_minor, currency (nullable) | Old only | Old | Drop new columns |
| 2. Dual-write | Same | Old + new, same transaction | Old | Deploy previous version |
| 3. Backfill | Same | Old + new | Old; shadow-read new and compare | Stop job; new columns ignored |
| 4. Switch reads | Same | Old + new | New (ramp 1% → 100% by flag) | Flip flag back |
| 5. Switch writes | Constraints validated on new | New + old (kept for rollback) | New | Flip reads back; old is current |
| 6. Contract | Drop amount_text after 2–4 week soak | New | New | None — the one-way door |
# Step 2: dual-write inside the existing transaction
def update_payment(p):
with db.transaction():
db.execute("""UPDATE payments
SET amount_text = %s, amount_minor = %s, currency = %s
WHERE id = %s""",
format_text(p), p.amount_minor, p.currency, p.id)
# Step 3: shadow read
def get_amount(id):
row = db.fetch("SELECT amount_text, amount_minor, currency FROM payments WHERE id=%s", id)
old = parse_text(row.amount_text)
if row.amount_minor is not None: # backfilled or dual-written
new = (row.amount_minor, row.currency)
if new != old:
metrics.incr("migration.shadow_mismatch", tags=["field:amount"])
log.warn("mismatch", id=id, old=old, new=new) # sampled, PII-safe
return old # users still see the old path
The finder's job: before step 2, list every writer of amount_text — application code, batch jobs, admin tools, other services' direct queries. Database audit logs or pg_stat_statements filtered by statement text over two weeks will find the ones code search misses. A missed writer is the most common cause of post-cutover divergence.
3. Throttled, Resumable Backfill — Keyset Batches Gated on Replica Lag#
The backfill is the slow part, so it must be safe to run for days: idempotent, resumable, and self-throttling.
BATCH = 5000
MAX_LAG_S = 1.0
TARGET_BATCH_MS = 200
cursor = checkpoint.load("payments_amount_backfill") or 0
while True:
lag = max(replica_lag_seconds(r) for r in replicas)
if lag > MAX_LAG_S:
metrics.gauge("backfill.paused_for_lag", 1); sleep(5s); continue
t0 = now()
rows = db.execute("""
UPDATE payments
SET amount_minor = parse_minor(amount_text),
currency = parse_currency(amount_text)
WHERE id > %s AND id <= %s + %s
AND amount_minor IS NULL -- idempotent: skips dual-written rows
RETURNING id""", cursor, cursor, BATCH)
cursor += BATCH
checkpoint.save("payments_amount_backfill", cursor) # resume here after a crash
metrics.gauge("backfill.cursor", cursor)
elapsed = now() - t0
BATCH = clamp(BATCH × TARGET_BATCH_MS / elapsed, 500, 20000) # keep each txn short
if cursor > max_id_at_start: break
# 2B rows at ~6K rows/s sustained ≈ 93 hours. Run on weekdays only if weekend batch load
# competes; report ETA = remaining_rows / rows_per_s_1h to the migration channel daily.
Why keyset, not OFFSET: OFFSET 1,000,000 rescans a million rows each batch, so batch cost grows linearly and the job slows to a crawl near the end. Walking primary-key ranges keeps every batch the same cost. Why amount_minor IS NULL: rows touched by dual-write after step 2 are already correct; the backfill must not overwrite them with a value derived from a stale read.
4. Datastore Migration — CDC, Dark Reads and Per-Slice Cutover#
Moving to a different store (monolith PostgreSQL → sharded cluster, one engine to another) runs two systems for weeks. Replicating from the old store's log catches every writer; cutting over one slice at a time limits blast radius.
cutover(slice):
assert verifier.mismatch_rate(slice, last=7d) == 0
assert checksum(old, slice) == checksum(new, slice)
directory.set(slice, state="PAUSED") # writers get a retryable 503 for ~5-30s
wait_until(cdc.lag(slice) == 0, timeout=60s) or abort_and_unpause()
directory.set(slice, store="new", state="ACTIVE")
start_reverse_replication(slice) # new -> old, keeps rollback cheap
monitor 24h: error rate, latency, reconciliation; then next slice
Order of slices: internal and test tenants first, then a few hundred small tenants, then by size, with the largest tenants last — after the procedure has run dozens of times. See Multi-Tenancy & Noisy Neighbours for directory-based placement, which is the same mechanism.
Technique Comparison
| Technique | Blocks Traffic? | Pausable | Verifies Semantics | Rollback | Best For |
|---|---|---|---|---|---|
| Plain DDL | Yes, while locked or rewriting | No | No | Another DDL | Small tables, metadata-only changes |
| Concurrent / NOT VALID DDL | No (brief lock only) | No, but restartable | No | Drop | Indexes and constraints on big tables |
| Online copy tool (gh-ost style) | Brief cut-over lock | Yes, throttles on lag | No | Keep old table until dropped | Table rewrites, type changes in MySQL |
| Expand / contract + dual-write | No | Yes | Yes, via shadow reads | Free until contract | Model refactors in one database |
| CDC replication + per-slice cutover | Seconds per slice | Yes | Yes, via dark reads | Reverse replication during soak | Datastore migrations, resharding |
| Maintenance window | Yes, entirely | N/A | Offline only | Restore | Tiny systems, regulated one-time cutovers |
Architecture Diagram#
How to narrate it: the old path stays the source of truth until a flag says otherwise. Data reaches the new path two ways — live changes (dual-write or CDC) and history (backfill) — and both are idempotent so they can overlap. Verification runs continuously on live traffic and in bulk before each switch. Replica lag is the throttle for everything that reads the old store hard. The control plane holds the step list, the rollback for each step and the named point of no return.
Failure Scenarios#
1. The Lock Queue — A One-Second ALTER Takes the Table Down for Four Minutes#
A team adds a nullable column to orders (800M rows, 4,000 queries/s) in a deploy hook. It's metadata-only and should take milliseconds. No lock_timeout is set.
14:02:00 An analyst's ad-hoc report has been running on the primary for 3 minutes (ACCESS SHARE).
14:02:01 Deploy hook runs ALTER TABLE orders ADD COLUMN gift_note text. Waits for ACCESS EXCLUSIVE.
14:02:01 Every new query on orders queues behind the ALTER. Throughput on orders: 4,000/s -> 0.
14:02:05 Connection pool exhausted (400 connections waiting). Checkout API timeouts.
14:02:30 Health checks fail; pods restart; deploy marked failed and retried. Retry queues again.
14:06:10 The report finishes. ALTER takes 20ms. Queued queries flood in; p99 spikes for 90s.
14:07:40 Recovered. 5.5 minutes of checkout outage for a 20ms change.
Detection: pg_locks waiters on orders > 50; db.lock_wait_seconds p99; queries in pg_stat_activity waiting on Lock with a DDL at the head; checkout error rate.
Blast radius: every service touching orders — checkout, order history, fulfilment — for 5.5 minutes.
Mitigation: cancel the DDL (pg_cancel_backend) to release the queue, then cancel or move the long report.
Prevention: lock_timeout = 3s with retry on every migration; DDL lint in CI; analysts on a replica, never the primary; a statement_timeout on the primary for interactive roles.
Owner: the owning team runs the migration; the database platform owns the runner defaults and lint.
🎯 Staff Insight: The migration was correct and the change was instant. The outage was the queue. "It's a metadata-only change" is exactly the sentence that precedes this incident.
2. The Backfill That Broke Read-Your-Writes — Replicas 18 Minutes Behind#
A backfill updates 1.2B rows with 50K-row batches as fast as it can, on a primary with three read replicas serving 70% of reads.
Mon 10:00 Backfill starts: 50K rows per UPDATE, no throttle. WAL rate 15 MB/s -> 120 MB/s.
Mon 10:20 Replica apply can't keep up. Replica lag 2s -> 90s.
Mon 10:45 Lag 18 minutes. Users save a setting, refresh, see the old value. Support tickets spike.
Mon 10:50 Read-after-write logic routes to primary for 30s after a write; lag exceeds that window.
Primary takes 3x normal read load; CPU 85%.
Mon 11:05 Backfill killed. Replicas take 25 minutes to catch up.
Detection: replica.lag_seconds > 1 (warn), > 30 (page); WAL generation rate vs baseline; backfill.rows_per_s alongside lag on one dashboard.
Blast radius: all users reading from replicas — stale reads for ~45 minutes; primary near saturation.
Mitigation: stop the backfill; shift reads to primary for critical paths until lag recovers.
Prevention: lag-gated throttling (pause above 1s), adaptive batch size targeting ~200ms transactions, WAL-rate budget per migration.
Owner: owning team runs the backfill; database platform owns the backfill framework's throttle defaults.
3. The Missed Writer — 0.3% of Rows Diverged, Found After the Drop#
A team moves user preferences from a JSON blob column to a normalized table via expand/contract. Shadow reads are clean for a week; reads and writes switch; the old column is dropped after 5 days.
Week 1 Dual-write in the preferences service. Shadow reads sampled at 1%: 0 mismatches.
Week 2 Reads switched, then writes. Old column still written.
Week 2 Day 5 after switch: old column dropped to reclaim 300 GB.
Week 4 Monthly job (owned by the notifications team) runs. It writes opt-outs
directly to the old JSON column via raw SQL. Column gone: job fails loudly. Worse,
its previous month's run had written 140K opt-outs only to the blob.
Week 4 Those 140K opt-outs were never in the new table. 0.3% of users received marketing
email after opting out for 3 weeks.
Detection: a writer inventory from database audit logs before dual-write; full reconciliation (not sampled) before the contract step; per-writer metrics on the old column. Blast radius: 140K users; a compliance issue for unwanted marketing email. Mitigation: restore the blob from a snapshot, extract opt-outs, apply to the new table, suppress sends. Prevention: inventory every writer from query logs over at least one full business cycle (a month, if monthly jobs exist); full checksum reconciliation before contract; soak period at least as long as the slowest writer's cycle. Owner: preferences team owns the migration; notifications team owns its job; the contract step required their sign-off and didn't get asked.
Operational Reality Matrix#
| Failure | Detection Signal | Blast Radius | Mitigation | Owner |
|---|---|---|---|---|
| DDL lock queue | pg_locks waiters > 50 behind a DDL; lock wait p99 | Every query on the table | Cancel DDL; lock_timeout + retry | Owning team + DB platform |
| Table rewrite under lock | DDL duration; I/O saturation | Table offline for the rewrite | Expand/contract or online copy tool | Owning team |
| Backfill replica lag | replica.lag_seconds > 1s | Stale reads; primary overload | Lag-gated throttle, smaller batches | Owning team |
| Missed writer | Writer inventory gaps; reconciliation mismatches | Silent divergence | Audit-log writer inventory; CDC instead of dual-write | Owning team |
| Shadow-read mismatch | migration.shadow_mismatch > 0 | None yet — that's the point | Fix transform; re-backfill affected range | Owning team |
| Cutover lag | CDC lag at pause > 60s | Writes paused longer than budget | Abort and unpause; retry slice later | Platform + owning team |
| Contract too early | Downstream job failures after drop | Data loss for unmigrated writers | Snapshot before drop; soak ≥ slowest cycle | Owning team + downstream sign-off |
| Mixed-version incompatibility | Error spike during rolling deploy | Half the fleet during the deploy | Every step compatible with adjacent versions | Owning team |
Beyond Staff: The Principal View#
Why L7 Sees This Problem Differently#
A Staff engineer runs one migration safely. A Principal engineer notices that the org runs 30–100 of them a quarter, that each team rebuilds the same backfill loop with the same missing throttle, and that the worst outages come from the migrations nobody thought were risky — the instant column add, the "small" index. Migrations are also where technical debt compounds or gets paid: the half-finished expand/contract with both columns still live two years later, the datastore migration stuck at 80% because the last 20% is the hard tenants, the deprecated system nobody can turn off because one job still reads it. At L7 the work is making safe migrations the path of least resistance (lint, runner, backfill framework, verifier as shared tooling), making migration completion a tracked commitment rather than an afterthought, and pricing the cost of running two systems in parallel so finishing gets prioritized.
🧭 Principal Move: "We have 14 migrations open, 6 of them over a year old, and we're paying for two clusters on three of them. I want a migration platform that makes the safe path the default — lint, lock-safe runner, throttled backfill, verifier — and a quarterly review where every open migration has an owner, a finish date and a cost of not finishing."
The Org-Level Fault Line#
A migration platform with guardrails vs each team running its own migrations.
| Option | What Works | What Breaks | Who Pays |
|---|---|---|---|
| Each team runs DDL and backfills its own way | Autonomy; no waiting on a platform | Same incidents repeat across teams; no lock safety by default; half-finished migrations accumulate | On-call for whichever team forgets lock_timeout; the business paying for parallel systems |
| Central DBA team approves and runs every change | Expertise concentrated; consistent safety | Bottleneck — teams wait days for a column; DBAs lack domain context for semantic changes | Product velocity; DBA burnout |
| Platform provides tooling and guardrails; teams own semantics | Lint, lock-safe runner, backfill framework, verifier and flags are self-service; DBAs review only flagged-risky changes | Platform must support every store in use; tooling lags new engines | Platform headcount (2–4 engineers) |
The Principal default is the third row: the platform makes the mechanical safety automatic, owning teams own the step sequence and verification of meaning, and only changes the lint flags as dangerous — rewrites, very large tables, contract steps — get human review.
Cost Model#
Assumptions: fully loaded engineer ~$25K/month; a datastore migration runs two systems in parallel for its duration; storage doubling during copies; "migration engineering" is the owning team's time.
| Scale | Migration Load | Machinery | Infra Cost of Parallel Running | People | Rough Cost |
|---|---|---|---|---|---|
| Startup | ~5 schema changes/month, small tables | Migration files in CI, lock_timeout defaults | Negligible | Part of feature work | ~$0 extra; ~1 day per risky change |
| Growth | ~40 changes/month; 1–2 model refactors per quarter | DDL lint, lock-safe runner, backfill framework, shadow-read helper | Temporary 2× storage on refactored tables (~$2–5K/month during runs) | 1.5 FTE platform; refactor ~1 engineer-month each | ~$40K/month platform + ~$25K per refactor |
| Large | ~300 changes/month; 2–3 datastore migrations a year | Self-service platform, per-slice cutover tooling, verifier service, migration review | Parallel clusters during datastore moves: $50–300K/month each for 3–9 months | 4-person platform team; datastore move ~6–15 engineer-years | ~$100K/month platform; $2–6M per datastore migration all-in |
The Principal observation: at scale, the dominant cost of a datastore migration is not engineering — it is running two systems for the long tail. A migration that reaches 90% in three months and then spends nine more months on the last 10% pays for two clusters the whole time. The highest-leverage decision is often to fund the hard 10% early (largest tenants, weirdest access patterns) rather than last, and to put a decommission date on the old system that leadership defends.
The 3-Year Evolution Path#
The Year 3 trigger is the one that costs the most: the first cross-store migration without per-slice tooling tends to stall at the largest tenants, and the parallel-running bill becomes the reason leadership finally funds the platform.
One-Way Doors vs Two-Way Doors#
| Decision | Door Type | Reversibility Cost |
|---|---|---|
| Adding a nullable column, an index, a NOT VALID constraint | Two-way | Drop it |
| Switching reads to the new path | Two-way | Flag flip, as long as the old path is still written |
| Switching writes to the new path (old still written) | Two-way | Flip reads back; old is current |
| Stopping writes to the old path | One-way-ish | Rollback now requires replaying writes since the stop |
| Dropping the old column, table or datastore | One-way | Snapshot restore plus manual reconciliation, if a snapshot exists |
| Changing a primary key or shard key | One-way-ish | A second full migration in the other direction |
| Publishing the new schema to other teams (events, APIs, CDC consumers) | One-way-ish | Their code depends on it; reverting needs their migration too |
The Standard I'd Write#
RFC: Production Schema and Data Migrations (v1)
Scope: Every DDL statement, backfill, or datastore change on a production system serving users or other teams.
MUST:
- Every DDL statement MUST run through the migration runner with
lock_timeout≤ 5s and automatic retry. DDL in application startup or deploy hooks without the runner is prohibited.- Changes flagged by the lint as table-rewriting or exclusive-locking on tables over 1M rows MUST use an online path (concurrent, NOT VALID, copy tool, or expand/contract).
- Backfills MUST be idempotent, resumable from a checkpoint, and throttled on replica lag (default pause above 1s).
- Semantic migrations MUST pass shadow-read verification (zero unexplained mismatches over ≥ 7 days) and a full checksum before reads switch.
- Contract steps (dropping columns, tables or systems) MUST wait a soak period ≥ the slowest known writer's cycle (default 30 days), take a snapshot, and have written sign-off from every team that reads or writes the old shape.
SHOULD: Inventory writers from database audit logs before dual-writing; prefer CDC over application dual-write across stores; cut over datastore migrations per slice, largest last.
Exceptions: Approved by the database platform team, with a named rollback plan.
Success metrics: zero lock-queue incidents; zero data loss from contract steps; median open-migration age < 90 days; parallel-running cost reported monthly per migration.
What I'd Tell the VP#
Changing how we store data is something we do hundreds of times a quarter, and a handful of those changes cause outages or silent data problems — most recently a five-minute checkout outage from a change that should have taken a fraction of a second. I'm proposing shared tooling that makes the safe way the default way, plus a rule that nothing gets deleted until every team that depends on it has confirmed. Separately, we're currently paying to run old and new systems side by side on three migrations; I want each one to have an owner and a finish date. It costs a small platform team of about three engineers. In exchange, schema changes stop being an outage risk, and we stop paying twice for systems we meant to retire.
Principal Interview Signals#
| Signal | What It Sounds Like |
|---|---|
| Safety by default | "The lint and the runner prevent the lock-queue outage. Review doesn't scale; defaults do." |
| Prices parallel running | "Two clusters cost $150K a month. Finishing the last 10% early is the cheapest thing we can do." |
| Completion as a commitment | "Every open migration has an owner, a finish date and a cost of not finishing." |
| One-way door discipline | "The drop is the only irreversible step, so it waits 30 days and needs every reader's sign-off." |
| Hard part first | "Cut over the largest tenants while the team still has momentum, not at month eleven." |
Staff answers that L7 interviewers find insufficient:
- "I'd use expand/contract with lock timeouts." — Correct for one migration; silent on the 40 teams that will each forget the timeout.
- "We'll run both databases until we're confident." — No end date, no cost, no owner for the long tail.
- "We'll verify with shadow reads." — Doesn't make verification a gate, so the next team skips it under deadline pressure.
How Real Companies Built It#
Stripe: Four-Phase Dual Writing for Hundreds of Millions of Objects#
Stripe's write-up on online migrations describes moving subscriptions from being stored on Customer objects to a separate Subscriptions table, across hundreds of millions of records, with a four-phase dual-writing pattern: dual-write to the old and new locations, change read paths to the new location, change write paths to the new location, then remove the old data and code. For the backfill, rather than querying production, they made database snapshots available to Hadoop, used MapReduce (via Scalding) to identify the records to migrate, and ran a large parallel migration with a fleet of processes. To check the new code paths, they used GitHub's Scientist library, which runs two code paths and alerts when their results differ in production (Online migrations at scale).
Staff insight: The order of phases is the interview answer: writes are the last thing to move, not the first. And the verification step uses live production comparisons, not a test environment — the only place the real data distribution and real code paths exist.
GitHub: gh-ost and Triggerless Online Schema Changes#
GitHub built gh-ost, an open-source online schema migration tool for MySQL, because trigger-based tools caused problems at its write concurrency: triggers add per-query interpretation overhead to every write on the migrated table, contend for locks in ways that grow with write concurrency, and cannot truly be paused. gh-ost instead creates a ghost table, copies rows in chunks, and applies ongoing changes by tailing the binary log, then swaps tables at cut-over. Because it is asynchronous it can throttle on replica lag, load metrics, custom queries or a flag file, suspend writes to the primary entirely while throttled, run test migrations on a replica, and have parameters like chunk size and maximum lag changed during the run (gh-ost: GitHub's online migration tool for MySQL).
Staff insight: gh-ost is the change-data-capture idea applied to a single table: copy history in chunks, replay live changes from the log, cut over briefly. Knowing why triggers were the problem — overhead and lock contention under concurrency, and no real pause — is a stronger signal than knowing the tool's name.
Notion: Double-Writes, Backfill and Dark Reads into 480 Shards#
When Notion moved from a single PostgreSQL database to a sharded fleet of 480 logical shards on 32 physical databases, it described three phases: double-writing incoming writes to both old and new databases (using an audit log), backfilling existing data from the monolith, and verification using dark reads that compared records from both systems before switching over. The final backfill script used all 96 CPUs of the instance provisioned for it and took about three days for production, and the switchover itself needed about five minutes of scheduled maintenance (Sharding Postgres at Notion).
Staff insight: This is the datastore-migration template end to end, with real numbers: days of backfill, verification on production reads, and a cutover measured in minutes. When an interviewer asks "how long would that take?", "three days to backfill on a 96-core machine, five minutes to cut over" is the shape of an honest answer.
Practice Drill#
Prompt: "Our orders table has 1.9 billion rows on a single PostgreSQL primary with four replicas, 8,000 writes/s at peak. Order IDs are signed 32-bit integers, the sequence is at 1.9 billion, and at current growth we hit 2^31 in about five months. Twenty services read the table, and two analytics teams query it directly. Plan the migration to 64-bit IDs with no downtime."
Staff Answer
A type change on the primary key is a full table rewrite under ACCESS EXCLUSIVE if done with ALTER COLUMN TYPE — hours of outage on 1.9B rows — so this is an expand/contract migration with the deadline as the main constraint. Step 1, expand: add id_new bigint (nullable, instant), and change the sequence to be bigint so new values can exceed 2^31 later; also add bigint columns in every table with a foreign key to orders. Step 2, dual-write: a trigger or the application sets id_new = id on insert and update — a trigger is acceptable here because it's a single assignment and catches every writer, including the analytics teams' tools. Step 3, backfill: keyset batches of 5K rows by id, WHERE id_new IS NULL, checkpointed, paused when any replica lags over 1s. At ~6K rows/s net of throttling that's ~88 hours — call it a week with pauses for peak traffic. Step 4: build a unique index on id_new with CREATE INDEX CONCURRENTLY (hours, no write blocking), add CHECK (id_new IS NOT NULL) NOT VALID, validate it, then SET NOT NULL. Step 5, verify: full chunked checksum of id = id_new across all rows and in every child table. Step 6, swap: in one short transaction with lock_timeout = 3s and retry, drop the child foreign keys and the old PK constraint, promote the id_new index with ADD PRIMARY KEY USING INDEX, rename columns, point the sequence at the new column, and re-add child foreign keys as NOT VALID (validated afterwards with a light lock) — milliseconds of heavy lock, because every expensive part is already done. Before the swap, every reader must tolerate a bigint ID: I'd inventory clients from pg_stat_statements and connection metadata, and check each service's ORM types, JSON serializers (JavaScript clients lose precision above 2^53, which is still far away, but string IDs in the API are safer) and the analytics teams' schemas. Timeline: two weeks of preparation and client audit, one week of backfill, a week of verification and child tables, swap in week five — leaving four months of margin. Keep id_old for 30 days after the swap, then drop with sign-off from the analytics teams. Metrics: backfill.cursor and ETA, replica lag, migration.lock_timeout count, sequence headroom (max_id / 2^31) alerting at 90%.
Why this is L6:
- Recognizes the PK type change as a table rewrite and decomposes it so every expensive operation runs without a heavy lock, leaving a millisecond swap.
- Sizes the backfill in days against the five-month deadline and throttles on replica lag.
- Audits every reader for 64-bit compatibility before the swap and treats the drop as the one-way door with sign-off.
What L7 adds:
- Notices other tables with 32-bit keys and proposes an org-wide audit with sequence-headroom alerts, so this never arrives with five months' notice again.
- Turns the steps into a reusable runbook in the migration platform, because "widen a primary key" will recur.
- Sets a policy that new tables use 64-bit or UUID keys by default, enforced in the schema lint.
Staff Interview Application#
How to Introduce This Pattern#
"I'll never change this in place. I'll add the new shape, move data and traffic to it in steps that each work with the code on either side, verify on production traffic before users depend on it, and only then remove the old shape. Every DDL runs with a short lock timeout, the backfill is throttled on replica lag, and the drop at the end is the only irreversible step, so it waits a month and needs the readers' sign-off."
Lead with the step sequence, then lock safety, then verification, then the point of no return and who signs off on it.
When NOT to Use This Pattern#
- Small or cold tables: a 50K-row lookup table can take a plain
ALTERin milliseconds. Uselock_timeoutand move on. - Pre-launch or single-tenant dev systems: no live traffic, no compatibility window. Migrate and redeploy.
- Derived data you can rebuild: search indexes, caches, materialized views. Build a new one alongside, swap an alias, delete the old — no dual-write or reconciliation needed. See Elasticsearch.
- True one-time regulated cutovers: where two live sources of truth are worse than a short planned window, run a pure shadow with reconciliation and switch in a declared window.
- When the problem is load, not schema: if the table is slow because of hot keys or write volume, a migration won't help — see Write-Heavy Systems and Hot Keys.
Follow-Up Questions to Anticipate#
| Interviewer Asks | What They Are Testing | How to Respond |
|---|---|---|
| "How do you rename a column with no downtime?" | Expand/contract fluency | "Add new, dual-write, backfill, switch reads, switch writes, soak, drop. Six deploys; each compatible with its neighbours." |
| "Why not just run the ALTER at night?" | Lock awareness | "One long query holding the table turns a 20ms ALTER into minutes of queued traffic. lock_timeout plus retry, any time of day." |
| "How long will the backfill take?" | Back-of-envelope with throttling | "Rows over sustainable rate with lag gating. 2B at 6K/s is about four days; I'd plan for a week." |
| "How do you know the new data is right?" | Verification design | "Shadow reads on live traffic with a zero-mismatch bar for a week, plus chunked checksums of every row before switching reads." |
| "What if you need to roll back after switching writes?" | Reversibility | "The old path is still written until the soak ends, so rollback is a read flip. After we stop writing it, rollback means replaying." |
| "How would you move to a different database entirely?" | Cross-store migration | "CDC from the old store, bulk backfill, dark reads, then per-tenant cutover with a seconds-long write pause each, reverse replication during the soak." |
Scorecard#
| Dimension | Senior (L5) | Staff (L6) | Principal (L7) |
|---|---|---|---|
| Framing | Writes the target schema and a script | Designs the step sequence and compatibility window | Treats migrations as recurring org work with a platform and completion tracking |
| Lock safety | Runs DDL at low traffic | Knows which forms lock or rewrite; lock_timeout, concurrent, NOT VALID | Lint and runner make safe DDL the default |
| Backfill | One big UPDATE | Keyset batches, idempotent, checkpointed, lag-gated | Shared backfill framework with throttle defaults |
| Verification | Spot checks after | Shadow reads and full checksums gate the read switch | Verification is a platform gate, not a team choice |
| Rollback | Backups | Per-step rollback; named point of no return; soak period | Org policy on soak and contract sign-off; parallel-running cost tracked |
Strong Hire Signals
| Signal | What It Sounds Like |
|---|---|
| Sequence over script | "Six steps, each compatible with the code on either side." |
| Knows the lock queue | "lock_timeout 3 seconds, or one report takes the table down." |
| Verifies before switching | "A week of clean shadow reads before any user sees new data." |
| Names the one-way door | "The drop is irreversible, so it waits 30 days and needs sign-off." |
Lean No-Hire Signals
| Signal | Why It Misses the Bar |
|---|---|
| "Rename the column and deploy" | Breaks every pod still on old code during the rollout |
| Unthrottled single-statement backfill | Replica lag, long transactions and a primary near saturation |
| No verification plan | Mismatches are found by users or auditors |
Common False Positives: Knowing gh-ost's flags ≠ knowing when a migration is semantic. Saying "blue-green" ≠ having a data compatibility plan — databases aren't stateless. "We have backups" ≠ a rollback plan.
Capacity Planning Quick Reference#
Sizing the Migration#
backfill_hours = rows / sustainable_rows_per_s / 3600 # 1B at 5K/s ≈ 55 h
sustainable_rate = rate at which replica lag stays < 1 s # measure, don't guess
batch_size ≈ rows per ~200 ms transaction # typically 1K–10K rows
wal_budget = backfill_wal_rate ≤ replica_apply_rate − normal_wal_rate
temp_storage ≈ table_size × 1.0–1.5 # copy tools and new columns
verification_reads ≥ 1M shadow comparisons and ≥ 7 days # before switching reads
soak_period ≥ slowest writer or reader cycle # 30 days if monthly jobs exist
parallel_run_cost = new_system_monthly × months_until_decommission
sequence_headroom = current_max_id / type_max # alert at 0.8, page at 0.9
Key Numbers Worth Memorizing#
| Number | Context |
|---|---|
| 2–5 s | lock_timeout for migration DDL, with retry |
| 1 s | Replica lag at which a backfill should pause |
| 1K–10K rows | Typical backfill batch; aim for ~200 ms per transaction |
| ~55 h | 1B rows at 5K rows/s — backfills are measured in days |
| ≥ 7 days | Clean shadow-read period before switching reads on important data |
| 1–4 weeks | Soak before contract; 30 days when monthly jobs exist |
| 2^31 ≈ 2.1B | Ceiling of a signed 32-bit ID; audit sequences before it bites |
| 3 days / 5 min | Notion's production backfill / switchover maintenance window |
| 4 phases | Stripe: dual-write, move reads, move writes, remove old |
| 1M rows | Rough threshold above which lint should require an online DDL path |
Common Pitfalls Checklist#
- Every DDL runs with
lock_timeoutand retry, never from an app startup hook - Table-rewriting forms on large tables go through expand/contract or a copy tool
- Each step works with the code version before and after it
- All writers of the old shape are inventoried from query logs, not just code search
- Backfill is idempotent, keyset-paginated, checkpointed and lag-gated
- Shadow reads and a full checksum gate the read switch
- Writes move last, and the old path stays written through the soak
- The point of no return is named, and the drop has sign-off from every dependent team
- A snapshot exists before any drop
- Open migrations have an owner, a finish date and a tracked parallel-running cost