Skip to content

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

0 of 2 checks done
  1. 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.

  2. The safe pattern moves in reversible steps:

    1. 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.
    2. Backfill history, without overwriting newer values, throttled so production isn't starved.
    3. Verify: sample rows, compare counts by range, and run dark reads that query both and alert on mismatches.
    4. Switch reads gradually, keeping dual writes.
    5. Switch writes so the new store becomes the source of truth.
    6. Clean up once nothing reads the old store.
  3. Check

    Why switch reads before writes?

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

Write at least 60 characters (0 so far). Write it as you would say it in a design review. You will compare it against the points a strong answer makes.

Where you practise it

Further reading

Engineers describing it in systems they run.

  • 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.