Hash sharding is the default way to split a large dataset across machines: hash the shard key, map the hash to a shard, and send the request there. It spreads keys evenly without anyone choosing boundaries, which is why it underpins Cassandra's partitioner, Redis Cluster, many document stores and most home-grown sharding layers. The idea fits in one line of code, and most of the production trouble comes from details of that line: which hash, which modulus, and what happens to the mapping when the number of shards changes.
This article is about the key-to-shard function itself. The surrounding machinery, such as the router, the versioned shard map and live shard splits, is covered in Database Sharding Architecture, and rings with virtual nodes are covered in Consistent Hashing in depth. Here we compare hash-mod-N, fixed logical slots and jump consistent hash by how much data they move, then cover what hashing costs you: range scans, hot keys and cross-key operations.
Why hash, and what the hash must guarantee
A shard key with a natural order, such as a timestamp or an auto-increment id, sends all new writes to whichever shard owns the newest range. Range sharding then needs active splitting to keep up. Hashing removes the order: consecutive ids land on unrelated shards, so write load spreads evenly from the first day, and nobody has to pick split points.
The hash must have two properties. It must be uniform, so that each shard gets close to 1/N of the keys. It must also be stable: every client, in every language and every release, must compute the same value for the same key forever, because the hash decides where data already lives. Fast non-cryptographic hashes such as MurmurHash3 or xxHash are the usual choice; Cassandra's default partitioner uses Murmur3. Do not use a language's built-in hash. Python randomises string hashing per process unless PYTHONHASHSEED is fixed, and many runtimes make no promise that the value will stay the same across versions.
import xxhash # pip install xxhash; any stable 64-bit hash works
def key_hash(shard_key: str) -> int:
# Encode explicitly: the same logical key must give the same bytes everywhere.
return xxhash.xxh64_intdigest(shard_key.encode("utf-8"), seed=0)Normalise keys before hashing. If one service hashes "User-42" and another hashes "user-42", the same user lives on two shards.
Hash mod N and the resize problem
The obvious mapping is shard = hash(key) % N. It is perfectly balanced and needs no state. Its problem shows up the first time N changes. A key stays put only if its hash gives the same remainder under both moduli, and for N to N+1 that is true for just 1/(N+1) of keys. Growing from 4 to 5 shards moves 80% of the data; from 10 to 11, about 91%. A simulation over 200,000 random 64-bit hashes gives 79.9% and 90.8%.
Moving most of the data to add one machine means a full rewrite of the dataset under live traffic, with double the storage during the copy. Mod-N is fine for data you can rebuild cheaply, such as a cache you are willing to cold-start, or for systems that genuinely never resize. For a primary store it is a trap.
Fixed logical slots and a slot map
The standard fix is to add a level of indirection. Hash keys into a fixed, large number of logical slots, chosen once and never changed, and keep a small map from slot to physical shard. The key-to-slot step is still mod, but the modulus never changes. Resharding moves whole slots, and the map records where each slot lives.
SLOTS = 1024 # fixed for the life of the data; pick 10-100x max shards
def slot_for(key: str) -> int:
return key_hash(key) % SLOTS
class ShardMap:
def __init__(self, assignment: list[int], version: int):
self.assignment = assignment # assignment[slot] = shard id
self.version = version
def shard_for(self, key: str) -> int:
return self.assignment[slot_for(key)]
def rebalance(self, new_shards: int) -> list[tuple[int, int, int]]:
"""Move the fewest slots so every shard owns floor or ceil of SLOTS/new_shards."""
target = [SLOTS // new_shards + (1 if s < SLOTS % new_shards else 0) for s in range(new_shards)]
owned = {s: [] for s in range(new_shards)}
for slot, shard in enumerate(self.assignment):
owned.setdefault(shard, []).append(slot)
spare = []
for shard, slots in owned.items():
cap = target[shard] if shard < new_shards else 0
spare += slots[cap:]
owned[shard] = slots[:cap]
moves = []
for shard in range(new_shards):
while len(owned[shard]) < target[shard]:
slot = spare.pop()
moves.append((slot, self.assignment[slot], shard))
owned[shard].append(slot)
return moves # (slot, from, to) for the migration jobRedis Cluster is the best-known example. It defines 16,384 hash slots, computes the slot as CRC16 of the key modulo 16,384, and assigns slot ranges to primaries. Adding a node means migrating some slots to it, and clients follow redirects when a slot has moved.
Pick the slot count once and pick it generously: enough that the largest cluster you can imagine still gets many slots per shard, so balance stays even and each move is small. Too many slots only costs a slightly bigger map. Too few is permanent, because changing the slot count rehashes everything.
Jump consistent hash and the other options
If you want minimal movement without a map, jump consistent hash (Lamping and Veach, 2014) maps a 64-bit key to one of N buckets with no stored state. Growing from N to N+1 moves only 1/(N+1) of keys, all of them to the new bucket, and balance is near-perfect.
def jump_hash(key: int, num_buckets: int) -> int:
"""Lamping & Veach jump consistent hash. key is a 64-bit unsigned int."""
b, j = -1, 0
while j < num_buckets:
b = j
key = (key * 2862933555777941757 + 1) & 0xFFFFFFFFFFFFFFFF
j = int((b + 1) * (float(1 << 31) / float((key >> 33) + 1)))
return bOn the same 200,000 hashes, going from 10 to 11 buckets moved 9.2% of keys, against 90.8% for mod-N, and 10 buckets each received between 19,856 and 20,153 keys. The catch is that buckets are numbered 0 to N-1 and you can only add or remove at the end. You cannot retire bucket 3 because its machine died, so jump hash maps keys to logical shards that are then mapped to replicated machines. It also cannot express weights, so a node twice as large cannot take twice the keys. When you need arbitrary membership or weights, use a ring with virtual nodes or rendezvous hashing.
| Scheme | Data moved, N to N+1 | State | Remove any shard | Weights |
|---|---|---|---|---|
| hash mod N | N/(N+1), e.g. 80% for 4 to 5 | none | no | no |
| fixed slots + map | about 1/(N+1), in whole slots | slot map | yes | yes, by slot count |
| jump hash | 1/(N+1) | none | only the last | no |
| ring with vnodes | about 1/(N+1) | token ring | yes | yes, by vnode count |
Choosing the shard key
Even placement of keys is not the same as even load. Three properties of the shard key decide whether the shards are actually balanced.
Cardinality. A key with few distinct values cannot spread. Sharding orders by country gives at most a couple of hundred values, with a handful holding most of the rows. Prefer keys with millions of values, such as user id, account id or device id.
Skew in data per key. If one tenant has 40% of all rows, hashing the tenant id puts 40% of the data on one shard whatever the hash does. Either shard on something finer, such as tenant id plus entity id, or give that tenant its own placement, which is what directory-based sharding is for.
Skew in traffic per key. A celebrity account or a flash-sale item can take thousands of times the average request rate. Hashing spreads keys, not requests, so a hot key is pinned to one shard. The usual mitigations are a read cache in front of the hot key and, for writes, key salting: write to item:991#k for a random k from 0 to 7 and read all eight sub-keys and combine them. Salting multiplies read cost by the salt count, so apply it only to keys you have measured as hot.
What hashing costs you
Hashing destroys key order, and that has three costs you should accept knowingly.
Range scans become scatter-gather. "All orders between two dates" no longer maps to a few adjacent shards; it must ask every shard and merge the results. A fan-out query is only as fast as its slowest shard. If each shard independently exceeds your latency target on 1% of requests, a query that touches 32 shards exceeds it on about 27.5% of requests, because 1 - 0.9932 is about 0.275. The common fix is a compound key: hash only the first part to choose the shard, and keep rows sorted by the second part within it. Cassandra's partition key plus clustering columns is exactly this design, so "orders for user 48213, newest first" stays a single-shard, ordered read.
Secondary lookups fan out. A query by email when the shard key is user id has to visit every shard if each shard indexes only its own rows. The alternative is a global index, itself hash-sharded by email, which turns the query into two point lookups but must be kept in sync with the base data, usually asynchronously.
Cross-key operations span shards. A transaction or join over two users on different shards needs a distributed protocol. Co-locate data that is used together by sharding it on the same key. Redis Cluster supports this with hash tags: if a key contains a non-empty substring in braces, only that substring is hashed, so {user1000}.cart and {user1000}.profile share a slot and can be used together in multi-key commands.
Worked example: adding two shards
A messaging service stores conversations keyed by conversation_id on 8 shards, using 1,024 slots, so each shard owns 128 slots. Messages are clustered by timestamp within each conversation, so loading the last 50 messages is one ordered read on one shard.
Growth pushes the shards to 70% disk, so the team adds two shards. The rebalance asks for 10 shards with 102 or 103 slots each, so each existing shard gives up 25 or 26 slots and 204 slots move in total: 19.9% of the data, close to the ideal 2/10. Each slot moves by copying a snapshot, streaming changes since the snapshot, briefly pausing writes to the slot, bumping the map version, and resuming. Routers that still hold the old map version get a "wrong shard" error with the new version and retry. Had the system used hash % 8 directly, the same expansion would have moved 80% of all data.
One problem remains. A single public broadcast conversation receives 30% of the write traffic, and its shard runs hot whatever the slot layout. Hashing cannot fix a single hot key, so the team salts that conversation's message writes across 16 sub-keys and merges them on read, and caches the most recent page for readers.
Failure modes
- Hash drift between clients. A new service uses a different hash or a different byte encoding of the key, reads miss and writes go to the wrong shard. Publish the key-to-slot function as a shared library with test vectors, and check them in CI.
- Too few slots. 64 slots on 48 shards leaves some shards with one slot and others with two, a 2x imbalance that cannot be fixed without rehashing everything.
- Unbounded fan-out. A dashboard query that hits every shard runs fine with 8 shards and saturates the cluster at 128. Track fan-out per query type and send analytical reads to a replica or a warehouse.
- Stale maps during moves. A router writes to a slot's old owner after cutover and the write is lost. Owners must reject writes for slots they no longer own, checked against the map version.
- Hot partitions mistaken for hot shards. Adding shards does nothing for a single hot key. Measure per-key traffic before resharding.
What to do next
- Write down your shard key, its cardinality, and the largest share of rows and of traffic held by a single key value.
- Replace any built-in or ad hoc hash with one stable 64-bit hash in a shared library, with test vectors checked in every client.
- If you use hash mod N on a primary store, introduce fixed logical slots now, with 10 to 100 times more slots than your largest planned shard count.
- List every query that does not include the shard key, measure its fan-out and p99, and decide on a compound key, a global index or an offline path for each.
- Find your top 10 keys by request rate and decide in advance whether each one gets a cache, salting or its own placement.
- Rehearse a slot move in staging, including a router with a stale map, and confirm that old owners reject writes after cutover.