Start here: the problem it solves#
TL;DRthe 30-second version
- Sharding splits one dataset across many machines. A shard is one piece; the partitioning strategy is the rule that decides which shard each row lives on.
- Hash partitioning sends hash(key) mod N to a shard. It spreads load evenly, but a resize changes N and moves almost every key.
- Range partitioning gives each shard a contiguous band of keys. Range scans are fast, but sequential keys (auto-increment ids, timestamps) pile into one shard.
- Directory partitioning keeps a lookup table of key to shard. Total control, at the cost of a table that grows with the data and one extra lookup per request.
- The hidden cost is resharding. Plain hash-mod moves about N/(N+1) of the data when you add a shard. A bucketed or consistent-hashing scheme moves only about 1/(N+1).
Say you're running a user database. It started on one server, and every query just ran. Then you grew. Now you have two billion users, each row is about 1 KB, and that's 2 TB before indexes. No single machine has the memory to cache it or the disk throughput to serve it.
The first instinct is to buy a bigger machine. That's scaling up, and it works right until it doesn't: the biggest server money can buy still has a ceiling. The alternative is to scale out and spread the data across many ordinary machines. That's sharding.
The obvious split is by first letter: usernames A–F on shard 0, G–M on shard 1, and so on. Now watch it break. Far more names start with 'S' or 'M' than with 'X' or 'Z', so the busiest shard ends up holding about three times what the quietest one does. It's three times as busy too, and it becomes the bottleneck for the whole system. The split used a property of the key that isn't evenly distributed. So the real question is the function that maps a key to a shard: it has to find any key without searching, and keep the pieces roughly equal.
The mechanism: a shard key and a routing rule#
Every sharding scheme has two parts. The shard key is the field you route on: the user id, the account id, whatever you look things up by. The routing rule turns that shard key into a shard number. Pick the shard key once, up front. It decides where every row lives, so changing it later means moving everything. Then choose one of three routing rules.
The first rule is hash partitioning. Hash the shard key to a number and take it mod N, where N is the shard count. The hash spreads keys evenly, so each shard gets about the same share. Neighbouring keys like user:1000 and user:1001 land on unrelated shards. This is the default when you want even load and you look rows up one at a time.
The second rule is range partitioning. Each shard owns a contiguous band of the key space: shard 0 owns keys 0 through 1999, shard 1 owns 2000 through 3999, and so on. The payoff is range scans. 'Every order from last Tuesday' touches one or two shards instead of all of them, because nearby keys live together. The risk is the flip side of that same property.
The third rule is directory partitioning. Keep an explicit table that records, for every key, which shard it's on. You can put any key anywhere and move a single hot key by itself. But the table has one row per key, so it grows with the data, and every request pays for the extra lookup.
PredictYou shard an events table on a timestamp using range partitioning. New events always have the current time. Which shard gets all the writes?
Hint: Where does 'now' always fall in a range that runs oldest to newest?
The last one, the shard that owns the most recent time range. Timestamps only ever increase, so every new write lands at the top of the key space. One shard takes the entire write load while the rest only serve old reads. This is the monotonic-key hotspot, and it's the classic reason not to range-partition on a timestamp or an auto-increment id.
When one shard runs hot#
A hotspot is a shard that carries far more than its fair share of data or traffic. Range partitioning creates one whenever the shard key climbs over time. Auto-increment ids and timestamps always fall in the top band, so the shard that owns that band takes every new write. You bought four machines and you're writing to one.
Hash partitioning scatters sequential keys, but it has its own hotspot: a single popular key. One celebrity account or one enormous tenant gets a thousand times the traffic of everyone else, and hashing sends all of it to whichever shard owns that key. Adding shards doesn't help, because the key is indivisible. The fix lives outside the sharding scheme: replicate the hot key to several shards, cache it in front of the database, or split the key itself into sub-keys.
Range partitioning has a graceful answer to its own hotspots. When a range gets hot, split it. Cut the busy shard's band in two at the median of its keys and hand the upper half to a new shard. Only that shard's keys move. This is how HBase splits regions and how DynamoDB's 'split for heat' relieves a hot partition: it picks the split point from where the traffic actually is.
The cost nobody plans for: adding a shard#
Balance is only half the story. The other half shows up the day you add a shard. This is resharding, and its cost is one number: how many keys have to physically move to a new machine. Every key that moves is a network copy, a cache invalidation, and a window where a request might look in the wrong place.
Plain hash-mod is expensive to grow. Routing is hash(key) mod N. Go from 4 shards to 5 and the divisor changes, so almost every key's answer changes. A key stays put only if hash(key) mod 4 and hash(key) mod 5 happen to agree, which is rare. About four out of every five keys move. In general, a resize from N to N+1 keeps about 1/(N+1) of the keys and moves about N/(N+1), so the bigger the cluster, the worse a single resize gets. Do that on a live 2 TB cluster and you're copying 1.6 TB while still serving traffic.
The fix is to stop hashing directly onto the shard count. Hash onto a large, fixed number of virtual buckets, say a few thousand, and keep a small map from buckets to shards. Routing becomes: hash to a bucket, then look up the bucket's shard. The shard count never appears in the hash, so growing the cluster doesn't rehash anything. To add a shard, reassign a handful of buckets to it and move only their keys: about one in N+1 of the data, roughly 20% going from 4 to 5. This is the fixed-partition idea from 'Designing Data-Intensive Applications'. Consistent hashing reaches the same bound with a ring instead of a bucket map; see that topic for the ring in full.
| Strategy | Keys moved when 4 → 5 shards | Why |
|---|---|---|
| Hash-mod (hash mod N) | ~80% (four in five) | The divisor N is baked into the hash, so changing N re-routes almost everything. |
| Bucketed / consistent hashing | ~20% (one in five) | Keys hash to fixed buckets; only the reassigned buckets' keys move. |
| Range (split a hot shard) | Only the split shard's keys | A targeted split touches one shard; the rest are untouched. |
| Directory | Only what you choose to move | The table forces no movement; you migrate keys deliberately. |
Choosing a strategy#
| Strategy | Load balance | Range scans | Resize cost | Best when |
|---|---|---|---|---|
| Hash | Even | Bad (hits all shards) | High unless bucketed | Point lookups by key, even load matters most |
| Range | Risky (hotspots) | Excellent | Cheap (split one shard) | Time-series, range queries, ordered scans |
| Directory | Whatever you enforce | Depends on layout | Whatever you choose | Few large tenants, need per-key control |
Most real systems combine these. A common pattern is a compound shard key: hash a high-cardinality prefix so load spreads, then keep a range component within each shard so scans inside a tenant still work. Another is to range-partition but pre-split the ranges and add randomness to the key so monotonic writes don't all land in one place.
The one decision that outlives all the others is the shard key itself. Choose a key with high cardinality (many distinct values, so load can spread), low skew (no single value dominates the traffic), and alignment with your most common query (so that query hits one shard, not all of them). Get those three right and the routing rule is a detail. Get the key wrong and no routing rule saves you.
The trades you're actually making#
- Hash trades range scans for even load. 'Every order from Tuesday' is spread across every shard, so the database has to ask all of them and merge the answers.
- Range trades safety for scans. An ordered read touches one or two shards, but a climbing key sends every new write to one shard.
- Directory trades simplicity for control. Every request pays one extra lookup, and the lookup service is a new component you have to keep fast and alive. If it's down, nobody can find anything.
- Underneath all three sits the resize trade. Hash-mod is the simplest rule to write, but it moves almost the whole dataset the day you add a shard. A bucket layer or a ring costs a little indirection on every request, and a resize moves only a small slice. You pay a constant, tiny tax so you never pay the enormous one.
How real systems shard
- MongoDB shards a collection on a shard key and offers both hashed and ranged sharding. It groups documents into chunks, splits chunks as they grow, and a balancer migrates chunks from busy shards to idle ones. Its own advice for a monotonically increasing key is hashed sharding, or a compound key with a high-cardinality prefix.
- DynamoDB hides the shard count. You pick a partition key and it spreads items across internal partitions by hashing it. When one partition gets hot, adaptive capacity isolates the busy items and can split the partition for heat. A single partition caps out around 3,000 reads or 1,000 writes per second, which is why AWS insists on high-cardinality partition keys.
- Vitess shards MySQL. A keyspace is split into N shards with non-overlapping ranges, and a vindex (typically a hash) maps each row's shard key to a shard. Its resharding tool copies and verifies data to the new shards while the old ones keep serving, then cuts over with only a few seconds of read-only time.
- Cassandra hashes the partition key onto a consistent-hashing ring with virtual nodes, so adding a node steals a fair, small slice from the existing ones.
Pitfalls
- Range-partitioning on a monotonic key. Timestamps and auto-increment ids always grow, so every write lands on the last shard. Hash the key, add a random prefix, or pre-split and salt.
- Rehashing on the shard count. hash(key) mod N moves ~80% of your data on a 4→5 resize. Hash onto fixed buckets or use consistent hashing.
- A single hot key. One celebrity or one giant tenant can overwhelm its shard no matter how many shards you add. Replicate it, cache it, or split it into sub-keys.
- Cross-shard queries and transactions. A join or aggregate that spans shards fans out to all of them and merges the results, and an atomic transaction across shards needs a coordinator. Design the shard key so your hottest query stays on one shard.
- Forgetting rebalancing is not free. Even a good scheme copies data over the network during a resize. Do it gradually, throttle it, and keep serving from the old layout until the new one is verified.
In an interview
When a design problem outgrows one machine, say the word sharding and go straight to the two decisions that matter: what's the shard key, and what's the routing rule. Justify the shard key against the dominant query ('we look users up by id, so shard on user id') and against skew. Interviewers are listening for whether you know the key is the important choice.
Reach for hash partitioning by default, and say you'd hash onto fixed buckets or use consistent hashing so you can grow without a full reshuffle. Reach for range partitioning for range scans or time-series, and in the same breath name its hotspot risk and how you'd dodge it. If someone raises a celebrity or a whale tenant, call it a single-hot-key problem and answer with replication or caching, not more shards.
The senior signal is talking about resharding before you're asked. Bring up the ~N/(N+1) movement of naive hash-mod, contrast it with the ~1/(N+1) of a stable scheme, and mention doing the migration live and throttled.
What's the difference between sharding and partitioning?
People use them almost interchangeably. Partitioning is the general act of splitting a dataset into pieces. Sharding usually means partitioning across separate machines for horizontal scale. A partitioned table can live on one server; a sharded one spans many.
Sharding vs replication, same thing?
No, and you usually want both. Sharding splits different data onto different machines for capacity. Replication copies the same data onto several machines for availability and read throughput. A production cluster shards for scale, then replicates each shard so a machine failure doesn't lose that slice.
How is this different from consistent hashing?
Consistent hashing is one specific routing rule, the ring, designed to minimise movement when the machine count changes. Sharding is the broader question of how to split and locate data at all. Consistent hashing, or fixed buckets, is the answer to the resharding half of it.
Can I change the shard key later?
Only by re-sharding the entire dataset, which is a major migration. The shard key decides where every row lives, so treat it as a long-term commitment and pick it for the queries you'll actually run at scale.
References
- Designing Data-Intensive Applications, Ch. 6: Partitioning — Kleppmann — partitioning by key range vs hash, rebalancing, and the fixed-number-of-partitions strategy.
- MongoDB Manual — Sharding — Shard keys, chunks, the balancer, hashed vs ranged sharding.
- Amazon DynamoDB — Partitions, hot keys, and split for heat — How DynamoDB partitions by hash and splits a hot partition.
- Vitess — Resharding — Splitting and merging shards on a live cluster with minimal downtime.