Online data migrations
Moving live data to a new schema or store without downtime: write to both, backfill the past, verify, switch reads, then switch writes, with a way back at every step.
Storage & state
Learn it
Data outgrows its home: a table needs a new shape, a database needs sharding, a store needs replacing. The system can't stop while billions of rows move, and rows keep changing during the copy. "Copy, then switch" loses every write made during the copy, and a big-bang cutover has no way back.
The safe pattern moves in reversible steps:
- Dual write: every new write goes to the old store (still authoritative) and the new one, via application code or, more reliably, a change log.
- Backfill history, without overwriting newer values, throttled so production isn't starved.
- Verify: sample rows, compare counts by range, and run dark reads that query both and alert on mismatches.
- Switch reads gradually, keeping dual writes.
- Switch writes so the new store becomes the source of truth.
- Clean up once nothing reads the old store.
Check
Why switch reads before writes?Think first
The backfill upserts every historical row unconditionally while dual writes are running. What goes wrong?
Quick reference
The same ideas, condensed for revision.
How it goes wrong
- Lost writes during the copy
- Writes that arrive between snapshot and switch never reach the new store.
- Backfill clobbers newer data
- An old snapshot value overwrites a newer dual-written one.
- Partial dual writes
- One store's write fails and the other's succeeds, so they diverge silently.
- No way back
- Writes moved to the new store with nothing keeping the old one current.
Instead, consider
- Maintenance window
- Brief downtime is acceptable and the dataset copies within it.
- Expand/contract schema changes
- The change is to columns in one database: add the new shape, migrate code, remove the old.
- Leave old data in place
- Only new data needs the new store; old data can be read from the old one until it ages out.
In practice
- gh-ost / pt-online-schema-change
- Shadow table plus change capture, then an atomic table swap (MySQL).
- Postgres logical replication
- Stream changes to a new cluster, then fail over.
- Scientist-style experiments
- Run old and new read paths side by side and report differences.
It assumes
- Writes can be captured completely (in code or from a log).
- Rows have a version or timestamp so the backfill can avoid clobbering newer data.
- The data can be compared between stores cheaply enough to verify.
Explain it in your own words
Where you practise it
Further reading
Engineers describing it in systems they run.
- The Great Re-shard: adding Postgres capacity (again) with zero downtime
Notion · Arka Ganguli and others · Post, Jul 2023
The payoff of fixed logical shards: growing to more machines by moving shards whole.
- How Discord Stores Trillions of Messages
Discord · Bo Ingram · Post, Mar 2023
The same data model six years on: hot partitions, a service layer that merges identical reads, and a migration of the full history to a new database.
- How we turned ClickHouse into our event mansion
PostHog · James Greenhill · Post, Nov 2021
Why event analytics outgrew Postgres, how they chose a column store, and the mistakes they made running it.
- Herding elephants: lessons learned from sharding Postgres at Notion
Notion · Garrett Fidalgo · Post, Oct 2021
Picking a shard key and a shard count, then moving a live database across with no data lost. Includes an honest list of what they would do differently.
Related concepts
- Replication
Keeping copies of data on several machines for durability, read capacity and locality, and living with copies that briefly disagree.
- Append-only logs
Recording changes as an ordered, immutable sequence of facts, from which current state, history and replicas can be derived.
- Partitioning
Splitting data or work by key so each part is handled independently: scaling out, and giving each key a single owner.
- Reconciliation
Periodically comparing your records with an authoritative source and repairing differences, the backstop for every message that was lost.