Every sharding article lists the strategies: hash, range, directory, consistent hashing. That list is necessary but it is not a decision. Two teams can pick hash sharding and get opposite results, because what matters is the key they hash and the queries they run against it.
This article is about making the decision from evidence. We look at what "balanced" really means, why some keys can never be balanced by any function, how to score candidate keys with a simulator fed by your own query log, and how co-location and secondary indexes change the answer. For the strategy-by-strategy mechanics, read sharding strategies compared alongside this page.
Two decisions: key and function
A sharding scheme is two choices. The partition key is the set of fields that decides where a row lives. The partition function turns the key into a shard: a hash modulo a slot count, a lookup in sorted range boundaries, or an explicit directory. Most systems add a second mapping from logical shards to physical nodes, so data can move without changing the function.
The function determines what you can do cheaply. Hashing spreads writes and destroys order, so range scans on the key hit every shard. Ranges keep order, so scans are cheap, but sequential keys pile up at the end. A directory can express anything but becomes a component that must be fast, consistent and available. The key determines which questions are single-shard. That is usually the more important choice, and it is much harder to change later.
What balanced really means
Engineers often assume a good hash gives equal shards. It gives equal shards on average, which is weaker. If n keys of equal weight are hashed uniformly into m shards, the expected load per shard is n/m, but the largest shard is reliably above that. A standard balls-into-bins result says that when n is much larger than m log m, the maximum is about n/m plus the square root of 2 (n/m) ln m.
Plug in numbers. With 64 shards and 640,000 equally busy customers, the mean is 10,000 per shard and the excess term is the square root of 2 x 10,000 x ln 64, about 290. The busiest shard carries roughly 3 percent more than average: fine. With 64 shards and 6,400 customers, the mean is 100 and the excess is about 29, so the busiest shard runs nearly 30 percent hot. Few keys means poor balance, no matter how good the hash. This is one reason systems use many logical slots or virtual nodes per machine, covered in the consistent hashing deep dive.
The larger problem is that keys do not have equal weight. Tenant sizes, product popularity and account activity follow heavy-tailed distributions. If your biggest tenant generates 8 percent of traffic and you have 64 shards, its shard carries at least 8 percent, five times the 1.6 percent average, whatever function you use. A function of the key cannot split one key. Only changing the key can.
Scoring candidate keys with a simulator
So the right question is not "hash or range" but "which key, scored against our workload". You can answer it before migrating anything. Take a sample of real queries with their filter fields, and a sample of rows with the fields you might shard on. Route both through each candidate key and measure three numbers: the max-to-mean load ratio across shards, the mean and 99th-percentile number of shards a query touches, and the fraction of write transactions that touch more than one shard.
import hashlib, statistics
from collections import Counter
M = 64 # logical shards
def h(*parts):
d = hashlib.blake2b("|".join(map(str, parts)).encode(), digest_size=8).digest()
return int.from_bytes(d, "big") % M
def shards_for(query, key_fields, route):
"""Shards a query must visit: one if it binds every key field, else all."""
if all(f in query["where"] for f in key_fields):
return {route(*[query["where"][f] for f in key_fields])}
return set(range(M))
def score(candidate, rows, queries, txns):
fields, route = candidate["fields"], candidate["route"]
load = Counter()
for r in rows: # row weight = expected request share
load[route(*[r[f] for f in fields])] += r["weight"]
loads = [load[s] for s in range(M)]
fanout = sorted(len(shards_for(q, fields, route)) for q in queries)
cross = sum(1 for t in txns
if len({route(*[w[f] for f in fields]) for w in t["writes"]}) > 1)
return {
"load_ratio": max(loads) / statistics.mean(loads),
"fanout_mean": statistics.mean(fanout),
"fanout_p99": fanout[int(0.99 * (len(fanout) - 1))],
"cross_shard_txn": cross / len(txns),
}
CANDIDATES = {
"hash(customer_id)": {"fields": ["customer_id"], "route": lambda c: h(c)},
"hash(order_id)": {"fields": ["order_id"], "route": lambda o: h(o)},
"hash(customer_id, order_month)": {"fields": ["customer_id", "order_month"],
"route": lambda c, m: h(c, m)},
}The simulator is deliberately simple. It treats a query as single-shard only when it binds every key field with an equality, which is how most routers behave. Extend it with range routing if a candidate uses range boundaries, and weight each row by its share of requests rather than counting rows, because hot rows matter more than many cold ones.
Worked example: an orders service
Consider an orders service. Its main queries are: list a customer's recent orders, fetch an order by id, and a nightly report by date. Writes create an order and its line items in one transaction. Imagine running the simulator on a workload shaped like this, with heavy-tailed customer activity where the largest customer produces 3 percent of traffic. The numbers below are illustrative values showing the typical pattern, not measurements or a benchmark.
| Candidate key | Load ratio | Fan-out mean | Fan-out p99 | Cross-shard txns |
|---|---|---|---|---|
| hash(customer_id) | 2.1 | 1.6 | 64 | 0% |
| hash(order_id) | 1.03 | 38 | 64 | 0% with items keyed by order |
| hash(customer_id, order_month) | 1.4 | 9.5 | 64 | 0% |
Read the trade-offs off the table. Hashing by order id balances perfectly but turns the most frequent query, a customer's orders, into a scatter across every shard. Hashing by customer keeps that query on one shard and keeps each order's transaction local, but the largest customer makes one shard about twice as busy as average. The compound key splits big customers across months, which spreads load but makes "recent orders" touch several shards.
The usual decision is customer id, plus two fixes for what it costs. Fetch-by-order-id is served by encoding the shard in the order id itself, a pattern described in sharded database architecture, so it becomes single-shard without a lookup. The heaviest customers are isolated: a directory overrides the hash for a short list of named tenants and places each on its own shard. The nightly report scatters, which is acceptable for a batch job and is better served from a replica or warehouse anyway.
Monotonic keys and hot tails
Monotonic keys deserve a special warning. Timestamps, auto-increment ids and time-ordered UUIDs all send every new row to the end of the key space. Under range partitioning that means one shard takes all inserts while the rest are idle, and splitting it just moves the hot spot to the new last shard. Under hashing the insert load spreads, but you lose the ability to scan by time.
Three fixes are common. Prefix the key with a small hash bucket, so writes spread over a fixed number of ranges and time scans read every bucket in parallel. Lead with a high-cardinality entity such as device or customer and put time second, so each entity's history is ordered but new writes come from many entities. Or accept the hot shard for a write-ahead stream and partition only historical data by time. Pick by which read you must keep cheap.
Co-location and secondary indexes
Rows that are read or written together should live together. If orders and line items share the customer id as their distribution key, a join between them for one customer is local and a transaction creating both is single-shard. Systems make this explicit in different ways: Citus has co-located tables that share a distribution column, and Spanner has interleaved child tables stored with their parent row. The idea is called an entity group: a root entity and everything that hangs off it.
Co-location constrains every child table to carry the root key, sometimes denormalised into tables that would not otherwise need it. That cost is worth paying for the tables in your hottest transactions, and not for tables read by unrelated access paths.
Queries that filter on something other than the key need a secondary index. A local index lives on each shard and indexes only its rows, so writes stay local but lookups scatter to every shard. A global index is itself sharded by the indexed value, so lookups touch one index shard, but each write updates a second shard, which needs either a distributed transaction or asynchronous maintenance with briefly stale results. Routing and global index mechanics are covered in sharding in system design; for the key decision, count how many of your queries would need each kind.
Keeping the decision honest in production
A simulator scores a snapshot. Workloads drift: a new feature adds a query that filters on a field outside the key, a marketing campaign creates a celebrity tenant, a product line grows tenfold. Keep the same three numbers as live metrics so you notice when the decision stops holding.
- Per-shard load: export request rate, CPU and storage per logical shard, and alert on the max-to-mean ratio rather than on absolute values, which hide imbalance as the fleet grows.
- Fan-out per query shape: tag each query with the number of shards it touched. A rising share of scatter queries usually means a new access path that needs an index or a different service.
- Top keys: sample the hottest keys per shard every few minutes. The first warning of a celebrity key is one key moving from a fraction of a percent of traffic to several percent.
When a key turns hot, act in order of cost. Cache its reads in front of the shard. Move it to a dedicated shard with a directory override. Only if neither works, change the key for that entity, for example by appending a bucket number, and accept that reads for it now fan out across its buckets. Re-keying a whole table is a migration measured in weeks of dual writes and backfills, which is why the simulator exercise up front is worth an afternoon.
Failure modes and trade-offs
| Failure mode | How it shows up | Prevention |
|---|---|---|
| Key chosen for balance, not access | Most queries scatter; p99 latency follows the slowest shard | Score fan-out on the real query log, not just distribution |
| Celebrity key | One shard saturated while others idle | Directory overrides or key salting for named hot keys |
| Too few logical shards | Uneven nodes, painful splits | Start with many logical slots mapped to few nodes |
| Monotonic insert hot spot | Last range shard pinned at 100% write load | Bucket prefix or entity-first compound key |
| Unplanned cross-shard writes | Distributed transactions in the hot path | Co-locate entity groups by root key |
The underlying trade-off never goes away: a key optimised for one access path makes others scatter. The goal is not to avoid scatter but to put it on the queries that can afford it.
What to do next
- Export a week of query logs with filter fields and frequencies, and rank queries by volume and latency sensitivity.
- Measure key frequency for each candidate field; find the share of traffic owned by the top 10 keys.
- Run the simulator for three to five candidate keys and record load ratio, fan-out and cross-shard write fraction.
- Pick the key that makes your top queries single-shard, then plan explicit fixes for the rest: id encoding, directory overrides, secondary indexes.
- Check for monotonic components and choose a bucket prefix or entity-first ordering.
- List the tables in your hottest transaction and co-locate them on the root key.
- Start with many more logical shards than nodes so you can rebalance by moving slots.