Core Concept

Consistent Hashing

Consistent hashing limits key movement when nodes join or leave, which is why distributed caches and Dynamo-style databases rely on a hash ring instead of modulo.


1. What It Is

When you're sharding a cache in an interview, you'll reach for consistent hashing because naive modulo remaps almost every key every time you add a node. It's the go-to answer for elastic distributed caches and Dynamo-style partitions.

What:

A hashing topology where resizing the server partition array only requires remapping 1/N of the keys on average β€” instead of nearly all of them.

Primary purpose:

Prevent massive data rebalancing cascades and cache invalidation storms when nodes scale up or down.

Usually used for:

Distributed caches, DynamoDB-style databases, and stateful reverse proxies.

2. Core Mental Model

Picture servers and keys on the same ring β€” ownership moves clockwise, so only a slice of keys shifts when membership changes:

β­• The Shared Ring Space

Map both physical servers and keys onto the exact same circular hash ring (from 0 to 2^64 - 1).

🧭 Clockwise Ownership

A key is assigned to the first server encountered moving clockwise. Removing a node only shifts its immediate segment.

🎭 Virtual Node Spans

Deploy 100-200 virtual points per physical machine to distribute keys uniformly across the ring.

In the room

Many candidates start with hash(key) % N. If the interviewer nods but asks what happens when you scale from 10 to 11 nodes, that's your cue to draw a ring and mention virtual nodes. Follow up with hot-key mitigation β€” the ring alone doesn't solve celebrity traffic.

3. Why It Matters in HLD

Consistent hashing matters the moment node membership becomes elastic β€” auto-scaling, rolling deploys, spot preemption. We frame it around three lenses:

Needed When:

Node membership changes frequently due to auto-scaling, spots, rolling deploys, or hardware crashes.

Avoids:

Naive modulo rebalancing (hash(key) % N) which remaps almost 100% of keys whenever N changes, crashing downstream DBs.

Optimizes For:

Minimal network rebalance blast radius, elastic scaling velocity, and predictable routing latencies.

4. Architecture & Data Flow

Narrate the ring as an interview walkthrough. Step 1 β€” Map the ring: hash both servers and keys onto the same 0…2^64 space. Step 2 β€” Clockwise ownership: a key belongs to the first server encountered moving clockwise. Step 3 β€” Node join: adding a server only steals the segment before it on the ring β€” about 1/N keys move. Step 4 β€” Node leave: the next clockwise server absorbs the departed node's segment. Step 5 β€” Virtual nodes: scatter 100–200 vNodes per physical machine so load stays uniform even with few servers.

Loading...

Virtual Nodes (vNodes) Allocation

To prevent hashing hotspots, physical servers map to multiple scattered points on the ring:

Loading...

In the room

When they ask what happens scaling from 10 to 11 nodes with modulo, draw the ring and say "about 1/11 of keys move." Then mention vNodes and hot-key mitigation β€” the ring alone does not solve celebrity traffic.

5. Key Characteristics

These numbers are what we cite when the interviewer asks about rebalance cost and memory overhead:

  • Minimal Key Remapping: Only 1/N of keys migrate on average when scaling out by one node.
  • Virtual node count trades memory for uniform key distribution on the ring:
V (Virtual Nodes)DistributionMemoryRebalance Cost
Low (10 vNodes)Uneven (Hotspots likely)Minimal (< 1 KB per server)Ultra-fast
Optimal (100–200 vNodes)Near-Uniform (Skew < 5%)Small (A few KB per server)Highly manageable
High (500+ vNodes)Flawlessly UniformLarger (Megabytes at scale)Slower update times

6. Strategic Tradeoffs

The ring solves remapping β€” not every distributed problem. We state the trade-off:

BenefitCost
Elastic Scaling (adding or removing a node only remaps 1/N of total keys)Memory Lookup Overhead (must traverse a sorted hash ring array in O(log S) time)
Load Balancing via vNodes (distributes keys uniformly even with low physical server counts)Rebalance Overhead (moving data on node changes still incurs disk/network migration IO)

7. Failure / Bottleneck Awareness

Consistent hashing does not eliminate hot keys or rebalance I/O β€” we volunteer these mitigations:

πŸ”₯ Hot Key on One Ring Segment

Problem: Even with uniform vNodes, a single hot key (e.g. celebrity_tweet_100) can concentrate all its traffic on one broker.

Mitigation: Local caches on the hot path, or replicate the key under salted suffixes so reads spread across multiple ring positions.

πŸ›‘ Rebalance Traffic Spikes

Problem: Adding a node still moves ~1/N of data β€” at Cassandra scale that can be gigabytes of background transfer competing with live traffic.

Mitigation: Rate-limit and phase rebalancing; run during low-traffic windows where possible.

8. Common HLD Usage

Name the ring when the problem involves elastic cache pools or Dynamo-style partitions:

ProblemUsage
Consistent Distributed CacheRouting cache keys across Memcached/Redis nodes to avoid full invalidation on restarts
Stateful WebSocket ConnectionsDirecting user chat sockets statefully to the exact same gateway server with minimum disconnects
DynamoDB / Cassandra ShardingDeciding which database partition owns which primary key hash clockwise

9. Decision Signals

Draw a hash ring when membership changes frequently and naive modulo would remap everything:

🎯 Think Consistent Hashing When:
  • You are designing dynamic distributed key-value stores (e.g. Cassandra, DynamoDB).
  • You must scale stateless or stateful websocket gateway arrays behind reverse proxies.
  • You want to dynamically scale distributed cache node clusters without completely invalidating current entries.

11. Deep Dive (Optional)

Why Modulo Hashing Fails First

The naive approach maps key K to server hash(K) mod N. It works until N changes. Add one cache node (N=4 β†’ N=5) and roughly 80% of keys remap to different servers β€” a cache stampede and thundering herd on the database. Remove a failed node and the same catastrophe repeats in reverse.

Consistent hashing exists because elastic infrastructure changes N constantly β€” auto-scaling, rolling deploys, spot instance preemption. With a hash ring, adding or removing one physical node remaps only about 1/N of keys on average. That is the entire reason interviewers expect you to graduate from modulo to a ring when discussing distributed caches or Dynamo-style partitions.

Hash Ring Array Search Mechanics

In production systems (e.g., Libketama, Cassandra), the hash ring is represented in memory as a sorted array of virtual node hashes mapped to physical server addresses. When a request arrives with key K:

  1. Calculate the hash value of the key: h = xxHash(K).
  2. Execute a Binary Search (O(log V) where V = N * vNodes) over the sorted ring array to find the first vNode hash greater than or equal to h.
  3. If no vNode hash is greater than h, wrap around and select the first element in the array (index 0).

Non-Cryptographic Hashing Efficiency

Never use heavy cryptographic hash functions like SHA-256 or MD5 for key routing on the ring. Consistent hashing is CPU bound. Production systems select ultra-fast, high-distribution non-cryptographic hashes like **Murmur3** or **xxHash** to execute routing lookups in under 10 nanoseconds.

**Rendezvous (Highest-Random-Weight) hashing** is an alternative when you want minimal key movement without maintaining a sorted ring: score each server as hash(server, key) and pick the highest score. Lookup is O(N) servers but simpler than vNode ring maintenance β€” common in small-to-medium cache clusters.

πŸ’¬Review

Help Us Improve

How helpful was this walkthrough?

Click a star to rate. We actively use this feedback to refine and update our system design content.

Placeholder
Optional but highly appreciated!

Discussion

Share your thoughts, ask questions, or help others.

Loading comments...