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.

Advertisement

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.

Query logsampled, real trafficRow samplekey frequenciesCandidate keysfields + functionSimulatorroute, count, scoreLoad ratiomax / mean per shardFan-outshards per queryCross-shard writesfraction of txnsDecisionkey, function, index planScore candidate keys against the workload you have, before data makes the choice permanent
Workload-driven shard-key selection. Real queries and a sample of rows are routed through each candidate key in a simulator, which scores balance, fan-out and cross-shard writes.

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.

Advertisement

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 keyLoad ratioFan-out meanFan-out p99Cross-shard txns
hash(customer_id)2.11.6640%
hash(order_id)1.0338640% with items keyed by order
hash(customer_id, order_month)1.49.5640%

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 modeHow it shows upPrevention
Key chosen for balance, not accessMost queries scatter; p99 latency follows the slowest shardScore fan-out on the real query log, not just distribution
Celebrity keyOne shard saturated while others idleDirectory overrides or key salting for named hot keys
Too few logical shardsUneven nodes, painful splitsStart with many logical slots mapped to few nodes
Monotonic insert hot spotLast range shard pinned at 100% write loadBucket prefix or entity-first compound key
Unplanned cross-shard writesDistributed transactions in the hot pathCo-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

  1. Export a week of query logs with filter fields and frequencies, and rank queries by volume and latency sensitivity.
  2. Measure key frequency for each candidate field; find the share of traffic owned by the top 10 keys.
  3. Run the simulator for three to five candidate keys and record load ratio, fan-out and cross-shard write fraction.
  4. Pick the key that makes your top queries single-shard, then plan explicit fixes for the rest: id encoding, directory overrides, secondary indexes.
  5. Check for monotonic components and choose a bucket prefix or entity-first ordering.
  6. List the tables in your hottest transaction and co-locate them on the root key.
  7. Start with many more logical shards than nodes so you can rebalance by moving slots.
Key takeaway: A sharding strategy is a partition key plus a partition function, and the key matters more. Even a perfect hash leaves the busiest shard noticeably above average when keys are few, and no function can split a single hot key. Score candidate keys against your real query log for load ratio, fan-out and cross-shard writes, choose the key that keeps your most important queries on one shard, and fix the rest with id encoding, hot-key overrides, bucket prefixes, co-location and secondary indexes.