Hiring BarSupport

Online Schema & Data Migrations

Pattern44 min read4 diagrams

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 TABLE forms need an ACCESS EXCLUSIVE lock; while one waits behind a 2-minute analytics query, every new read on that table queues behind it. Set lock_timeout to 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.

IntentConstraintStrategyFailure ModeCorrectness Bar
Schema change on a live table (add column, index, constraint, type change)Minutes to hours; same database; no downtimeNative online DDL where it exists; lock timeouts; concurrent index builds; NOT VALID then VALIDATE; copy tools for rewritesLock queue blocks all traffic; table rewrite saturates I/OZero 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 changeExpand → dual-write → backfill → shadow-read → switch reads → switch writes → contractMissed writer path; semantic mismatch found after switchMismatch 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 parallelDual-write or CDC replication, bulk backfill, dark reads, per-tenant or per-shard cutoverDivergence during the long parallel run; cutover lag; cost of two systemsPer-entity checksums match; cutover reversible per slice
Rebalancing / resharding (move ranges or tenants between nodes)Continuous or periodic; same storeCopy + change-stream catch-up + brief write pause + directory flipWrites lost in the flip window; hot destinationNo 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#

StrategyWhat WorksWhat BreaksWho Pays
Maintenance window (stop writes, migrate, restart)Simplest; no compatibility window; easy to verifyDowntime proportional to data size; a 2B-row rewrite doesn't fit in any window; rollback means another windowUsers and the business during the outage; on-call doing it at 3am
In-place DDL (plain ALTER TABLE)One statement; instant for metadata-only changesRewrites and exclusive locks on large tables; lock queue blocks reads; not pausableEvery 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 renameDoubles storage temporarily; trigger-based variants add write overhead; cut-over still needs a brief lockDBA or platform team; disk headroom
Expand / contract in the applicationEvery step small, independent and reversible; works on any storeMany deploys over weeks; code carries both paths; easy to forget the contract stepThe owning team's velocity; the codebase until cleanup
Dual write + backfill + shadow readVerifies correctness on real traffic before users depend on itDual write is itself a consistency risk; every writer must be found; doubles write loadThe owning team; the new store's capacity
CDC replication + cutoverCaptures every writer, including ones you didn't know about; no app dual-writeLag at cutover; transformation logic lives in the pipeline; needs log accessPlatform team running CDC; cutover coordination
Big-bang copy and switchFast to execute; one cutoverNo verification on live traffic; rollback loses writes made after the switchUsers 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 ALTER or 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 INSTANT column add in MySQL 8.0, doesn't rewrite the table. Still set lock_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#

