Lesson 15 / 26

Load Balancing Algorithms and Consistent Hashing

Compare round robin, least load and consistent hashing and know when each fits.

How to pick a server

Round robin rotates through servers: simple, but unfair if requests differ in cost. Weighted variants give bigger servers more traffic. Least connections / least load picks the server doing the least work, which handles uneven request durations better. Random with two choices (pick two at random, use the less loaded) gets most of the benefit cheaply and scales well. Consistent hashing maps a key (user ID, session, cache key) onto a ring of servers so the same key keeps landing on the same server (useful for caches and stateful backends), and adding or removing a server moves only a small fraction of keys, unlike plain hash % N which reshuffles almost everything. Pair any algorithm with health checks so dead servers leave the pool.

Round robin vs least load, run

I ran this plain-Python model. Eight requests where the first takes 10 units and the rest 1 each: round robin gives server a 13 units of work and b 4, while least-load gives 10 and 7, a far fairer split.

import hashlib, bisect, math, random

def round_robin(n_req, servers): return [servers[i % len(servers)] for i in range(n_req)]
def least_conn(durations, servers):
    busy = {s: 0 for s in servers}; picks = []
    for d in durations:
        s = min(servers, key=lambda x: busy[x]); picks.append(s); busy[s] += d
    return picks, busy
dur = [10, 1, 1, 1, 1, 1, 1, 1]
rr = round_robin(8, ["a", "b"]); rr_load = {s: sum(d for d, p in zip(dur, rr) if p == s) for s in "ab"}
lc, lc_load = least_conn(dur, ["a", "b"])
print("round robin load", rr_load, "| least-load load", lc_load)

Output:

round robin load {'a': 13, 'b': 4} | least-load load {'a': 10, 'b': 7}

Consistent hashing vs modulo, run

I ran this plain-Python model. Going from 3 to 4 servers with 1,000 keys: consistent hashing moves 257 keys (close to the ideal 25%), while hash % N moves 770. For a cache, that is the difference between a small dip and a stampede to the database.

import hashlib, bisect, math, random

def h(s): return int(hashlib.md5(s.encode()).hexdigest(), 16)
def ring(nodes, vnodes=100): 
    pts = sorted((h(f"{n}#{i}"), n) for n in nodes for i in range(vnodes)); return [p for p, _ in pts], [n for _, n in pts]
def lookup(r, key):
    keys, owners = r; i = bisect.bisect(keys, h(key)) % len(keys); return owners[i]
keys = [f"user-{i}" for i in range(1000)]
r3 = ring(["n1", "n2", "n3"]); r4 = ring(["n1", "n2", "n3", "n4"])
moved = sum(1 for k in keys if lookup(r3, k) != lookup(r4, k))
naive = sum(1 for k in keys if h(k) % 3 != h(k) % 4)
print("consistent moved", moved, "of 1000 | modulo moved", naive, "of 1000")

Output:

consistent moved 257 of 1000 | modulo moved 770 of 1000

Quick check: Why use consistent hashing for a cache cluster?

  • Adding or removing a node moves only a small share of keys
  • It makes every key unique
  • It removes the need for hashing
  • It encrypts keys
Answer

Adding or removing a node moves only a small share of keys — Limited key movement avoids mass cache misses when the cluster changes.