Shard a Live Database Without Downtime, stage 4 of 9: change it
Move billions of rows, live
Notion captured writes with an application-level audit log rather than Postgres logical replication, because replicating the initial snapshot of tables this large through logical replication was too slow for them at the time.
System so far· 5 parts
Select a component to see what it is responsible for and which state it owns.
- 1Users → Application servers: Requests
- 2Application servers → Monolith Postgres: Reads and writes until cutover
- 3Application servers → Connection poolers: Queries by shard
- 4Connection poolers → Shard hosts: Schema on host
What you need to know
0 of 1 checks done
Moving live data has one shape, whatever the tools. See Online data migrations:
- Capture every new write (here, an audit log).
- Backfill existing rows.
- Catch up by replaying captured writes until the copy trails by seconds.
- Verify the copy against the original.
- Switch reads and writes, briefly.
- Keep the old store as a fallback for a while.
Check
Why start capturing writes before the backfill begins?