Sharding Strategies

Splitting a dataset across machines without regret: key-based, range-based, and directory-based sharding, consistent hashing with virtual nodes, and the true cost of resharding.

Advanced · 22 min read

Why this matters

Your database was fast when it fit on one machine. Then the dataset grew 100×, and no amount of vertical scaling (bigger box, faster disks) keeps up — physics and budgets both tap out. Sharding (horizontal partitioning) splits the data across many machines so each one holds a fraction. It is the only strategy that scales writes indefinitely, and it is also the strategy most likely to ruin your quarter if you botch the shard key. Choose well and shards are invisible; choose badly and one shard melts while the others idle.

The analogy is the same library you split across branch buildings back in the partitioning lesson — same three layouts: shelve by author's last initial (key-based), by acquisition-date ranges per building (range-based), or keep a master catalog saying which building holds what (directory-based). This lesson doesn't re-teach those layouts; it goes one level down. Here you get the resize math that breaks naive key sharding, the mechanics of consistent hashing with virtual nodes, the true cost of resharding, and how to pick a shard key when the constraints fight each other.

Video: SD 11 | Database Sharding | System Design Course | Uplatz — Uplatz
Uplatz walks through why one database hits a wall, what sharding is, and when to reach for it.

Strategy 1: Key-based (hash) sharding

The key-based layout, from the partitioning lesson: shard = hash(key) mod N. Deterministic, uniform-ish distribution, trivially computable. No lookup needed. Cassandra, DynamoDB, and Redis Cluster all hash the key.

flowchart LR
    K["key: user_8472"] --> H["hash → 0x3f…"]
    H --> M["mod 8 → shard 3"]
    M --> S3["Shard 3"]

The catch is the mod N: adding a shard renumbers almost everything. With naive modulo hashing, going from 8 to 9 shards moves ~8/9 of your data — a full reshuffle for one shard of new capacity. This is the problem consistent hashing was invented to solve (see below). Key-based sharding also scatters related data: a user's orders land on different shards than the user, so joins become application-level scatter-gather.

Video: Partitioning & Sharding - Redis for Developers — The Software Mentor
Redis Cluster's 16384 hash slots show how hashing keys spreads rows evenly, plus hash tags to co-locate related keys.

Strategy 2: Range-based sharding

The range layout, from the partitioning lesson: contiguous key ranges per shard — shard 1 holds A–F, shard 2 holds G–M, and so on. HBase, Bigtable, CockroachDB, and MongoDB (range sharding) work this way. Range queries are beautiful — "all orders in March" hits one or two shards instead of all of them.

The failure mode has a name: the hotspot. Time-series data, auto-increment IDs, or anything monotonically increasing piles all writes onto the last shard — one machine does 100% of the writes while the others nap. The standard fix is a salted prefix: prepend hash(key) mod K to the key so sequential writes spray across K shards, at the cost of making range scans touch K shards. There is no free lunch; there is only lunch whose price you've itemized.

Video: What Is a Partition Key? Routing, Pruning, Hot Partitions, and How to Choose One — Kandi Brian
Walks through range-based sharding, the shard-by-date hotspot, and spreading fixes.

Strategy 3: Directory-based sharding

The card-catalog layout, from the partitioning lesson: keep a lookup service mapping each key (or key range) to its shard. Maximum flexibility: move any key anywhere, split hot shards surgically — at the cost of the directory itself: it's a critical-path lookup on every request (cache it aggressively) and a new thing to keep highly available (as the partitioning lesson says, the catalog must itself be replicated — replica set, quorum — never a lone process). Useful when your sharding needs are irregular enough that no formula works, e.g., multi-tenant SaaS where one whale customer deserves their own shard.

Video: What is Database Sharding? — Anton Putra
A phone-book ledger maps each key to its shard — Anton covers directory sharding's flexibility and its single-point-of-failure price tag.

Consistent hashing: adding shards without the apocalypse

