Shard a Live Database Without Downtime, stage 7 of 9: change it
From 32 hosts to 96
Each of the 32 hosts holds 15 of the 480 logical shards.
System so far· 7 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 → Write audit log: Log every write
- 4Backfill and catch-up → Monolith Postgres: Copy existing rows
- 5Backfill and catch-up → Write audit log: Replay logged writes
- 6Backfill and catch-up → Shard hosts: Write rows by shard
- 7Application servers → Connection poolers: Queries by shard
- 8Connection poolers → Shard hosts: Schema on host
- Request / response
- Asynchronous
What you need to know
0 of 2 checks done
With fixed logical shards, adding hosts means moving whole schemas. Postgres logical replication can copy a schema to a new host and keep it in sync. Cutover is a brief pause at the connection poolers while replication drains, then a routing change.
No workspace changes shard, so no row-level migration code is needed.
Check
Why not change the hash to spread workspaces over 96 shards?Work it out
Each pooler held 50 connections per host to 32 hosts. With 96 hosts, how many connections per pooler?