When splitting a table's rows across independent database servers, the placement function can either hash the chosen key or assign contiguous ranges of it. Compare hash-based and range-based placement, and say which workloads suit each.
answer
- Hash = uniform, no order
- Range = ordered scans, cheap splits
- Monotonic key + range = write hotspot
- Hash fixes sequential skew, not frequency skew
- Virtual buckets, not raw modulo
basics
~20 sHashing spreads keys uniformly, so load balances well but ordered and range scans must hit every shard. Range placement keeps neighbouring keys together, so range scans are local and splits are cheap, but sequential keys create a hot shard on the newest range.
solid answer
~60 s**Hash placement** applies a hash to the key and maps the result to a shard (usually via a fixed set of virtual buckets rather than plain modulo, so adding a shard moves buckets instead of rehashing everything). It gives near-uniform distribution regardless of key shape, which is why it is the safe default for monotonic keys like ids and timestamps. The cost: order is destroyed, so `WHERE created_at BETWEEN ...` or an ordered page becomes a scatter-gather across all shards, and you cannot cheaply split one busy region of the key space. **Range placement** assigns contiguous key intervals to shards. Range and prefix scans stay on one shard, results come back ordered, and a hot range can be split in half or moved without touching anything else — this is how systems that auto-balance tablets work. The cost: any monotonically increasing key (auto-increment id, `now()`) makes the last range absorb 100% of inserts, and uneven key density means you must actively rebalance boundaries. Rule of thumb: hash when access is by point key and the key grows monotonically; range when access is by ordered scan and you want cheap online splitting.
go deeper
State the two schemes and their one-line tradeoff: hashing spreads evenly but loses order; ranges keep order but can hotspot.
Explain which query shapes prune under each, why monotonic keys break range placement, and why virtual buckets exist.
Pick a scheme for a stated workload, discuss rebalancing mechanics, single-hot-key limits under hashing, and hybrid hash-plus-time layouts.
Reason about elasticity and operational cost over years: which scheme lets you grow without downtime, what you monitor to detect skew early, and the cost of being wrong for each option.
## The two placement functions Once you have picked a shard key, something must turn a key value into a shard number. There are two mainstream answers. **Hash placement**: `shard = f(hash(key))`. The hash deliberately destroys any meaning in the key so that adjacent values land far apart. Almost nobody uses raw `hash(key) % N` in production, because increasing `N` from 8 to 9 remaps roughly 8/9 of all rows. Instead the hash space is cut into a large fixed number of **virtual buckets** (say 4096 or a consistent-hashing ring), buckets are assigned to physical shards, and growing the cluster moves whole buckets. The mapping stays stable; only the bucket→shard table changes. **Range placement**: shard 1 holds keys `[a, m)`, shard 2 holds `[m, s)`, and so on. Boundaries are metadata, stored in a routing table or a coordinator, and can be edited: a range that grows too big or too hot is split at its midpoint into two ranges that can live on different servers. ## Distribution and hotspots Hashing gives you uniformity for free. It does not care whether your key is an auto-increment integer, a UUID, an email address or a timestamp — the output is close to uniform, so rows and, usually, load spread evenly. This is its dominant virtue, and it is why hash placement is the default in Citus's `create_distributed_table`, in Vitess's default hash vindex, and in most home-grown shard layers. Hashing does **not** protect you from *value* skew: if 30% of your rows carry `tenant_id = 42`, hashing puts 100% of those rows on one shard, because the same key always hashes the same way. Hashing fixes *sequential* skew, not *frequency* skew. Range placement inherits whatever shape the key has. That is fatal for monotonic keys: with ranges on `created_at` or on an auto-increment id, every new row lands in the topmost range, so exactly one shard takes all inserts while the rest serve cold history. It is fine — even ideal — for keys whose density is roughly known and stable, or for systems that continuously split and move ranges automatically. ## Query shapes This is where the two diverge most sharply. With range placement, a predicate like `key BETWEEN x AND y` maps to a contiguous span of ranges, so the router touches only the shards that can possibly contain matches, and each shard returns its rows already ordered — a merge of a few sorted streams. Ordered pagination, time-window queries and prefix scans are all cheap. With hash placement, the same predicate is unanswerable without asking every shard, because logically adjacent keys are physically scattered. Only **equality** on the shard key (or `IN` on a small list) prunes. Ordered pagination becomes a global sort across a scatter-gather result. Both handle point lookups on the shard key equally well. ## Rebalancing and elasticity Range systems rebalance by splitting a range at a chosen point and handing half to another node; the split is local, cheap, and can be triggered automatically by size or heat. This makes range placement attractive for systems that must absorb unpredictable growth without a human deciding shard counts. Hash systems rebalance by moving virtual buckets. That works and is well understood, but a single hot **key** cannot be split at all — its bucket is atomic, and every row with that key is on one shard by definition. If one tenant outgrows a machine under hash placement, your only options are to give it a dedicated shard, or to change the key (for example to a composite `(tenant_id, sub_id)` hash) and move data. ## Composite and layered schemes Real systems mix the two. A common pattern is hash on the tenant, range on time *within* the shard's local table partitioning — you get even tenant spread across servers and cheap time-window pruning and retention inside each server. Another is a two-level scheme: hash placement for the physical bucket, with an explicit directory that can override placement for a handful of outsized keys. ## Choosing Ask two questions. *Is the shard key monotonic?* If yes, range placement will hotspot on writes and you want hashing (or a deliberately randomised prefix). *Do queries scan ordered ranges of the key?* If yes, hashing turns your main access path into a fan-out and you want ranges. For the classic OLTP multi-tenant workload — point and small-set access by tenant, no ordered scans across tenants — hash on tenant id is almost always right. For time-series or log-shaped data queried by window, range on time is natural, with the standing caveat that the newest range is the write hotspot and needs a plan (more granular ranges, a salted prefix, or accepting that the hot range gets its own hardware). ## How to answer Name both functions, give the one-line strength and the one-line failure of each (hash: uniform but no ordered pruning; range: ordered and splittable but hotspots on monotonic keys), then pick one for a stated workload and say what you would monitor to know you were wrong.
- Why do hash-sharded systems map keys into a few thousand virtual buckets instead of computing hash(key) modulo the number of servers?Plain modulo ties placement to the server count, so going from N to N+1 servers changes the target of nearly every key and forces a full data reshuffle. With a fixed large bucket count, the key-to-bucket mapping never changes and scaling only edits the bucket-to-server assignment, so adding a server moves roughly 1/(N+1) of the data. It also lets you move individual buckets to relieve a hot server without touching the hash function.
- You are range-sharding on an auto-increment id and all inserts pile onto the last shard. What options do you have short of re-sharding everything?Short-term you can give the tail range its own, larger machine and split it more aggressively so the frontier moves faster. Structurally you either prefix the key with a small hashed salt so inserts spread across a fixed set of ranges while keeping approximate ordering within each, or switch placement of the write-hot table to hashing while leaving read-heavy historical tables on ranges. Both are real migrations; the honest answer is that a monotonic key and range placement are a mismatch.
saying these in an interview costs you the question
- Claiming hash sharding eliminates hotspots, ignoring that one very frequent key value still lands on one shard
- Using hash(key) % shard_count and treating cluster growth as a trivial config change
- Choosing range placement on a timestamp or auto-increment id without acknowledging the write hotspot
- Believing range predicates can be pruned under hash placement
- Treating the choice as universal rather than derived from the workload's access patterns