7.3 Sharding of Key-Value Data
Adding one node
- 16/24
- keys moved
- 67% / 20%
- vs ideal
67% of keys move when the ideal is 20%. Hashing mod N reshuffles almost everything because changing N changes the modulus for every key — this is why the book calls it unusable for a growing cluster.
The goal: spread data AND query load EVENLY. With a fair share each, 10 nodes should handle 10× the data and 10× the throughput (ignoring replication). Adding/removing a node should allow rebalancing.
Vocabulary of unfairness:
| Term | Meaning |
|---|---|
| Skewed | Some shards have more data or queries than others — makes sharding much less effective |
| Hot shard / hot spot | A shard with disproportionately high load. In the extreme, all load on one shard → 9 of 10 nodes idle and your bottleneck is the single busy node |
| Hot key | One key with particularly high load — e.g. a celebrity in a social network |
We need an algorithm taking a record's PARTITION KEY and returning its shard. In a key-value store the partition key is usually the key, or the first part of the key; in a relational model it might be a column of a table — not necessarily its primary key. The algorithm must be amenable to rebalancing in order to relieve hot spots.
3.1 Sharding by key range
Assign a contiguous range of partition keys (min to max) to each shard — like the volumes of a paper encyclopedia.
“Simply having one volume per two letters would lead to some volumes being much bigger than others. To distribute data evenly, the shard boundaries need to adapt to the data.”
Who chooses the boundaries:
| Mode | Systems |
|---|---|
| Manual | Vitess (a sharding layer for MySQL) |
| Automatic | Bigtable, HBase, MongoDB (range option), CockroachDB, RethinkDB, FoundationDB |
| Both | YugabyteDB (manual and automatic tablet splitting) |
Within each shard, keys are stored in SORTED order (B-tree or SSTables, Ch 4). Two advantages:
- Range scans are easy
- You can treat the key as a CONCATENATED INDEX to fetch several related records in one query — e.g. sensor data keyed by timestamp, where a range scan gets all readings from a particular month
The downside — hot shards from nearby writes:
The new cost: fetching multiple sensors within a time range now needs a separate range query for each sensor.
Rebalancing key-range shards:
- Pre-splitting — HBase and MongoDB let you configure an initial set of shards on an empty database. This requires that you already have some idea what the key distribution will look like
- Splitting — as volume/throughput grow, split an existing shard into two or more, each holding a contiguous subrange. Similar to what happens at the top level of a B-tree
- Merging — if large amounts of data are deleted, merge adjacent small shards
- Split triggers — a configured size (HBase default: 10 GB) or, in some systems, write throughput persistently above a threshold. Thus a HOT shard may be split even if it isn't storing much data, so its write load spreads
✔ The number of shards adapts to the data volume — small data → few shards → small overheads; huge data → each shard capped at a configurable maximum.
✗ Splitting a shard is EXPENSIVE — it requires all its data to be rewritten into new files, similarly to a compaction. And a shard that needs splitting is often ALSO one under high load, so the cost of splitting can EXACERBATE that load, risking it becoming overloaded.
3.2 Sharding by hash of key
If you don't care whether partition keys are near each other (e.g. tenant IDs), hash the partition key first.
A good hash function takes SKEWED data and makes it UNIFORMLY DISTRIBUTED. A 32-bit hash returns a seemingly random number from 0 to 2³²−1; even very similar input strings hash to evenly distributed values — but the same input always produces the same output.
For sharding, the hash function need NOT be cryptographically strong: MongoDB uses MD5; Cassandra and ScyllaDB use Murmur3.
⚠️ Language built-in hashes may be unsuitable: in Java's
Object.hashCode()and Ruby'sObject#hash, THE SAME KEY MAY HAVE A DIFFERENT HASH VALUE IN DIFFERENT PROCESSES.
(a) Hash mod N — the naive approach, and why it fails
Three nodes, key assignment hash(key) % 3:
- node 0 — hashes 0, 3, 6, 9, 12 …
- node 1 — hashes 1, 4, 7, 10 …
- node 2 — hashes 2, 5, 8, 11 …
Add a fourth node and the assignment becomes hash(key) % 4:
| Hash | Was on | Now on | |
|---|---|---|---|
| 3 | node 0 | node 3 | ✗ moved |
| 4 | node 1 | node 0 | ✗ moved |
| 5 | node 2 | node 1 | ✗ moved |
| 6 | node 0 | node 2 | ✗ moved |
| 7 | node 1 | node 3 | ✗ moved |
| 8 | node 2 | node 0 | ✗ moved |
| 9 | node 0 | node 1 | ✗ moved |
“Most of the keys have to be moved from one node to another.”
mod Nis easy to compute but leads to VERY INEFFICIENT REBALANCING because of a lot of unnecessary movement. We need an approach that moves AS LITTLE DATA AS POSSIBLE.
(b) Fixed number of shards
Create many more shards than nodes and assign several shards to each node.
hash(key) % 1000, and the system separately tracks which shard is on which node.Only entire shards are moved — cheaper than splitting shards. The number of shards does not change, and the assignment of keys → shards does not change. Only the assignment of shards → nodes changes.
Note the cutover detail: reassignment is not immediate — it takes time to transfer a large amount of data over the network — so THE OLD ASSIGNMENT IS USED FOR ANY READS AND WRITES THAT HAPPEN WHILE THE TRANSFER IS IN PROGRESS.
Two practical tips:
- Choose a shard count divisible by many factors, so the dataset can be evenly split across various numbers of nodes — not requiring the node count to be a power of 2
- Account for mismatched hardware: assign MORE shards to more powerful nodes so they take a greater share of load
Used by Citus (sharding layer for PostgreSQL), Riak, Elasticsearch, Couchbase.
The limitations:
It works well AS LONG AS YOU HAVE A GOOD ESTIMATE of how many shards you'll need when you first create the database. You can then add/remove nodes easily — subject to the limitation that YOU CAN'T HAVE MORE NODES THAN SHARDS.
If the number turns out wrong, an EXPENSIVE RESHARDING is required: split each shard and write it out to new files, using a lot of additional disk space. Some systems DON'T ALLOW RESHARDING WHILE CONCURRENTLY WRITING, which makes it difficult to change the shard count WITHOUT DOWNTIME.
The Goldilocks problem: since each shard holds a fixed fraction of total data, shard size grows proportionally to total data.
- Shards too large → rebalancing and recovery from node failures become expensive.
- Shards too small → too much overhead.
“The best performance is achieved when the size of shards is just right” — hard to achieve if the shard count is fixed but the dataset size varies.
(c) Sharding by hash range — the best of both
Combine key-range sharding with a hash function so each shard contains a range of HASH VALUES rather than a range of KEYS.
- ✔ A shard can be split when it becomes too big or too heavily loaded, so the number of shards adapts to the volume of data rather than being fixed.
- ✗ Range queries over the partition key are not efficient — keys in the range are now scattered across all shards.
The saving grace for range queries:
If keys consist of two or more columns and the partition key is only the FIRST of them, you can still perform efficient range queries over the SECOND and later columns. As long as all records in the range query have the SAME partition key, they will be in the SAME shard.
(So (user_id, timestamp) hashed on user_id still gives you efficient "all events for user X between T1 and T2".)
Used by: YugabyteDB, DynamoDB; an option in MongoDB. Cassandra and ScyllaDB use a variant:
The hash space 0–1024 is split into contiguous ranges with random boundaries, with several ranges assigned to each node. The figure shows 3 ranges per node; the actual defaults are 16 per node in Cassandra and 256 per node in ScyllaDB.
“Some ranges are bigger than others, but by having multiple ranges per node, those imbalances tend to even out.”
Adding node 3: node 1 transfers parts of two of its ranges and node 2 transfers part of one of its ranges, so the new node gets an approximately fair share without transferring more data than necessary from one node to another.
Aside — the warehouse equivalent: BigQuery: the partition key determines the partition, "cluster columns" determine sort order within it. Snowflake assigns micro-partitions automatically but lets you define cluster keys. Delta Lake supports both manual and automatic partition assignment plus cluster keys. Clustering improves range-scan performance AND compression AND filtering.
(d) Consistent hashing
A consistent hashing algorithm maps keys to a specified number of shards satisfying two properties:
- The number of keys mapped to each shard is roughly equal
- When the number of shards changes, AS FEW KEYS AS POSSIBLE are moved
⚠️ "Consistent" here has NOTHING to do with replica consistency (Ch 6) or ACID consistency (Ch 8) — it describes the tendency of a key to stay in the same shard if possible.
Cassandra/ScyllaDB's algorithm is similar to the original definition. Other algorithms: highest random weight (rendezvous hashing) and jump consistent hashing.
With these approaches, rather than a small number of existing shards being SPLIT INTO SUBRANGES to create new shards for a new node, THE NEW NODE IS INSTEAD ASSIGNED INDIVIDUAL KEYS that were previously scattered across all the other nodes. Which is preferable depends on the application.
3.3 Skewed workloads and relieving hot spots
Consistent hashing ensures keys are uniformly distributed across nodes — but that DOESN'T MEAN THE ACTUAL LOAD IS UNIFORMLY DISTRIBUTED.
Skew means: much more data under some partition keys than others, or the request rate to some keys is much higher. You can still end up with some servers overloaded while others sit almost idle.
The canonical case: a post by a celebrity with millions of followers causes a storm of activity — a large volume of reads and writes to THE SAME KEY (the celebrity's user ID, or the ID of the action people are commenting on).
Three mitigations:
① Dedicated shard. A system that defines shards by ranges of keys (or hashes) can put an individual hot key in a shard BY ITSELF — perhaps even assigning it a dedicated machine.
② Application-level key splitting (salting).
Hot key: celebrity_123
- Write — append two random digits, giving
celebrity_123_00…celebrity_123_99. That splits the writes evenly across 100 keys, which are then distributed to different shards. - Read — must now read all 100 keys and combine them.
“The volume of reads to each shard of the hot key is not reduced; only the write load is split.”
The bookkeeping cost: it makes sense to salt only the small number of hot keys — for the vast majority of low-throughput keys this would be unnecessary overhead. So you also need a way to TRACK WHICH KEYS ARE SPLIT, and a PROCESS FOR CONVERTING a regular key into a specially managed hot key.
And it's dynamic: a viral post may be hot for a couple of days and then calm down. Some keys may be hot for WRITES while others are hot for READS, necessitating DIFFERENT STRATEGIES.
③ Automated heat management. Some cloud services do this automatically — Amazon calls it heat management or adaptive capacity.
3.4 Automatic vs manual rebalancing
| Mode | Examples |
|---|---|
| Fully automatic | DynamoDB — promoted as able to automatically add and remove shards to adapt to big load changes within a matter of minutes |
| Fully manual | Explicitly configured by an administrator |
| Middle ground | Couchbase and Riak generate a suggested shard assignment automatically but REQUIRE AN ADMINISTRATOR TO COMMIT IT |
Why automatic is attractive: less operational work; systems can even autoscale to adapt to workload changes.
Why automatic is dangerous:
- ① Rebalancing is expensive — it reroutes requests and moves a large amount of data. Done carelessly it can overload the network or the nodes and harm the performance of other requests.
- ② The system must keep processing writes while rebalancing. “If a system is near its maximum write throughput, the shard-splitting process might not even be able to keep up with the rate of incoming writes.”
- ③ The cascading-failure loop — automation plus automatic failure detection:
For that reason, it can be good to have A HUMAN IN THE LOOP for rebalancing. It's slower than a fully automatic process, but it can help prevent operational surprises.
Manual rebalancing is also useful for PREEMPTIVELY rebalancing when a surge is expected from a known event — Cyber Monday sales, or ticket sales for the World Cup.
(This is the same argument as Ch 6's manual-failover preference and Ch 2's "autoscaling is cool, but predictable load may prefer manual" — a consistent theme: automation that reacts to failure signals can amplify failures.)