wtf( )unctionsystem design, drawn
← all problemsDatabasesHard

Picking the column the whole system hangs on

One database can no longer hold the orders table, so it's being split across four shards. Every read and write will be routed by one column, and changing that choice later means migrating everything.

A good shard key spreads data evenly and keeps the rows a single query needs on one shard. A bad one gives you four databases with the traffic of one.

Choose the column to shard on.
Components — tap one, then tap a slot on the diagram
?Every read and write will be routed by this column. Changing it later means migrating every row.

Boundaries, outermost first: Shards: Shard 3, Shard 1, Shard 2 Outside every boundary: Order service, an empty slot for the route on this column Connections: Order service calls route on this column — route by key (step 1) route on this column calls Shard 1 (step 2) route on this column calls Shard 2 route on this column calls Shard 3

Shard 3
Shard 1
Shard 2
Order service