Database Sharding

Sharding splits one big table across several databases (shards) by a shard key, so data and load are spread over many machines.

Hash of user id

Shard 1

(id × 7) % 3 = 0 · 0 writes

    Shard 2

    (id × 7) % 3 = 1 · 0 writes

      Shard 3

      (id × 7) % 3 = 2 · 0 writes

        Each user goes to shard (id × 7) % 3 — a simple hash that scatters neighbouring ids.

        Step 1 / 16
        Users
        0
        Shards queried
        0

        Step by step

        The example above, written out — the same steps the animation plays.

        1. 1Hash of user id. Each user goes to shard (id × 7) % 3 — a simple hash that scatters neighbouring ids.
        2. 2Insert Asha (id 1) → shard 2. (1 × 7) % 3 = 1.
        3. 3Insert Ravi (id 2) → shard 3. (2 × 7) % 3 = 2.
        4. 4Insert Neha (id 3) → shard 1. (3 × 7) % 3 = 0.
        5. 5Insert Kiran (id 4) → shard 2. (4 × 7) % 3 = 1.
        6. 6Insert Meera (id 5) → shard 3. (5 × 7) % 3 = 2.
        7. 7Insert Omar (id 6) → shard 1. (6 × 7) % 3 = 0.
        8. 8Insert Priya (id 7) → shard 2. (7 × 7) % 3 = 1.
        9. 9Insert Sana (id 8) → shard 3. (8 × 7) % 3 = 2.
        10. 10Insert Tara (id 9) → shard 1. (9 × 7) % 3 = 0.
        11. 11Insert Vikram (id 10) → shard 2. (10 × 7) % 3 = 1.
        12. 12Insert Gita (id 11) → shard 3. (11 × 7) % 3 = 2.
        13. 13Insert Farhan (id 12) → shard 1. (12 × 7) % 3 = 0.
        14. 14SELECT … WHERE id = 7. The shard key tells us exactly where id 7 lives: shard 2. One shard, one query.

        …and 2 more steps — press Play above to watch them all.

        What's happening?

        1. Each row goes to a shard chosen from its shard key — by hashing it, or by ranges of it.
        2. A query that includes the shard key goes to exactly one shard.
        3. A query on anything else must ask every shard and merge the answers (scatter-gather).

        Where you'll meet it

        Very large user, order or message tables — at the scale where one database can no longer hold the data or the write load (Instagram, Discord, Uber).

        Common mistake

        Range-sharding on an ever-increasing key such as an auto-increment id or a timestamp: every new write lands on the newest shard — a hot spot.

        FAQ

        Hash or range sharding?

        Hash spreads writes evenly but makes range queries touch every shard; range keeps related rows together but can create hot spots.

        How do you choose a shard key?

        From the most common query — usually something like user_id — so most queries hit one shard.

        What happens when you add a shard?

        Data must move. Consistent hashing keeps the amount that moves small.