Search

How Consistent Hashing Distributes Data Across Servers

The short answer

Quick answer: Consistent hashing is a way of assigning keys to servers so that adding or removing a server moves only a small fraction of the keys. Both servers and keys are hashed onto the same circular number line, the hash ring. Each key belongs to the first server found moving clockwise from the key's position. When a server joins or leaves, only the keys in the arc next to it change owner. With N servers, about 1/N of the keys move, instead of nearly all of them, as happens with the simple hash(key) % N approach.

The problem with hash mod N

Suppose you spread cache keys across 4 servers:

server = hash(key) % 4

It distributes keys evenly, and every client computes the same answer. Now add a fifth server and the formula becomes hash(key) % 5.

Key hash% 4% 5Moved?
1020Yes
1131Yes
1202Yes
1313Yes
2000No

Going from 4 to 5 servers reassigns about 80% of all keys. For a cache, that means nearly every lookup misses at once and the database behind it is flooded. For a storage system, it means copying most of the data across the network. And servers change all the time: scaling up, scaling down, crashing.

This is the same weakness noted in sharding explained.

The hash ring

Consistent hashing, introduced by David Karger and colleagues at MIT in 1997 for distributed web caching, changes what is hashed.

  1. Take the output range of a hash function, say 0 to 2^32 - 1, and imagine it bent into a circle.
  2. Hash each server (by name or address) to a point on the circle.
  3. Hash each key to a point on the same circle.
  4. A key belongs to the first server clockwise from its position.
            S1 (hash 10)
          /              \
   S4 (hash 80)        S2 (hash 35)
          \              /
            S3 (hash 60)

key at 20  -> S2      key at 50  -> S3
key at 70  -> S4      key at 90  -> wraps round to S1

Adding a server

Add S5 at position 45, between S2 and S3. Keys between 35 and 45 used to belong to S3; now they belong to S5. Nothing else moves.

Removing a server

If S2 fails, its keys (those between 10 and 35) pass to the next server clockwise, S3. Again, nothing else changes.

In both cases the disruption is limited to one arc of the ring, roughly 1/N of the keys. The Wikipedia article gives the formal analysis.

The problem with the basic ring

With only a few points on the ring, placement is lumpy:

  • Uneven load. Random positions leave some servers with long arcs and others with short ones. One server might own 40% of the keys and another 10%.
  • Unfair failover. When a server dies, its entire load lands on a single neighbour, which may then be overloaded too.

Virtual nodes

The fix is to place each physical server on the ring many times, at different positions, by hashing server1#1, server1#2, and so on. Each position is a virtual node.

With 100 to 200 virtual nodes per server:

  • Load evens out. Each server owns many small arcs, and the totals average out.
  • Failures spread. A dead server's arcs are scattered round the ring, so its load is shared by many servers instead of one.
  • Capacity can be weighted. Give a machine with twice the capacity twice as many virtual nodes.

The cost is a bigger lookup table. Finding the owner of a key is a binary search over the sorted list of positions, which is fast.

Replication on the ring

Storage systems keep each key on several servers. A simple rule: store the key on the next N distinct physical servers clockwise from its position. If one fails, the others still have the data.

This is the design described in the Amazon Dynamo paper, which brought consistent hashing into mainstream database design. See how database replication works for the quorum rules that go with it.

Where it is used

  • Distributed caches. Memcached client libraries use it so that adding a cache server does not wipe out the whole cache. See caching strategies.
  • Distributed databases. Cassandra, DynamoDB and Riak partition data on a ring.
  • CDNs. Mapping content to cache servers was the original use; Akamai was founded on it.
  • Load balancers. Routing each user or URL to the same backend, with minimal disruption when backends change. See what a load balancer does.
  • Sharded real-time services. Assigning chat rooms or users to servers.

Variations

TechniqueIdeaNotes
Rendezvous (highest random weight) hashingFor each key, score every server with hash(key, server) and pick the highestNo ring to maintain; simple; work per lookup grows with the number of servers
Jump consistent hashA tiny algorithm that maps a key to one of N numbered bucketsVery fast and even; buckets can only be added or removed at the end
Maglev hashingA precomputed lookup tableUsed in Google's network load balancer
Bounded-load consistent hashingCap each server's load and overflow to the nextPrevents a hot server

A different approach to the same problem is to use a fixed number of slots. Redis Cluster hashes each key into one of 16,384 slots and assigns slots to nodes. Rebalancing means moving whole slots, tracked in an explicit table, rather than recomputing positions on a ring.

What it does not solve

  • Hot keys. If one key is extremely popular, consistent hashing sends all of its traffic to one server. Spreading keys evenly does not spread requests evenly.
  • Moving the data. It tells you which keys must move, but something still has to copy them while the system stays online.
  • Agreement on membership. Every client must use the same list of servers. If they disagree, they route the same key to different places. Membership is usually shared through a coordination service or a gossip protocol.
  • Range queries. Hashing scatters neighbouring keys, so scanning a range means asking every server.

A sketch in code

import bisect, hashlib

def h(value):
    return int(hashlib.md5(value.encode()).hexdigest(), 16)

class Ring:
    def __init__(self, servers, vnodes=100):
        self.ring = sorted(
            (h(f"{s}#{i}"), s) for s in servers for i in range(vnodes)
        )
        self.points = [p for p, _ in self.ring]

    def lookup(self, key):
        i = bisect.bisect(self.points, h(key)) % len(self.ring)
        return self.ring[i][1]

The hash function here only needs to spread values evenly; it is not being used for security. The general idea of hashing is covered in how hash maps work.

Frequently asked questions

What problem does consistent hashing solve?

It lets you add or remove servers in a hashed, distributed system while moving only a small share of the keys, instead of almost all of them.

What is a virtual node?

One of many positions on the hash ring assigned to a single physical server. Virtual nodes even out load and spread the effect of a failure.

How many keys move when a server is added?

About 1/N of them, where N is the new number of servers.

Is consistent hashing the same as sharding?

It is one method of deciding which shard holds which key. Sharding is the broader practice of splitting data across servers.

Conclusion

Consistent hashing replaces "divide by the number of servers", which breaks whenever that number changes, with "walk round the ring to the next server", which barely notices. Add virtual nodes for balance and replicas for safety, and you have the partitioning scheme behind many distributed caches and databases.

Related articles

Sources and further reading

Usama Muneer

Usama Muneer

Coder, Blogger, Tech Speaker & Web Technologies Enthusiast. Passionate about working on open-source Programming languages & Tools while utilizing my Product Development skills.

Your experience on this site will be improved by allowing cookies Cookie Policy