Shard a Live Database Without Downtime, stage 3 of 9: decide
How many shards?
Today 32 hosts are enough. In a few years it may be 96, or more. Moving rows between shards is a migration in its own right; moving a whole shard from one host to another is much easier.
System so far· 3 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
What you need to know
0 of 2 checks done
Separate two layouts:
- Logical: which shard a row belongs to. Fixed forever, computed from the workspace ID.
- Physical: which host a shard lives on. Changed whenever you like.
With many small logical shards (say 480), growing from 32 to 96 hosts means moving whole shards between hosts, not re-hashing rows.
Check
Why choose 480 logical shards rather than 500?Work it out
480 logical shards on 96 hosts. How many shards per host?