BehaviorSenior (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?"
LocksRuns ALTER TABLE in a deploy hookKnows which DDL forms rewrite or lock; sets lock_timeout; uses concurrent and NOT VALID forms; copy tools for rewritesLints migrations in CI for dangerous forms; schema changes go through a platform with automatic throttling
BackfillOne UPDATE ... SET new = f(old)Resumable, primary-key-ordered batches throttled on replica lag, checkpointed, idempotentShared backfill framework with lag-aware throttling and progress reporting every team uses
VerificationSpot-checks rows after the migrationShadow reads and full checksums before switching reads; mismatch budget of zero, investigatedMakes 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 returnSets org policy on soak periods and who signs off on one-way steps
OwnershipThe team that wrote the scriptOwning team runs it; downstream readers listed and signed off; DBA reviews locksMigration 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 LineThe Tension
1Database-Level Online DDL vs Application-Level Expand/ContractOne tool-driven operation vs many small deploys that change semantics safely
2Application Dual-Write vs Log-Based ReplicationExplicit control and same-transaction writes vs catching every writer without touching code
3Big Switch vs Gradual CutoverOne coordinated moment vs percentage or per-tenant ramps with two paths live
4Sampled Verification vs Full ReconciliationCheap continuous confidence vs exhaustive proof that costs a full scan
5Keep the Old Path Writable vs Cut Over One-WayFree 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 SayWhat Interviewers HearWhat 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#

Diagram: 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 INDEX without CONCURRENTLY, ADD CONSTRAINT without NOT VALID on tables over 1M rows, any DDL without lock_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.

StepSchemaWritesReadsRollback
1. ExpandAdd amount_minor, currency (nullable)Old onlyOldDrop new columns
2. Dual-writeSameOld + new, same transactionOldDeploy previous version
3. BackfillSameOld + newOld; shadow-read new and compareStop job; new columns ignored
4. Switch readsSameOld + newNew (ramp 1% → 100% by flag)Flip flag back
5. Switch writesConstraints validated on newNew + old (kept for rollback)NewFlip reads back; old is current
6. ContractDrop amount_text after 2–4 week soakNewNewNone — 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.

Diagram: 4. Datastore Migration — CDC, Dark Reads and Per-Slice Cutover
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

TechniqueBlocks Traffic?PausableVerifies SemanticsRollbackBest For
Plain DDLYes, while locked or rewritingNoNoAnother DDLSmall tables, metadata-only changes
Concurrent / NOT VALID DDLNo (brief lock only)No, but restartableNoDropIndexes and constraints on big tables
Online copy tool (gh-ost style)Brief cut-over lockYes, throttles on lagNoKeep old table until droppedTable rewrites, type changes in MySQL
Expand / contract + dual-writeNoYesYes, via shadow readsFree until contractModel refactors in one database
CDC replication + per-slice cutoverSeconds per sliceYesYes, via dark readsReverse replication during soakDatastore migrations, resharding
Maintenance windowYes, entirelyN/AOffline onlyRestoreTiny systems, regulated one-time cutovers

Architecture Diagram#

Diagram: 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#

FailureDetection SignalBlast RadiusMitigationOwner
DDL lock queuepg_locks waiters > 50 behind a DDL; lock wait p99Every query on the tableCancel DDL; lock_timeout + retryOwning team + DB platform
Table rewrite under lockDDL duration; I/O saturationTable offline for the rewriteExpand/contract or online copy toolOwning team
Backfill replica lagreplica.lag_seconds > 1sStale reads; primary overloadLag-gated throttle, smaller batchesOwning team
Missed writerWriter inventory gaps; reconciliation mismatchesSilent divergenceAudit-log writer inventory; CDC instead of dual-writeOwning team
Shadow-read mismatchmigration.shadow_mismatch > 0None yet — that's the pointFix transform; re-backfill affected rangeOwning team
Cutover lagCDC lag at pause > 60sWrites paused longer than budgetAbort and unpause; retry slice laterPlatform + owning team
Contract too earlyDownstream job failures after dropData loss for unmigrated writersSnapshot before drop; soak ≥ slowest cycleOwning team + downstream sign-off
Mixed-version incompatibilityError spike during rolling deployHalf the fleet during the deployEvery step compatible with adjacent versionsOwning 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.

OptionWhat WorksWhat BreaksWho Pays
Each team runs DDL and backfills its own wayAutonomy; no waiting on a platformSame incidents repeat across teams; no lock safety by default; half-finished migrations accumulateOn-call for whichever team forgets lock_timeout; the business paying for parallel systems
Central DBA team approves and runs every changeExpertise concentrated; consistent safetyBottleneck — teams wait days for a column; DBAs lack domain context for semantic changesProduct velocity; DBA burnout
Platform provides tooling and guardrails; teams own semanticsLint, lock-safe runner, backfill framework, verifier and flags are self-service; DBAs review only flagged-risky changesPlatform must support every store in use; tooling lags new enginesPlatform 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.

ScaleMigration LoadMachineryInfra Cost of Parallel RunningPeopleRough Cost
Startup~5 schema changes/month, small tablesMigration files in CI, lock_timeout defaultsNegligiblePart of feature work~$0 extra; ~1 day per risky change
Growth~40 changes/month; 1–2 model refactors per quarterDDL lint, lock-safe runner, backfill framework, shadow-read helperTemporary 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 yearSelf-service platform, per-slice cutover tooling, verifier service, migration reviewParallel clusters during datastore moves: $50–300K/month each for 3–9 months4-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#

Diagram: 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#

DecisionDoor TypeReversibility Cost
Adding a nullable column, an index, a NOT VALID constraintTwo-wayDrop it
Switching reads to the new pathTwo-wayFlag flip, as long as the old path is still written
Switching writes to the new path (old still written)Two-wayFlip reads back; old is current
Stopping writes to the old pathOne-way-ishRollback now requires replaying writes since the stop
Dropping the old column, table or datastoreOne-waySnapshot restore plus manual reconciliation, if a snapshot exists
Changing a primary key or shard keyOne-way-ishA second full migration in the other direction
Publishing the new schema to other teams (events, APIs, CDC consumers)One-way-ishTheir 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:

  1. 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.
  2. 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).
  3. Backfills MUST be idempotent, resumable from a checkpoint, and throttled on replica lag (default pause above 1s).
  4. Semantic migrations MUST pass shadow-read verification (zero unexplained mismatches over ≥ 7 days) and a full checksum before reads switch.
  5. 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#

SignalWhat 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 ALTER in milliseconds. Use lock_timeout and 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 AsksWhat They Are TestingHow 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#

DimensionSenior (L5)Staff (L6)Principal (L7)
FramingWrites the target schema and a scriptDesigns the step sequence and compatibility windowTreats migrations as recurring org work with a platform and completion tracking
Lock safetyRuns DDL at low trafficKnows which forms lock or rewrite; lock_timeout, concurrent, NOT VALIDLint and runner make safe DDL the default
BackfillOne big UPDATEKeyset batches, idempotent, checkpointed, lag-gatedShared backfill framework with throttle defaults
VerificationSpot checks afterShadow reads and full checksums gate the read switchVerification is a platform gate, not a team choice
RollbackBackupsPer-step rollback; named point of no return; soak periodOrg policy on soak and contract sign-off; parallel-running cost tracked

Strong Hire Signals

SignalWhat 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

SignalWhy It Misses the Bar
"Rename the column and deploy"Breaks every pod still on old code during the rollout
Unthrottled single-statement backfillReplica lag, long transactions and a primary near saturation
No verification planMismatches 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#

NumberContext
2–5 slock_timeout for migration DDL, with retry
1 sReplica lag at which a backfill should pause
1K–10K rowsTypical backfill batch; aim for ~200 ms per transaction
~55 h1B rows at 5K rows/s — backfills are measured in days
≥ 7 daysClean shadow-read period before switching reads on important data
1–4 weeksSoak before contract; 30 days when monthly jobs exist
2^31 ≈ 2.1BCeiling of a signed 32-bit ID; audit sequences before it bites
3 days / 5 minNotion's production backfill / switchover maintenance window
4 phasesStripe: dual-write, move reads, move writes, remove old
1M rowsRough threshold above which lint should require an online DDL path

Common Pitfalls Checklist#

  • Every DDL runs with lock_timeout and 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
  1. Loading the index…