A team assigned rows to database shards with hash(key) mod N, where N is the number of servers. Why does adding a server become painful, and what assignment schemes avoid that?
answer
- mod N: change N, ~all keys move
- Doubling moves half, but only buys 2x
- Ring: new node steals from clockwise neighbour → K/N
- Virtual nodes even out arcs and spread the takeover
- Fixed B buckets + directory = stable hash, movable unit
basics
~20 sChanging N changes the modulus, so nearly every key maps somewhere new and almost the whole dataset has to move. Consistent hashing (a hash ring) moves only about 1/N of keys; the common relational approach is fixed virtual buckets — hash into e.g. 4096 slots and move whole slots between servers.
solid answer
~60 sWith `hash(key) mod N`, the mapping is a function of N. Go from 8 to 9 servers and a key's destination changes unless `hash mod 8 == hash mod 9` — roughly 8/9 of all rows relocate. Every addition is effectively a full migration, so teams either never add capacity or only ever double. Two standard fixes: **Consistent hashing.** Hash both keys and nodes onto a ring; a key belongs to the next node clockwise. Adding a node steals keys only from its neighbour — about K/N of the data. Virtual nodes (many ring positions per server) keep the distribution even and make the load transfer spread across all existing servers rather than dumping on one. **Fixed virtual buckets** — the scheme most relational shard platforms actually use. Hash into a large, *never-changing* number of logical slots (1024, 4096, 65536), and keep a small directory mapping slot to physical shard. Adding a server means reassigning a subset of slots and moving only their data. The hash function never changes, the directory is tiny and cacheable, and it gives you an explicit unit for rebalancing and for pinning a noisy tenant to its own server.
code
text · 10 linesmodulo-N (N: 8 -> 9)
key A: h=1000 -> 1000%8=0 -> 1000%9=1 MOVES
key B: h=1001 -> 1001%8=1 -> 1001%9=2 MOVES
key C: h=1008 -> 1008%8=0 -> 1008%9=0 stays
=> ~8/9 of all rows relocate
fixed buckets (B=4096, directory bucket -> shard)
key A: h=1000 -> bucket 1000 -> shard 2 (bucket never changes)
add shard 9: reassign ~455 of 4096 buckets in the directory,
copy only those buckets' rowsgo deeper
Know that modulo-N reshuffles almost everything when N changes, and that consistent hashing exists to keep the movement small.
Quantify it — roughly 8/9 of keys move going 8→9 under modulo, versus about K/N on a hash ring — and describe fixed virtual buckets plus a directory as the practical relational scheme.
Argue the operational case for buckets: stable hash function, tiny cacheable versioned directory, bucket as the unit of migration, verification, throttling and rollback, plus explicit pinning for noisy tenants.
Treat the bucket count and the hash-versus-range choice as effectively irreversible commitments, and reason about them against the largest credible fleet size, the scan patterns the product needs, and the rebalancing automation you are willing to own.
## What goes wrong with modulo-N Assign shard as `shard = hash(key) % N`. It is uniform, stateless, and trivially computed by any client — which is why it keeps getting chosen. The defect only shows up at the first capacity change. The mapping is a function of N, so when N changes, so does almost every key's destination. Growing 8 → 9 relocates roughly 8/9 of all rows; only the keys where `h mod 8` happens to equal `h mod 9` stay put. There is no incremental version of that migration: you must copy nearly the whole dataset, keep both mappings live during the move, and cut over atomically. The usual coping strategy is **doubling**: N → 2N. If you use the hash's low bits, every key either stays where it is or moves to exactly one new sibling shard, so only half the data on each shard moves, and each shard's departing half goes to a single, known destination. That is genuinely simpler to operate. Its cost is granularity — you can only ever buy 100% more capacity, which is expensive at 32 shards and absurd at 256, and you cannot rebalance a single hot shard without touching everything. ## Consistent hashing Map the hash output onto a circular key space (say 0 to 2^32−1). Place each server at one or more positions on that ring by hashing its identifier. A key belongs to the first server encountered walking clockwise from the key's hash. Now adding a server inserts one new position on the ring. It takes over only the arc between itself and the previous position — the keys that used to belong to its clockwise neighbour. On average that is K/N keys, where K is the total. Removing a server hands its arc to the next one clockwise. The rest of the mapping is untouched. With one position per server, arc sizes are uneven (some servers get much larger arcs by luck) and a departing server dumps its entire load on exactly one neighbour. **Virtual nodes** fix both: give each server, say, 200 ring positions. Arcs average out, and when a server leaves, its 200 arcs are absorbed by many different neighbours. Virtual node counts can also be weighted, letting a bigger machine take proportionally more of the ring. ## Fixed virtual buckets (the practical relational answer) Most sharded relational deployments use a simpler variant that gives the same property with less machinery. Choose a large, fixed bucket count — 1024, 4096, 65536 — and compute `bucket = hash(key) % B`. **B never changes.** A separate small directory maps bucket → physical shard. Adding a server means updating the directory for some buckets and physically moving those buckets' rows. Why this tends to win in practice: - **The hash is stable forever.** No client ever needs to recompute placements; only the directory changes. - **The directory is tiny.** A few thousand entries: cacheable in every router, versioned, cheap to distribute. - **Buckets are the unit of everything.** Migration, verification, throttling, and rollback all operate per bucket, so a reshard is a sequence of small, independently retryable moves rather than one enormous copy. - **Explicit placement is possible.** A tenant generating disproportionate load can have its buckets pinned to dedicated hardware — a purely directory-level change. Pure consistent hashing gives you no such handle. - **Fractional growth works.** Moving from 8 shards to 10 is just a reassignment of about a fifth of the buckets. The design decision is B. Too small and buckets become too coarse to balance (and too large to move); too large and the directory and per-bucket bookkeeping get noisy. A common heuristic is a few dozen to a few hundred buckets per shard at the largest fleet size you can foresee, since B is effectively permanent. ## Range sharding and chunk splits Hashing scatters adjacent keys, which kills range scans. If the workload needs ordered scans, shard by **key ranges** instead, and let the system split a range into two **chunks** when it exceeds a size or load threshold, then move chunks between servers to balance. This gives elasticity at a finer grain than any static hash and keeps ranges scannable. The costs are a hotspot on the highest range when the key is monotonically increasing (timestamps, sequential ids), and a real rebalancer that has to be trusted and observed. ## Choosing Hash with fixed virtual buckets is the default for OLTP fleets that mostly do point lookups by an entity key. Range with chunk splits is right when ordered scans matter or key distribution is unpredictable. Pure `mod N` is right only for a fixed-size deployment that will provably never change — which, in practice, describes nothing.
- Why do virtual nodes matter in a consistent-hashing ring?With one ring position per server, arc sizes vary widely by luck, so load is uneven, and when a server is removed its whole arc lands on a single neighbour, which can then tip over. Giving each server many positions averages the arc sizes and spreads both the takeover and the departure across many peers. Position counts can also be weighted so heterogeneous hardware gets proportional load.
- How do you pick the bucket count B, given that it is effectively permanent?Size it for the largest fleet you can plausibly reach, allowing a few dozen to a few hundred buckets per shard at that size — so a few thousand buckets for a fleet that might reach tens of shards. Too few and you cannot balance finely or move data in small pieces; too many and the directory, the per-bucket metadata, and the migration bookkeeping become noisy. Because changing B means rehashing everything, err on the generous side.
- When would you shard by key ranges with chunk splitting instead of by hash?When the workload needs ordered scans — time windows, alphabetical listings, cursoring in key order — which hashing destroys by scattering adjacent keys. Ranges also adapt to skewed key distributions, since a hot or large range can be split and moved. The trade-off is that a monotonically increasing key sends all writes to the last range, so you need either a compound key or a rebalancer you actively monitor.
Modulo-N is seating guests by 'name length mod number of tables' — add a table and everyone stands up. Buckets are numbered seating cards assigned to tables on a chart: to add a table you move a few cards, and nobody's card changes.
saying these in an interview costs you the question
- Believing hash(key) mod N only moves 1/N of the data when N changes
- Thinking consistent hashing eliminates data movement rather than bounding it to roughly K/N
- Omitting virtual nodes and then being surprised by uneven load or a neighbour collapsing on removal
- Choosing a bucket count equal to the current shard count, which reduces the scheme back to modulo-N
- Assuming hash sharding still supports efficient range scans