# Consistent Hashing — System Design Interview Prep

Source: https://www.skillbyai.com/en/system-design-interview/b-hash

> Add servers without reshuffling everything.

## A ring with virtual nodes

Assigning keys with `hash(key) % N` moves almost every key when N changes, causing a storm of cache misses or data moves. **Consistent hashing** places servers and keys on a ring; each key belongs to the next server clockwise, so adding a server moves only about 1/N of keys. **Virtual nodes** (many points per server) smooth out the load. It is used in distributed caches, Dynamo-style databases and partitioned queues.

## Modulo versus a ring when adding a fifth server, run

I ran this with Python 3.12.3 using only the standard library; inputs are fixed or seeded, so the output is reproducible. With modulo hashing about 80% of 10,000 keys move; with a 100-virtual-node ring about 20% move, which is the minimum possible, and load stays roughly balanced.

```python
# Modulo hashing vs consistent hashing when adding a 5th server
import bisect, hashlib

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

keys = [f"user:{i}" for i in range(10_000)]

def modulo(n):
    return {k: h(k) % n for k in keys}

class Ring:
    def __init__(self, nodes, vnodes=100):
        self.points = sorted((h(f"{n}#{v}"), n) for n in nodes for v in range(vnodes))
        self.hashes = [p for p, _ in self.points]
    def node(self, key):
        i = bisect.bisect(self.hashes, h(key)) % len(self.points)
        return self.points[i][1]

before, after = modulo(4), modulo(5)
moved = sum(before[k] != after[k] for k in keys)
print(f"modulo 4 -> 5 servers : {moved / len(keys):.1%} of keys move")

r4, r5 = Ring(["s0", "s1", "s2", "s3"]), Ring(["s0", "s1", "s2", "s3", "s4"])
moved = sum(r4.node(k) != r5.node(k) for k in keys)
print(f"ring   4 -> 5 servers : {moved / len(keys):.1%} of keys move (ideal 20%)")
load = {}
for k in keys:
    load[r5.node(k)] = load.get(r5.node(k), 0) + 1
print("keys per server:", dict(sorted(load.items())))
```

Output:

```
modulo 4 -> 5 servers : 79.6% of keys move
ring   4 -> 5 servers : 19.9% of keys move (ideal 20%)
keys per server: {'s0': 1818, 's1': 2069, 's2': 2006, 's3': 2122, 's4': 1985}
```

## Mention virtual nodes

Without them, a few servers can own much larger arcs of the ring and become hot.

**Quiz:** Roughly what fraction of keys moves when going from 4 to 5 servers with consistent hashing?

- [ ] All of them
- [ ] About 80%
- [x] About 20%
- [ ] None

*Answer:* About 20%. Only the new server's share moves.
