Skip to content

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
12CLIENTUsersSERVICEApplicationserversDATABASEMonolithPostgres

Select a component to see what it is responsible for and which state it owns.

  1. 1Users → Application servers: Requests
  2. 2Application servers → Monolith Postgres: Reads and writes until cutover

What you need to know

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

  2. Check

    Why choose 480 logical shards rather than 500?