When choosing a shard key strategy for a data store, what are the differences between lookup-based, range-based, and hash-based sharding, and what does each cost you?
answer
- range = locality, risks hotspot on monotonic keys
- hash = even load, breaks range queries
- lookup/directory = manual control, extra hop
- scatter-gather on hash for range queries
- monotonic key = tail-of-range hotspot
basics
~20 sRange sharding groups nearby keys (like consecutive IDs) on the same shard, which is great for range scans but can pile new writes onto one shard. Hash sharding spreads keys evenly using a hash function, avoiding that but breaking range queries. Lookup sharding uses an explicit table to say which shard each key lives on, giving full control at the cost of an extra hop.
solid answer
~50 sRange-based sharding stores contiguous key ranges together (e.g., order_id 1–1M on shard A), which makes range scans and ordered reads cheap because they stay on one shard, but it risks a hot shard when writes concentrate at one end of the range — a monotonically increasing ID or timestamp means all new writes land on whichever shard owns the tail. Hash-based sharding runs the key through a hash function before assigning a shard, spreading writes evenly and avoiding that hotspot, but it destroys locality: a range query now has to fan out to every shard and merge results. Lookup-based (directory) sharding keeps an explicit key-to-shard mapping table, giving full manual control over placement — useful for isolating a noisy tenant or rebalancing selectively — at the cost of an extra lookup hop and a mapping store that itself must stay consistent and highly available. Teams pick hash for even load, range when ordered/range queries dominate, and lookup when they need fine-grained placement control.
go deeper
Should recognize there's more than one way to decide which shard a row goes to, even without full trade-off detail.
Should explain the core trade-off of each strategy — range's hotspot risk vs. locality, hash's even load vs. broken range queries, lookup's control vs. extra hop.
Should reason about which strategy fits a given workload's access pattern and proactively flag the monotonic-key hotspot risk before it happens in production.
Should discuss combining strategies (e.g., hashed prefix + range suffix) and how the choice constrains future query patterns and migration paths years down the line.
## The three strategies A shard key strategy is the rule that maps a row's shard key value to a physical shard, and the three common strategies trade **locality**, **load balance**, and **operational control** against each other in different ways. - **Range-based sharding** assigns contiguous blocks of key values to the same shard — shard A might own customer IDs 1 through 1,000,000, shard B owns 1,000,001 through 2,000,000, and so on. Because rows with nearby key values live together, this strategy makes range-oriented operations cheap: "give me all orders between ID 500,000 and 500,100" or "give me events from the last hour" (when time is the key) can usually be satisfied by touching a single shard or a small contiguous set of them, without any cross-shard coordination. - **Hash-based sharding** instead applies a hash function to the key and uses the hash output (often mod the shard count, or via a hash ring) to pick a shard. Because a good hash function scatters similar or sequential inputs to unrelated outputs, writes for consecutive keys land on essentially random shards, which spreads load very evenly across the cluster. - **Lookup-based** (also called **directory-based**) **sharding** keeps an explicit mapping — a table or service that records "this key, or this key range, lives on this shard" — and every request consults that directory before routing. ## Why each one exists Each strategy exists to solve a different failure mode of the others. 1. **Range sharding** exists because ordered/range access patterns are extremely common (time-series data, paginated listings, ID-based backfills), and hash sharding makes those patterns expensive by scattering related rows across every shard. 2. **Hash sharding** exists specifically to defeat the load-skew problem that plagues naive range sharding: many real-world key spaces are not written to uniformly — IDs are usually assigned in increasing order, timestamps always increase, and both mean that under pure range sharding, essentially 100% of new writes target the single shard that currently owns the highest range, while older shards go idle. 3. **Lookup/directory sharding** exists for cases where neither pure math (range or hash) gives you the placement control you need — for example, deliberately isolating one enormous or noisy tenant onto its own dedicated shard, or supporting incremental, selective rebalancing without having to touch every key's mapping at once. ## What each one costs The trade-offs are concrete and show up immediately in system design. | Strategy | What you pay | What you gain | |---|---|---| | **Range sharding** | With range sharding, you pay in potential hotspots on monotonic keys | you gain cheap, single-shard range scans and simple, predictable shard boundaries that are easy to reason about and split | | **Hash sharding** | With hash sharding, you pay by losing locality — any range query, any "give me all rows for this related group" query, has to be sent to every shard and merged in the application layer (scatter-gather), which is slower and more complex than a local scan | but you gain near-uniform load distribution regardless of key access patterns, which is usually the bigger win for write-heavy systems | | **Lookup sharding** | With lookup sharding, you pay an extra network hop and a new critical dependency (the directory itself must be fast, consistent, and highly available, since every request needs it) | but you gain the ability to place data anywhere, override the default algorithm for exceptional keys, and rebalance incrementally by just updating map entries rather than rehashing the whole keyspace | ## Failure signatures in production In production, these strategies show up as distinct failure signatures. - A **range-sharded** system with a monotonic key shows one shard pegged at high CPU/IO while siblings idle, and the fix usually involves either changing the key (adding a randomized or reversed prefix) or manually splitting the hot range. - A **hash-sharded** system shows unexpectedly high latency and cost on range/list queries, because what should be a single-shard scan is silently fanning out cluster-wide, and the tail latency of these scatter-gather queries grows with shard count since it's bound by the slowest responder. - A **lookup-sharded** system's failure mode is usually availability or staleness of the directory itself — if the mapping service is down or serving a stale mapping after a rebalance, requests can be routed to the wrong shard or fail outright. ## Where it shows up A well-known real-world example: DynamoDB and Cassandra default to hash-based partitioning specifically to avoid the classic monotonic-key hotspot problem, at the documented cost that range queries across the partition key are not supported — you must model access patterns around the partition key up front or add a secondary index. Systems like Vitess, by contrast, support range-based (and directory) sharding for MySQL specifically because many workloads still need efficient range scans on the shard key, accepting the hotspot risk as something the schema designer must actively manage (e.g., by choosing a non-monotonic shard key) rather than something the sharding scheme solves automatically.
- Why does range sharding with a monotonically increasing key (like an auto-increment ID or a timestamp) tend to create a hot shard?Because every new row has a key value greater than everything written before it, so every new write targets whichever shard currently owns the top of the range. That one shard absorbs 100% of new write traffic while shards holding older ranges go idle. The fix is usually to change the key — add a random or hashed prefix, or shard on something other than a strictly increasing value.
- How would you support an efficient range query, like 'all orders from the last 7 days,' on a hash-sharded store?You typically fan the query out to every shard in parallel (scatter-gather), filter and merge the results in the application or routing layer, and accept higher latency than a single-shard scan would give you. Some teams also maintain a separate range-friendly index or a dedicated time-series store alongside the hash-sharded primary store specifically to serve these queries cheaply.
Range sharding is like alphabetizing files into cabinets by last name — easy to pull everyone whose name starts with 'M', but if everyone you hire this year happens to be named 'Zheng,' the Z cabinet overflows while A through Y sit empty. Hash sharding is like assigning each file a random locker number instead — no cabinet ever overflows from a naming trend, but now finding 'everyone hired this year' means checking every locker in the building.
saying these in an interview costs you the question
- Doesn't know hash sharding breaks efficient range queries
- Claims range sharding never has hotspot problems
- Confuses lookup/directory sharding with hash sharding
- Can't name a concrete cost for any of the three strategies
- Suggests you can get even load and cheap range queries with the same simple strategy