learninfra · Linux · Networking · Kubernetes · System Design · AI Infrastructure · Exam blueprints · Drills

Sharding

Split a database that no longer fits on one machine: shard keys, routers, modulo versus range versus consistent hashing, and the queries and hot keys that sharding cannot fix.

An interactive System Design lesson: 18 steps, about 30 minutes, on a live simulation in your browser.

A messaging app keeps one row per user, and every message sent updates the sender's row. The single Postgres leader is at its limit: the table no longer fits on its disks, the rows people touch no longer fit in its memory, and it cannot take any more writes per second. Read replicas do not help. They multiply reads, and every replica still applies every write (see "Database Replication").

So the rows are split across four databases, shard-a to shard-d. Each shard holds some of the users and nothing else. Send 200 writes spread over the 40 users: shard-a takes 53, shard-b 17, shard-c 60, shard-d 70. No machine sees all the writes, and adding a machine adds write capacity.

What you will learn

  1. When one leader is not enough

    • Too many writes for one machine: Replication copies all the data to more machines; sharding splits the data so each machine holds and writes only its part.
    • The shard key and the router: Every query that names the shard key goes to exactly one shard; the router turns the key into an address.
  2. Modulo and the cost of growing

    • Add a fifth shard: With hash mod N, changing N moves almost every key, not just the new shard's share.
    • Drill: how many keys move?
  3. Range partitioning

    • Split by key range instead: Range partitioning keeps neighbouring keys together, so range scans stay on one shard, and so does any burst of neighbouring writes.
    • New sign-ups: Sequential keys under range partitioning send every new write to the last shard.
    • Split the hot range
  4. Consistent hashing

    • Put the shards on a ring
    • A sixth shard on the ring: Consistent hashing moves about 1/N of the keys when a shard joins or leaves, and only to or from that shard.
  5. Choosing the shard key

    • What makes a good shard key: A shard key needs many values, evenly used, and it must appear in the queries you run most.
    • A query without the shard key: A query without the shard key goes to every shard and is as slow as the slowest one; more shards make that tail worse, not better.
    • One key nobody can split: Sharding spreads keys, not traffic. A single hot key is served by a single shard, however many you have.
    • Drill: pick the shard key
  6. What sharding costs

    • Two rows, two shards: Anything that touches two shards (a transaction, a join, a unique constraint) is a distributed-systems problem, and you pay for it in latency and complexity.
    • One shard goes down: A dead shard takes exactly its own slice of the keys with it; everything else keeps working.
    • Where you will meet this
  7. Recap & playground

    • Cheat sheet
    • Playground