You met the ring idea in the partitioning lesson (and saw requests routed on one back in load balancers) — here are the mechanics that make it production-grade. The elegant fix for the mod N reshuffle: hash both keys and shards onto a ring; each key belongs to the first shard clockwise from it.

flowchart TD
    R["Hash ring 0 … 2³²"]
    R --> K1["key A → shard 1<br/>(first clockwise)"]
    R --> K2["key B → shard 3"]
    R --> ADD["Add shard 4 between 2 and 3:<br/>only keys mapping to that<br/>arc move — ≈ 1/N of data"]

The math that matters: adding one shard to N moves only ~1/N of keys (vs. ~(N-1)/N with naive modulo). With 10 shards, a scale-up moves ~10% of data, not ~90%. But plain consistent hashing distributes unevenly when N is small — so production systems use virtual nodes: each physical shard claims, say, 100–200 points on the ring, smoothing distribution and making heterogeneous hardware natural (big machine gets more virtual nodes). Dynamo, Cassandra, and Riak all do this.

Interactive diagram: StepThrough (loads in the app)

Video: Consistent Hashing | The Backend Engineering Show — Hussein Nasser
Builds the hash ring from the modulo-hashing failure, showing why adding a server moves only ~1/N keys; covers virtual nodes.

The numbers: resharding cost and shard sizing

Some back-of-envelope arithmetic before you shard:

Video: Database Sharding Explained: Partitioning, Hot Keys, and Safe Resharding — DevLevel13
Walks through resharding without an outage via dual-write and shadow reads, plus how hot keys drive migration cost.

Choosing the shard key (the decision you'll live with)

The partitioning lesson taught you what makes a key good — matches your access pattern, high cardinality, no monotonic hot keys. This section is about choosing when the constraints fight each other. The shard key determines data distribution, query patterns, and your future resharding pain. Good keys have high cardinality (many distinct values), even distribution, and colocate related data (a user's rows on one shard = single-shard transactions). Bad keys: timestamps (hotspot), booleans (two shards, one idle), tenant ID when one tenant is 80% of traffic (the whale problem — give the whale its own shard via directory-based routing).

Changing the shard key later is a full data migration. The most expensive schema change there is. Spend the design meeting. Future you, paged at 3 AM because shard 7 of 8 is at 100% CPU, will be grateful.

Video: Choosing a Shard Key — InterSystems Developers
Covers shard key criteria — cardinality, even distribution, query alignment — and when the default key is good enough.

Takeaways

  1. Key-based = simple and uniform but reshuffles on scale-up; range-based = great scans but hotspot-prone; directory-based = flexible but adds a critical lookup.
  2. Consistent hashing with virtual nodes makes adding a shard move ~1/N of data instead of nearly all of it.
  3. Budget the real costs: shard move time (GB ÷ network), connection explosion (S×C×P), and the hot-key ceiling no sharding fixes.
  4. The shard key is the hardest decision — high cardinality, even distribution, colocated related data — because changing it is a full migration.

Check your understanding

  1. What is the main problem with naive hash(key) mod N sharding when adding a shard?

    • The hash function gets slower as N grows
    • Keys can no longer be looked up without a directory
    • Almost all data must move because the modulus changed
    • Range queries become impossible forever
  2. What failure mode does range-based sharding suffer with monotonically increasing keys?

    • The hotspot: all writes pile onto the last shard
    • Data lands too evenly to be cached effectively
    • Keys hash to the wrong shard after every rebalance
    • The directory service becomes the bottleneck
  3. With consistent hashing, adding one shard to N approximately moves what fraction of keys?

    • Exactly half of them
    • About 1/N of them
    • Nearly all of them
    • None — keys stay exactly where they are
  4. Why do production consistent-hashing systems use virtual nodes?

    • To encrypt the keys sitting on the ring
    • To reduce the number of hash computations per lookup
    • To eliminate the need for replication entirely
    • To smooth uneven distribution with few shards and support heterogeneous hardware

Go deeper

Want to keep pulling this thread? These talks and tutorials go further than we did here:

Sources & further reading