skip to content

When you split a dataset across multiple nodes, what's the difference between range partitioning and hash partitioning, and what kind of query does each make fast or slow?

level: juniorimportance: must knowfreq 70%

answer

  1. range=sorted, contiguous
  2. hash=scrambled, uniform
  3. hot range vs even load
  4. scan fast vs scan fans out
  5. monotonic keys hot-spot ranges

basics

~20 s

Range partitioning keeps keys in sorted order across nodes (A-M on node 1, N-Z on node 2), so range scans are fast but one busy range can overload a node. Hash partitioning scrambles keys by a hash function, spreading load evenly but killing range scans.

solid answer

~50 s

Range partitioning assigns contiguous key ranges to each partition, preserving sort order across the cluster. That makes range queries (scan all users created in March) fast since the answer lives on one or few nodes, and supports efficient ordered iteration. The cost is uneven load: if writes cluster around one part of the keyspace (recent timestamps, a popular prefix), that partition becomes a hot spot while others sit idle. Hash partitioning runs each key through a hash function and assigns it based on the hash value, which scatters keys pseudo-randomly across nodes regardless of their natural ordering. This gives near-uniform load distribution for point lookups and writes, but destroys locality — a range scan now has to fan out to every partition and merge results, and there's no cheap way to iterate keys in sorted order. Systems often pick per use case: Bigtable/HBase use range partitioning for scan-heavy workloads; DynamoDB and Cassandra default to hash partitioning for write-heavy, point-lookup workloads.

go deeper

for a junior

Should state the basic mechanism of each (ordered ranges vs hash scatter) and give one example of when each is preferable — range for scans, hash for even load — without needing to discuss mitigation techniques.

for a middle

Should additionally identify the hot-spot failure mode for range partitioning under monotonic keys and know that hash partitioning breaks range-scan efficiency, ideally naming at least one real system that uses each.

for a senior

Should discuss the salting/prefix-hashing hybrid mitigation, composite-key designs that combine both approaches, and reason about which scheme fits a given workload's read/write pattern.

for a principal

