MentorNode
Start free
Foundations & Core PatternsMediumdesign-consistent-hashing

Design a Consistent-Hashing Shard Router

Design the routing layer that maps billions of keys onto a changing set of storage nodes, so that adding or losing a node moves the minimum possible data.

Consistent HashingVirtual NodesRebalancingHot Shard Mitigation
Traffic & Capacity Estimates:

500 storage nodes · 10B keys · node add/remove should move ~1/N of keys

Functional Requirements

  • •Deterministically map any key to a storage node from any client, with no central lookup.
  • •Adding or removing a node relocates only the keys it owns, not the whole keyspace.
  • •Replicate each key to the next R distinct physical nodes on the ring.
  • •Expose a membership view so clients learn about topology changes within seconds.

Non-Functional Requirements

  • •Key-to-node resolution in microseconds, computed locally on the client.
  • •Load imbalance across nodes under 10% despite non-uniform key popularity.
  • •Rebalancing must be incremental and throttled so it never saturates the network.

Back-of-the-Envelope Math

  • 500 physical nodes * 200 virtual nodes = 100,000 ring positions.
  • Removing one node relocates ~0.2% of the keyspace instead of ~100% with modulo hashing.

Key Architectural Trade-offs

  • Consistent hashing with virtual nodes vs rendezvous (HRW) hashing vs static range partitioning — rebalancing cost against routing simplicity.
  • Client-side ring computation (fast, but every client must learn membership) vs a routing proxy tier (one more hop, one more thing to scale).
  • Hot key on a single shard: split it with a per-key salt, or replicate it to every node and accept the write amplification.

Click or drag a component onto the canvas, then connect the handles to draw the data flow.

3 nodes · 2 edges

Components · 35

Client & Edge4
Compute & Gateway7
Storage & Caching11
Messaging & Streaming6
Coordination & Ops5
Intelligence2
Canvas overview