skip to content

When sharding a database by a chosen partition key, what are the trade-offs between range-based, hash-based, and directory-based (lookup-table) partitioning strategies?

level: middleimportance: must knowfreq 88%

answer

  1. range = locality but hotspots
  2. hash = even load, loses ordering
  3. directory = flexible, extra hop plus critical dependency
  4. consistent hashing minimizes rebalance churn
  5. DynamoDB hash partition key + range sort key

basics

~20 s

Range partitioning splits data by value ranges (like A-M, N-Z), keeping related data together but risking hot shards. Hash partitioning spreads data evenly using a hash function but breaks up ranges, making sorted scans slow. Directory-based partitioning uses a lookup table to map keys to shards, giving flexibility at the cost of an extra lookup step.

solid answer

~50 s

Range partitioning assigns contiguous key ranges to shards (e.g. user IDs 1-1M on shard A), preserving natural ordering so range scans and sorted queries are efficient, but concentrating load on whichever shard currently holds the 'hot' range (e.g. newest or sequential data) and needing active rebalancing as ranges grow unevenly. Hash partitioning applies a hash function to the key to pick a shard, distributing load and volume evenly and avoiding hotspots from sequential keys, but destroying locality so range scans have to fan out to every shard. Directory-based partitioning keeps an explicit mapping table consulted before routing a query, offering maximum flexibility (arbitrary reassignment, per-tenant placement, incremental rebalancing) at the cost of an extra lookup hop and the mapping table becoming a critical, must-be-highly-available component. Real systems often combine these, e.g. consistent hashing plus virtual nodes for rebalancing flexibility.

go deeper

for a junior

Can describe range vs hash partitioning at a high level with one trade-off each.

for a middle

Can pick an appropriate strategy given a stated query pattern and explain the hotspot risk of sequential keys.

for a senior

Designs compound shard keys (hash prefix plus range suffix) and reasons about consistent hashing and directory-service availability trade-offs.

for a principal

Sets shard-key policy across a platform, anticipates future query pattern shifts, and plans migration paths between strategies without downtime.

