Once a dataset outgrows one machine, you split it across several, and the first design decision is how a row finds its home. Engineers often frame this as choosing a shard key, which matters, and the sharding overview covers that choice. But two systems with the same key can behave completely differently depending on the strategy that maps the key to a machine. That strategy decides which queries are cheap, where hot spots form, and above all what it costs to add capacity next year.
This article compares the strategies directly: range, hash modulo N, token rings, fixed logical slots, directory lookup, and entity or geographic sharding. For each it explains the mechanism, which real systems use it, how it reshards and how it fails. It ends with a decision procedure, a small router in code and a worked sizing. Partitioning inside one database, which shares many ideas but no network, is covered in the partitioning article.
The two mappings
Every sharding strategy is a composition of two functions. The first maps a key to a logical shard. The second maps each logical shard to a physical node. Some designs collapse them into one, and that is where most resharding pain comes from: if the node count is part of the first function, adding a node changes where almost every key lives.
The strategies differ in where they put flexibility. A range scheme can split shards anywhere in key order. A slot scheme fixes the first mapping forever and moves only the second. A directory makes the first mapping an editable table. Keep this lens in mind; it explains every row of the comparison table later.
Range sharding
Range sharding assigns contiguous key ranges to shards: keys from A to F on one, G to M on the next. Ranges are split when they grow too large or too busy, and moved to balance load. Bigtable and HBase call them tablets and regions; CockroachDB calls them ranges and, by default, splits a range when it reaches 512 MiB, a value set in its zone configuration; MongoDB's ranged shard keys group documents into chunks that a balancer moves.
The advantage is order. A range scan over a key prefix, such as all orders for one customer between two dates, touches one shard or a few adjacent ones. Splits are local: one shard becomes two and nothing else moves. The weakness is monotonic keys. If the key is a timestamp or an auto-increment ID, every insert lands at the end of the key space, on the last range, and one node takes all the writes while the rest idle. Range systems answer with key design: prefix the key with something well distributed, or hash a prefix and keep order within it.
Hash modulo N
The simplest hash scheme computes hash(key) % N, where N is the number of shards, usually equal to the number of nodes. Distribution is even for any key that hashes well, and the router needs no metadata at all.
The cost appears on resize. Going from N to N plus 1 shards changes the result of the modulo for about N out of N plus 1 keys; growing from 10 to 11 nodes moves roughly 91 percent of your data. Resizing becomes a full migration, so teams postpone it, run hot, and then do it in an emergency. Range queries are gone too: adjacent keys scatter across every shard, so a scan becomes a query to all shards. Use plain modulo only when N will never change, such as a fixed fleet of caches you are willing to cold-start, and even then prefer one of the next two schemes.
Token rings and virtual nodes
Consistent hashing places nodes at points on a hash ring and assigns each key to the next node clockwise. Adding a node takes over only the keys between it and its predecessor, so about one in N plus 1 keys move instead of almost all. The mathematics, including the variance argument for virtual nodes and the ring-free alternatives, is worked through in the consistent hashing article and the rendezvous hashing article.
Cassandra and Dynamo-style stores use this design. Each Cassandra node owns several tokens, called virtual nodes, to even out ownership; Cassandra 4.0 lowered the default num_tokens to 16 and pairs it with an allocation algorithm that places tokens for a target replication factor. Operationally, a ring still couples the two mappings: adding a node changes which keys it owns, so data streams to it from neighbours, and replicas are chosen by walking the ring. That is fine, but it means capacity moves in units set by token placement rather than units you choose.
Fixed logical slots
The slot approach hashes keys into a large, fixed number of logical shards, chosen once and far larger than the node count, and then assigns slots to nodes with a small, editable table. Redis Cluster hashes every key with CRC16 modulo 16384 slots; a node owns a set of slots, and resharding moves slots between nodes, key by key, with redirects telling clients the new owner. Citus distributes a table across a fixed number of shards per table, 32 by default, and its rebalancer moves whole shards between workers. Elasticsearch fixes the number of primary shards at index creation; the split API can multiply it only by a factor allowed by the index.number_of_routing_shards setting, which is set when the index is created. Vitess maps each row through a vindex to a keyspace ID and assigns contiguous keyspace-ID ranges to shards named like -80 and 80-, which can be split.
Slots give you the best of both earlier schemes for point lookups: even distribution, and rebalancing that moves only the slots you choose. The decision you cannot undo cheaply is the slot count. Too few and a single slot becomes larger than a node can hold, or too hot to be served by one; too many and per-shard overhead (files, connections, metadata, planning time) adds up. Choose a count that leaves each node holding at least tens of slots at your largest expected cluster size.
Directory sharding
A directory, or lookup, strategy stores the mapping explicitly: a table says tenant 42 lives on shard 7. The router consults it, usually through a cache. Any placement is possible, so you can put a large customer on dedicated hardware, keep a regulated customer in one region, or move one tenant without touching anyone else.
The directory is now critical infrastructure. It must be highly available and strongly consistent during moves, and every router must handle a stale cache: the usual answer is a version number per entry, with shards rejecting writes for keys they no longer own so the router refreshes and retries. Directories also need the granularity to be coarse, such as per tenant or per bucket of keys, or the table grows as large as the data.
Entity and geographic sharding
Entity sharding picks a key that is the natural unit of the application, usually the tenant, account or user, so that all of that entity's rows live together. Citus calls this co-location: tables distributed on the same column and shard count place matching rows on the same node, so joins and transactions within one tenant stay local. Geographic sharding is a special case keyed by region, often for data residency rules.
The risk is skew. Tenants are not equal; one enterprise customer can be a thousand times the median. A shard that holds that customer is a hot shard by construction. Combine entity sharding with a directory override so whales can be isolated, and see the hot key mitigation article for techniques when a single key, not just a shard, is hot.
Side by side
| Strategy | Point lookup | Range scan | Add capacity | Hot spot risk | Examples |
|---|---|---|---|---|---|
| Range | One shard via metadata | Cheap, few shards | Split and move, local | Monotonic keys | CockroachDB, HBase, MongoDB ranged |
| Hash mod N | Computed, no metadata | All shards | Moves most data | Low | Simple caches |
| Token ring | Computed from ring | All shards | Stream from neighbours | Low, vnodes smooth it | Cassandra |
| Fixed slots | Computed plus slot map | All shards | Move chosen slots | One hot slot | Redis Cluster, Citus, Vitess |
| Directory | Lookup, cached | Depends on layout | Move any entry | Directory itself | Many SaaS platforms |
| Entity or geo | One shard per entity | Within an entity | Move entities | Large tenants | Citus co-location |
Resharding: what each strategy costs
Moving data while serving traffic follows the same pattern everywhere: copy a snapshot of the moving unit to its new home, stream the changes made during the copy, briefly block or fence writes, switch the mapping, then drain and delete the old copy. The strategies differ in how big the unit is and how many units must move. With slots, ranges and directories, you move a chosen handful of units; with a ring, you move a slice from each neighbour; with modulo, you move almost everything.
The cut-over is where data is lost. A router with a stale mapping can write to the old owner after the switch. Defend in depth: version the mapping, have each shard check ownership on writes, and keep the old owner in a read-only, forwarding state until every router has acknowledged the new version. The router below shows the slot-plus-directory hybrid many teams converge on.
import zlib
NUM_SLOTS = 4096 # fixed forever; choose generously
class WrongOwner(Exception): # raised by a shard that no longer owns the key
def __init__(self, current_version):
self.current_version = current_version
class Router:
def __init__(self, slot_map, overrides, version):
self.slot_map = slot_map # list[node] of length NUM_SLOTS
self.overrides = overrides # {tenant_id: node} for isolated whales
self.version = version
def slot(self, tenant_id: str) -> int:
return zlib.crc32(tenant_id.encode()) % NUM_SLOTS
def node_for(self, tenant_id: str) -> str:
return self.overrides.get(tenant_id) or self.slot_map[self.slot(tenant_id)]
def execute(self, tenant_id, sql, params, pool, refresh):
for _ in range(2):
node = self.node_for(tenant_id)
try:
return pool[node].run(sql, params, expect_version=self.version)
except WrongOwner as e: # shard fenced: mapping moved under us
self.__dict__.update(refresh(min_version=e.current_version).__dict__)
raise RuntimeError(f"routing for {tenant_id} did not converge")
Worked example: choosing for a multi-tenant orders service
A B2B platform stores orders for 30,000 tenants: 6 TB today, growing about 60 percent a year, with almost every query scoped to one tenant and filtered by date. The largest tenant holds 4 percent of all rows; the median holds a few megabytes. Nodes are sized to hold about 1.5 TB comfortably.
Tenant is the obvious key, since queries are tenant-scoped and joins stay local. Range sharding on tenant would cluster new tenants, which get sequential IDs, onto the last range. Hash modulo node count would make each expansion a full migration, and this system needs to grow from four nodes to roughly sixteen within three years. Fixed slots with a directory override fit best. Pick 4,096 slots: at sixteen nodes that is 256 slots per node, fine-grained enough to balance, while 30,000 tenants over 4,096 slots puts about seven tenants in each slot. Within a tenant, order rows by date so that date filters stay index range scans on one node.
The 4 percent tenant is 240 GB today and about 1 TB in three years: it fits on one node but would make its slot's node far heavier than the others. Give it a directory override and its own node now, before it forces an emergency move. Expansion then becomes moving 256 chosen slots per new node. Each slot holds about 1.5 GB today (6 TB over 4,096 slots) and about 6 GB at the three-year size, so a move is a series of small, throttleable copies rather than one bulk migration.
Failure modes
- Monotonic keys on range shards. All writes hit one range. Prefix or hash the leading key component.
- Too few logical shards. Slots or shards larger than a node, or too hot to split, freeze the layout. Over-provision the count.
- Celebrity entities. One tenant dominates a shard. Isolate it through a directory override before it saturates.
- Scatter-gather creep. New features add queries without the shard key; each touches every shard and p99 latency becomes the slowest shard's. Review queries for the key.
- Stale routing during moves. Writes land on the old owner. Version the mapping and fence old owners.
- Cross-shard transactions. Work spanning shards needs two-phase commit or sagas; design the key so common transactions stay on one shard.
What to do next
- Write down your top ten queries and check that each carries the proposed shard key.
- Decide which mapping you want to stay fixed; default to fixed slots or ranges, not modulo of the node count.
- Size the logical shard count for your three-year node count, with at least tens of shards per node.
- Measure entity size skew and plan directory overrides for the largest tenants now.
- Build the resharding path (copy, stream, fence, switch, drain) and rehearse it on staging before you need it.
- Monitor per-shard size, request rate and p99 latency, and alert when any shard drifts well above the median.