Should reason at the schema/product level: anticipate future query patterns before choosing a partition key, weigh the cost of a partition-key migration (usually a full data reshuffle), and evaluate vendor-specific constraints (e.g., DynamoDB's partition-key/sort-key model) against evolving access patterns.

## What partitioning has to decide Partitioning (also called **sharding**) splits a dataset across multiple nodes so that no single machine has to hold or serve the whole dataset. The central design question is: given a key, which partition does it live on? **Range partitioning** and **hash partitioning** are the two dominant answers, and they trade off in almost opposite ways. ## Range partitioning: contiguous intervals Range partitioning divides the keyspace into contiguous intervals and assigns each interval to a partition — for example: - usernames A-M go to partition 1, N-Z to partition 2 - or timestamps in January go to partition 1, February to partition 2 Because the partitions are defined by boundaries in the natural ordering of the key, the data within a partition stays **sorted**, and the partitions themselves are sorted relative to each other. This means a range query — 'give me every order placed between 09:00 and 10:00' — can be answered by touching only the one or few partitions whose range overlaps the query, and within that partition the answer can be produced by a sequential scan rather than a scatter-gather across the whole cluster. Systems built for scans and time-series access, like **Bigtable**, **HBase**, and many log-oriented stores, favor range partitioning for exactly this reason. ## Hash partitioning: scattered by hash value The mechanism for hash partitioning is different: each key is passed through a hash function (e.g., `MD5`, `MurmurHash`), and the resulting hash value — not the key itself — determines the partition, typically via something like `hash(key) mod N` or a position on a hash ring. Because a good hash function distributes its output pseudo-randomly and uniformly regardless of patterns in the input, keys that were adjacent in the original keyspace (user IDs 1001, 1002, 1003) end up on completely different, effectively random partitions. This is precisely why it's chosen: write and read load spreads evenly across the cluster even if the application's key-generation pattern is skewed or monotonic. ## The trade-off The trade-off is symmetric and unavoidable given the two mechanisms. - **Range partitioning risks hot partitions.** If the workload writes keys correlated with time or another monotonically increasing value — order IDs, event timestamps, an auto-incrementing primary key — then all recent writes land on the single 'newest' partition while older partitions go cold. This single-partition bottleneck can cap the entire cluster's write throughput at the capacity of one node, no matter how many nodes you add. - **Hash partitioning avoids this by construction, but pays for it by breaking locality.** Any query that depends on key order (range scans, prefix scans, 'give me the next 100 keys after X') now has to either fan out to every partition and merge results, or isn't supported efficiently at all. There is also no cheap way to co-locate related keys — two keys that are logically adjacent will typically live on unrelated nodes. ## Failure modes in production In production these failure modes show up differently. 1. A **range-partitioned** system under a monotonic write pattern shows one node pegged at high CPU/IO while its neighbors are idle — classic hot-partition symptoms, often with write-throttling on just that shard while cluster-wide capacity looks fine in aggregate. 2. A **hash-partitioned** system under a scan-heavy query pattern instead shows uniformly moderate load everywhere but high query latency and high inter-node network traffic, because every scan touches every partition. Neither failure looks like 'the disk is full' — they look like the wrong partitioning scheme for the workload. ## Where it shows up in real systems A concrete real-world example: - **HBase/Bigtable** use range partitioning on the row key, and teams are explicitly warned against monotonically increasing row keys (like timestamps) for exactly the hot-region problem described above — the standard mitigation is to prefix the key with a hash or reverse the timestamp bits to scatter writes across regions, effectively borrowing hash partitioning's load-distribution property while keeping a range-partitioned engine underneath. - **DynamoDB**, by contrast, hash-partitions on the partition key by default; it only supports efficient range queries within a single partition key's sort-key range (via composite primary keys), because global range scans across partition keys aren't cheap in a hash-partitioned system. - **Cassandra** similarly hashes the partition key onto a ring and only preserves ordering within a partition via clustering columns. ## Choosing one The practical takeaway is that this choice is driven by the dominant query pattern: | Scheme | Workload | |---|---| | **Range partitioning** | pick when the workload needs ordered scans and you can control or randomize the key prefix to avoid hot spots | | **Hash partitioning** | pick when the workload is point-lookup/write-heavy and even load distribution matters more than scan efficiency | Many production schemas hedge by hashing a prefix of a composite key while keeping a range-ordered suffix, getting load distribution and local ordering at once.

  • If a range-partitioned table is showing one hot partition because of monotonically increasing keys, what's a common fix that doesn't abandon range partitioning entirely?
    Prefix the key with a hash bucket, a random salt, or reverse the timestamp's bit order so writes scatter across the keyspace while still preserving local ordering within each bucket. This is the standard HBase/Bigtable pattern: you sacrifice some global scan efficiency in exchange for spreading write load, effectively hybridizing hash and range partitioning.
  • Can you get both fast range scans and even load distribution in one system?
    Only partially, and only by narrowing the scan scope: composite keys where you hash the high-cardinality prefix (e.g., customer ID) but range-order a low-cardinality suffix (e.g., order date) give you even load across customers and fast range scans within a single customer. A true global range scan across all customers still has to fan out, because the top-level partitioning is hash-based.
  • Why can't you just add more nodes to fix a hot range partition?
    Because the hot range is owned by one partition, and that partition is bound to specific nodes without further splitting — adding unrelated nodes doesn't relieve the partition serving the hot range. You'd need the system to dynamically split that range further, which is a separate rebalancing mechanism, not just adding capacity.

Range partitioning is like shelving library books alphabetically by title — great for browsing 'all books starting with S', but if everyone donates books titled 'The...' that shelf overflows. Hash partitioning is like assigning each book a random locker number by lottery — lockers fill evenly, but finding 'all books starting with S' means checking every locker.

saying these in an interview costs you the question

  • Says hash partitioning is 'always better' or 'always faster' without naming what it sacrifices
  • Doesn't recognize that monotonic keys cause range hot-spots
  • Believes range queries are equally efficient under hash partitioning
  • Can't name a mitigation for range hot-spots (salting/prefixing)
  • Confuses partitioning with replication

context