## The decision you cannot easily take back Sharding (horizontal partitioning) splits one logical dataset across multiple physical databases or nodes so no single machine has to hold or serve all the data. The central design decision is the **partitioning strategy**: the function that decides, given a record's shard key, which shard it lives on. Three classic strategies dominate, and picking the wrong one for your access pattern is one of the most consequential and hardest-to-reverse mistakes in a system's data layer. ## Range partitioning Range partitioning assigns each shard a contiguous span of key values - for example, shard 0 holds user IDs 1 through 1,000,000, shard 1 holds 1,000,001 through 2,000,000, and so on. To route a query, you compare the key against the boundaries and pick the matching shard. - **The major benefit is locality.** Because keys near each other in value are stored near each other physically, range scans (give me all orders between two dates, all users alphabetically from M to P) and sorted queries can be served by touching only one or a few shards, and ordering doesn't require a distributed merge across every shard. - **The cost is uneven load.** If the shard key correlates with insertion order - a common case, since many systems shard on an auto-incrementing ID or a timestamp - then all new writes land on whichever shard currently owns the highest range, turning that one shard into a hotspot while older shards sit idle. Bigtable and HBase use range partitioning (calling the units 'regions' or 'tablets') and specifically warn against monotonically increasing row keys for exactly this reason. ## Hash partitioning Hash partitioning applies a hash function to the shard key and uses the result (often mod the shard count, or via consistent hashing) to pick a shard. Because a good hash function scatters even sequential or clustered input keys uniformly across the output space, this evens out both data volume and write load - no single shard becomes a target for all new inserts. The trade-off is the mirror image of range partitioning's strength: **locality is destroyed**. A query like 'all orders from the last hour' now has no single shard to go to, because consecutive timestamps hash to essentially random shards; the query must fan out to every shard, gather results, and merge them, which is more expensive and has tail latency determined by the slowest shard. Cassandra and DynamoDB default to hash-based placement for exactly this even-distribution property. ## Directory-based (lookup-table) partitioning Directory-based (lookup-table) partitioning takes a different approach: instead of a formula, a service maintains an explicit table mapping keys or key-ranges to shard IDs, and every read or write first consults that directory to find the right shard. This is **the most flexible option** - you can place a specific high-traffic tenant on its own dedicated shard, move a single key's data without touching anything else, or implement custom placement logic such as data residency (keep EU users' data on EU shards) - but it adds an extra network hop and a new critical dependency: the directory service itself must be fast, highly available, and consistent, or every query in the system is affected. Vitess (YouTube/Slack's MySQL sharding layer) and many multi-tenant SaaS platforms use directory-based routing for tenant-level placement flexibility. ## Range and hash, mirror images | Strategy | The major benefit | The cost | |---|---|---| | **Range** | locality: range scans and sorted queries can be served by touching only one or a few shards | uneven load, if the shard key correlates with insertion order - then all new writes land on whichever shard currently owns the highest range | | **Hash** | evens out both data volume and write load, because a good hash function scatters sequential or clustered input keys uniformly | locality is destroyed: a query like 'all orders from the last hour' must fan out to every shard, and tail latency is determined by the slowest shard | ## Blends used in practice In practice, production systems often blend these ideas rather than picking one purely. - **Consistent hashing** is a refinement of hash partitioning designed to minimize data movement when the shard count changes: instead of `hash(key) mod N`, which reshuffles almost everything when N changes, keys and shards are placed on a hash ring, and only the keys between the old and new node's position need to move. - Adding **virtual nodes** (each physical shard owns many small points on the ring) smooths out load distribution further and makes rebalancing more granular. - Some systems **layer range partitioning within a hash bucket** - DynamoDB hashes the partition key to pick a physical partition, but within that partition, items with the same partition key are stored in sort-key order, giving even distribution across partitions and locality within a partition, which is exactly why choosing a good partition key plus sort key pair matters so much in DynamoDB schema design. ## Let the query pattern choose The choice ultimately follows the query pattern, not the data model in isolation: if your dominant queries are point lookups by a well-distributed key (get user by ID), hash partitioning is close to a free win; if your dominant queries are range scans over a naturally ordered key (get all events in a time window), range partitioning avoids expensive fan-out, but you must actively manage hotspots, often by choosing a compound key that combines a random or bucketed prefix with the natural range so writes spread across shards while each shard still holds a contiguous range internally.

  • Why do systems like Bigtable/HBase warn against using a monotonically increasing row key (like a timestamp or auto-increment ID) with range partitioning?
    Because all new writes would target the single shard/region currently holding the highest key range, making that one node a write hotspot while the rest sit idle. The common fix is to prefix the key with a hash, shard number, or reversed timestamp so writes spread across regions.
  • What operational problem does consistent hashing solve compared to plain hash(key) mod N partitioning?
    With mod N, adding or removing a single shard changes the target shard for nearly every key, forcing a massive data reshuffle. Consistent hashing places shards and keys on a ring so only the keys adjacent to the changed shard need to move, dramatically reducing rebalancing cost.
  • In a directory-based sharding scheme, why does the directory/lookup service itself become a critical design concern?
    Because every single query has to consult it before it can even reach the data, so if it's slow, unavailable, or inconsistent, the entire system is affected regardless of how healthy the underlying shards are. It typically needs its own replication, caching, and high-availability story, often more rigorous than the shards it routes to.

Range partitioning is like a library that shelves books by call-number ranges - great for browsing a section, but everyone wanting the newest bestsellers crowds one shelf. Hash partitioning is like assigning each book a random locker by lottery number - lockers fill evenly, but you can't browse 'all books starting with S' without checking every locker. Directory-based partitioning is like asking a librarian who keeps a card catalog mapping every book to its exact shelf - flexible placement, but now you depend on the librarian being available and fast.

saying these in an interview costs you the question

  • picks a shard key without asking what the dominant query pattern is
  • doesn't recognize that a sequential key plus range partitioning causes hotspots
  • thinks hash partitioning preserves range-scan efficiency
  • forgets the directory table itself needs to be highly available
  • assumes resharding is free once you pick hash partitioning

context