System Design (High-Level Design) Module 4 – Networking for system designers
Load-balancing algorithms, from round robin up
Round robin, least connections, latency-aware and two-random-choices balancing, and hashing for cache locality, raced on one fleet with a slow server.
What you will learn
- Compare round robin, least connections, latency-aware and random-two-choices policies
- Use consistent hashing when the same key should keep reaching the same server
- Measure the effect of each policy on tail latency in a simulation you can change
Before you start
On this page
A load-balancing algorithm decides which healthy server gets the next request. The simplest ones ignore load: round robin takes the servers in turn and random picks any of them. The load-aware ones look before they send: least connections picks the server with the fewest requests in progress, latency-aware policies prefer servers that have been answering quickly, and the power of two random choices compares just two servers picked at random and sends to the less busy one. Hashing policies send the same key to the same server, so that its cache stays warm. Which one is right depends on three questions: do requests cost about the same, do the servers run at the same speed, and how fresh is the balancer’s view of the load? This lesson races six policies on the same requests, on a fleet with one slow server and on a mixed fleet, so you can see each answer in numbers.
The policies, one by one
Round robin sends each request to the next server in a fixed rotation. It is the default of many balancers,
Amazon’s Application Load Balancer among them, and it is the right choice when requests cost about the same and the
servers are alike (ALB routing algorithms).
Equal counts are not equal work, though. Google’s SRE book reports requests whose CPU cost differs by a factor of
1,000 or more, and round robin leaving up to twice as much CPU use on the busiest task as on the idlest
(SRE book, chapter 20). Weighted round robin gives
bigger servers more turns: with weight=3 on one of three nginx servers, it gets 3 of every 5 requests
(nginx load balancing).
Random picks a server uniformly at random. On average it spreads load like round robin, with more variation from moment to moment; Envoy notes that it generally beats round robin when no health checks are configured, because it has no bias towards the server that follows a failed one in the rotation (Envoy load balancers).
Least connections, called least requests or least outstanding requests on HTTP balancers, sends each request to
the server with the fewest in progress. It reacts to slow servers without measuring their speed: by Little’s law a
server that takes longer holds more requests at once, so it looks busier and gets fewer new ones. HAProxy recommends
leastconn for long sessions such as LDAP and SQL and says it suits short HTTP sessions less well; Amazon suggests
least outstanding requests when requests vary in cost or targets vary in capacity
(HAProxy manual, balance). Its trap is a server that fails fast:
it answers errors so quickly that it has almost nothing in progress, looks like the least loaded server and pulls in
more traffic, which the SRE book calls sinkholing.
Latency-aware policies keep a moving average of each server’s response times, usually an exponentially weighted moving average (EWMA), and favour servers that answer quickly. Linkerd uses EWMA to send requests to the fastest endpoints, and Finagle offers a “peak EWMA” that reacts sharply to slow responses and weighs the average by each server’s outstanding requests (Linkerd, Finagle). Used alone, a latency average herds traffic, as the race below shows.
The power of two random choices samples two servers at random and sends the request to the one with fewer
requests in progress. It needs no sorted list and no global view, and it avoids the herd behaviour of always picking
the single least loaded server. It is a default in several proxies: Envoy’s least request balancer compares 2 random
hosts by default when their weights are equal, HAProxy’s balance random makes 2 draws by default, nginx offers
random two least_conn, and Finagle’s default balancer is P2C, short for power of two choices
(Envoy,
HAProxy manual,
nginx upstream module).
Hashing turns a key of the request (the client’s address, a user ID, the URL) into a server, so that the same key keeps reaching the same server while the set of servers stays the same. It trades even load for locality, and the section on hashing below covers how to keep that trade cheap when servers come and go.
In short, policy by policy:
- Round robin looks at nothing but its rotation. A slow server keeps its full share, many balancers can share a fleet without harm, and it suits requests and servers that are alike.
- Least connections looks at requests in progress. A slow server gets less at once, but stale counts make many balancers herd, so it suits one balancer with fresh counts, and long sessions.
- Latency (EWMA) looks at recent response times. A slow server gets less once its average rises; used alone it herds, so it is combined with load, for servers of different speeds.
- Two random choices looks at requests in progress on two random servers. A slow server gets less, it stays robust with many balancers or old counts, and it is the usual default for proxies and clients.
- Consistent hashing looks at a key of the request. A slow server keeps its keys, and it suits caches and per-key state.
The race: six policies, one slow server
This simulation sends the same 36,079 requests to 10 servers under each policy. Each server works on one request at a time, most requests take 4 or 10 ms, and one in 200 takes 100 ms. In fleet 1, server 8 is three times slower than the rest, as a server with a noisy neighbour or a failing disk might be; in fleet 2, half the servers are twice as fast as the other half. The last two lines repeat fleet 1 with load counts that are 100 ms old, which is what each balancer sees when many balancers share a fleet and learn the load from periodic reports. Change the numbers at the top and run it again.
// A load-balancer race: six policies send the same requests to the same fleet of 10 servers.
// Each server works on one request at a time, in arrival order. Change FLEETS, RATE or STALE_MS and run it again.
const RATE = 900; // requests a second, arriving at random (about Poisson)
const SECONDS = 40;
const STALE_MS = 100; // the last part: load reports this old, as many separate balancers would see them
// Work per request in ms on a normal server: 70 % take 4 ms, 25 % 10 ms, 4.5 % 25 ms, 0.5 % 100 ms.
const COSTS = [[0.7, 4], [0.95, 10], [0.995, 25], [1, 100]];
const FLEETS = [
{ name: "Fleet 1: server 8 is 3 times slower", speed: [1, 1, 1, 1, 1, 1, 1, 3, 1, 1], watch: [7], label: "to 8" },
{ name: "Fleet 2: servers 1-5 are twice as fast", speed: [0.5, 0.5, 0.5, 0.5, 0.5, 1, 1, 1, 1, 1], watch: [0, 1, 2, 3, 4], label: "to 1-5" },
];
function random(seed) {
// mulberry32: a small seeded generator, so every run prints the same numbers
let a = seed >>> 0;
return () => {
a = (a + 0x6d2b79f5) >>> 0;
let t = Math.imul(a ^ (a >>> 15), a | 1);
t ^= t + Math.imul(t ^ (t >>> 7), t | 61);
return ((t ^ (t >>> 14)) >>> 0) / 4294967296;
};
}
function workload() {
const r = random(42);
const requests = [];
for (let tick = 0; tick < SECONDS * 10000; tick++) {
if (r() < RATE / 10000) {
const u = r();
requests.push({ at: tick / 10, cost: COSTS.find(([upTo]) => u < upTo)[1] });
}
}
return requests;
}
const POLICIES = {
"round robin": (s) => s.turn++ % s.n,
random: (s) => Math.floor(s.r() * s.n),
"least connections": (s) => {
let best = Infinity;
let ties = [];
for (let i = 0; i < s.n; i++) {
if (s.load[i] < best) [best, ties] = [s.load[i], [i]];
else if (s.load[i] === best) ties.push(i);
}
return ties[Math.floor(s.r() * ties.length)];
},
"two random choices": (s) => {
const [a, b] = twoDifferent(s);
return s.load[b] < s.load[a] ? b : a;
},
"lowest latency": (s) => s.latency.indexOf(Math.min(...s.latency)), // the best moving average
"latency x load": (s) => {
// two random choices again, comparing average latency times (requests in progress + 1)
const [a, b] = twoDifferent(s);
return s.latency[b] * (s.load[b] + 1) < s.latency[a] * (s.load[a] + 1) ? b : a;
},
};
function twoDifferent(s) {
const a = Math.floor(s.r() * s.n);
const b = Math.floor(s.r() * (s.n - 1));
return [a, b >= a ? b + 1 : b];
}
function race(requests, speed, policy, staleMs = 0) {
const n = speed.length;
const s = { n, r: random(7), turn: 0, load: new Array(n).fill(0), latency: new Array(n).fill(10) };
const done = speed.map(() => []); // finish times of each server's requests, in order
const next = new Array(n).fill(0); // the first of them not finished yet
const freeAt = new Array(n).fill(0);
const sent = new Array(n).fill(0);
const waits = [];
let reportAt = 0;
for (const q of requests) {
for (let i = 0; i < n; i++) {
while (next[i] < done[i].length && done[i][next[i]].end <= q.at) {
const d = done[i][next[i]++];
s.latency[i] += 0.2 * (d.end - d.at - s.latency[i]); // moving average of response times
}
if (!staleMs) s.load[i] = done[i].length - next[i]; // requests in progress right now
}
if (staleMs && q.at >= reportAt) {
for (let i = 0; i < n; i++) s.load[i] = done[i].length - next[i]; // a fresh report, then none for a while
reportAt = (Math.floor(q.at / staleMs) + 1) * staleMs;
}
const i = POLICIES[policy](s);
const end = Math.max(q.at, freeAt[i]) + q.cost * speed[i];
freeAt[i] = end;
done[i].push({ at: q.at, end });
sent[i]++;
waits.push(end - q.at);
}
waits.sort((x, y) => x - y);
const at = (p) => waits[Math.floor(p * waits.length)];
return { p50: at(0.5), p99: at(0.99), sent };
}
const ms = (x) => (x < 100 ? x.toFixed(1) : String(Math.round(x))).padStart(7);
const requests = workload();
console.log(`10 servers, ${RATE} requests a second for ${SECONDS} s`);
console.log(`(${requests.length} requests, the same for every policy)`);
for (const fleet of FLEETS) {
console.log(`\n${fleet.name}`);
console.log(`${"policy".padEnd(19)}${"p50 ms".padStart(7)}${"p99 ms".padStart(7)}${fleet.label.padStart(7)}`);
for (const policy of Object.keys(POLICIES)) {
const { p50, p99, sent } = race(requests, fleet.speed, policy);
const share = (100 * fleet.watch.reduce((sum, i) => sum + sent[i], 0)) / requests.length;
console.log(`${policy.padEnd(19)}${ms(p50)}${ms(p99)}${`${share.toFixed(1)}%`.padStart(7)}`);
}
}
console.log(`\nFleet 1, load reports ${STALE_MS} ms old`);
for (const policy of ["least connections", "two random choices"]) {
const { p50, p99 } = race(requests, FLEETS[0].speed, policy, STALE_MS);
console.log(`${policy.padEnd(19)}${ms(p50)}${ms(p99)}`);
} Output
10 servers, 900 requests a second for 40 s (36079 requests, the same for every policy) Fleet 1: server 8 is 3 times slower policy p50 ms p99 ms to 8 round robin 8.6 30145 10.0% random 11.6 32365 9.8% least connections 4.0 60.3 4.5% two random choices 7.6 100 4.8% lowest latency 7831 29773 2.2% latency x load 10.0 102 2.2% Fleet 2: servers 1-5 are twice as fast policy p50 ms p99 ms to 1-5 round robin 4.0 85.8 50.0% random 5.0 111 50.1% least connections 4.0 25.0 58.5% two random choices 4.0 34.6 55.7% lowest latency 3271 16956 65.1% latency x load 4.0 49.3 67.1% Fleet 1, load reports 100 ms old least connections 31.1 471 two random choices 13.5 182
Recorded with Node.js 24.21.0 on macOS 26 arm64. To run it yourself: mise exec node@24.21.0 -- node lb_race.mjs
Runs on this device, in your browser. The first run downloads JavaScript (about 0.6 MB), which is kept for the next runs.
Your run, in this browser
Read fleet 1 first. Round robin and random give the slow server its 10 %, about 90 requests a second, while it can finish only about 48, so its queue grows for as long as the run lasts and the p99 reaches 30 s. Both load-aware policies cut its share to under 5 % and keep the p99 near 60 to 100 ms. Least connections wins here, because one balancer with exact counts can always find an idle server; two random choices comes close with far less information. “Lowest latency”, which always picks the best moving average, is the warning: an average changes only when requests finish, so every new request goes to whichever server looked fastest until its queue has grown, then to the next one, and the median wait reaches 7.8 s. Comparing latency times load between two random servers (“latency x load”) fixes that, and in fleet 2 it sends the fast servers 67.1 % of the requests, close to their two thirds of the capacity.
The last part is the surprise. With counts 100 ms old, least connections sends every request between two reports to the same few servers that looked idle, and its p99 rises to 471 ms; two random choices spreads the same requests over random pairs and stays at 182 ms. Mitzenmacher showed this in queueing models: with old load information, sending tasks to the apparently least loaded server can significantly hurt, while the less loaded of two random servers stays effective across a wide range of settings (How Useful Is Old Information?). Each of many balancers that counts only its own requests is in a similar position, and that is how sidecar proxies usually see the load.
Exercise · Easy · Python
Write the choices of four balancing policies
Write the decision each policy makes, in policies.py. Servers are numbered from 0, and in_flight lists how many requests each server has in progress.
- round_robin(n, turn) returns the server for request number turn (counting from 0) when there are n servers. - least_loaded(in_flight, start=0) returns the server with the fewest requests in flight. Look at the servers in order from start, wrapping around after the last one, and on a tie keep the first one you found, so that different values of start spread ties over different servers. - two_choices(in_flight, a, b) returns whichever of servers a and b has fewer requests in flight, and a on a tie (the two servers were already picked at random). - update_average(average, sample, weight=0.2) returns a moving average of response times after one more sample: the old average moved weight of the way towards the sample.
For example, least_loaded([3, 1, 4, 1], start=2) is 3, and update_average(10, 20) is 12.0.
Starter code · policies.py
def round_robin(n, turn):
"""The server for request number `turn` of n servers."""
# Replace this line with your code.
return 0
def least_loaded(in_flight, start=0):
"""The server with the fewest requests in flight; ties go to the first found from `start`."""
# Replace this line with your code.
return 0
def two_choices(in_flight, a, b):
"""The less loaded of servers a and b; a on a tie."""
# Replace this line with your code.
return 0
def update_average(average, sample, weight=0.2):
"""The moving average after one more sample."""
# Replace this line with your code.
return 0 The sample tests · test_policies.py
import math
from policies import least_loaded, round_robin, two_choices, update_average
def test_round_robin():
"""takes the servers in turn and starts again after the last"""
assert [round_robin(3, t) for t in range(7)] == [0, 1, 2, 0, 1, 2, 0]
assert round_robin(10, 25) == 5
def test_least_loaded():
"""finds the server with the fewest requests in flight"""
assert least_loaded([5, 2, 7, 3]) == 1
assert least_loaded([0, 4, 4]) == 0
assert least_loaded([9, 8, 7, 6]) == 3
def test_least_loaded_ties():
"""breaks ties by the first server found from start, wrapping around"""
assert least_loaded([3, 1, 4, 1]) == 1
assert least_loaded([3, 1, 4, 1], start=2) == 3
assert least_loaded([3, 1, 4, 1], start=3) == 3
assert least_loaded([2, 2, 2], start=1) == 1
assert least_loaded([1, 5, 1], start=2) == 2
def test_two_choices():
"""keeps the less loaded of the two, and the first on a tie"""
assert two_choices([4, 1, 6], 0, 1) == 1
assert two_choices([4, 1, 6], 2, 0) == 0
assert two_choices([4, 4, 6], 1, 0) == 1
def test_update_average():
"""moves the average part of the way towards each sample"""
assert math.isclose(update_average(10, 20), 12.0)
assert math.isclose(update_average(10, 20, weight=0.5), 15.0)
assert math.isclose(update_average(50, 10, weight=0.1), 46.0) A hint
turn % n is round robin. For least_loaded, walk k from 0 to len(in_flight) - 1 and look at server (start + k) % len(in_flight); replace the best so far only when a server has strictly fewer requests in flight. The moving average is average + weight * (sample - average).
Results of the sample tests
| Test | Result | Details |
|---|
What your code printed
The sample tests run on this device, in your browser (Pyodide): nothing is sent to mysmartcopilot.com. The first run downloads Python (about 13.5 MB), which is kept for the next runs. A check in your browser is feedback for you, not proof that the code is right for every input.
Why two choices are enough
The classic model is balls thrown into bins. Throw n balls into n bins at random and the fullest bin ends up with about log n / log log n balls; let each ball look at d ≥ 2 random bins and take the emptiest, and the fullest bin holds about log log n / log d plus a constant (Mitzenmacher, Richa and Sitaraman). The jump from one choice to two is exponential; a third choice only improves things by a constant factor, which Mitzenmacher also found for queues where tasks wait for service (The Power of Two Choices). Here is that model run:
"""Balls into bins: the fullest bin with one random choice, two choices and three choices.
Each ball goes to one random bin, or to the emptier of d bins chosen at random. Five seeded runs each.
"""
import random
BINS = 1_000
def fullest(balls, choices, seed):
rng = random.Random(seed)
load = [0] * BINS
for _ in range(balls):
best = rng.randrange(BINS)
for _ in range(choices - 1):
other = rng.randrange(BINS)
if load[other] < load[best]:
best = other
load[best] += 1
return max(load)
for balls in (1_000, 100_000):
average = balls // BINS
print(f"{balls:,} balls in {BINS:,} bins ({average} per bin)")
print(f"{'':10}{'fullest bin, 5 runs':>20}{'over avg':>9}")
for choices in (1, 2, 3):
runs = [fullest(balls, choices, seed) for seed in range(5)]
label = "1 choice" if choices == 1 else f"{choices} choices"
print(f"{label:10}{''.join(f'{r:>4}' for r in runs)}{max(runs) - average:>9}") Output
1,000 balls in 1,000 bins (1 per bin)
fullest bin, 5 runs over avg
1 choice 5 5 6 6 5 5
2 choices 3 3 3 3 3 2
3 choices 2 2 2 2 2 1
100,000 balls in 1,000 bins (100 per bin)
fullest bin, 5 runs over avg
1 choice 138 142 135 129 132 42
2 choices 102 102 102 102 102 2
3 choices 101 102 102 102 101 2
Recorded with Python 3.14.8 on macOS 26 arm64. To run it yourself: mise exec python@3.14.8 -- python3 p2c.py
Runs on this device, in your browser. The first run downloads Python (about 13.5 MB), which is kept for the next runs.
Your run, in this browser
With 1,000 balls in 1,000 bins, one random choice left a bin with 5 or 6 balls and two choices kept every bin at 3 or fewer. With 100 balls per bin on average, one choice left the fullest bin 29 to 42 balls above the average, while two choices kept it 2 above, in every run. A third choice made almost no difference. For a load balancer, the bins are servers and the balls are requests: one extra look at the load is what removes the hot spots.
Hashing when the same key should reach the same server
Hashing is the policy to use when serving a key from the same server is worth more than perfectly even load: a cache
tier, where a second server would miss; per-user state that is expensive to rebuild; or a shard. HAProxy’s balance uri, for example, hashes the URL so that proxy caches see each object on one server and keep their hit rate high
(HAProxy manual).
The obvious way, hash(key) mod N, breaks when N changes. HAProxy’s manual says of its map-based hashing that most
mappings change when the server count changes, and nginx warns that adding or removing a server may remap most keys
(nginx upstream module). For a cache tier, that is a
mass miss at the moment you added capacity. Consistent hashing places servers and keys on a ring and gives each key
to the next server point after it, so a server that joins or leaves moves only the keys next to its points: about
1/N of them, as Envoy puts it for its ring hash. Try it on the ring below, and compare the keys that move with hash mod
n:
Add and remove servers, change the number of virtual nodes, and see which keys move. Each key belongs to the first server point clockwise from it.
At the start: 3 servers with 4 virtual nodes each hold 200 keys: A 69 (35%), B 41 (21%), C 90 (45%).
Maglev hashing, from Google’s software network load balancer, fills a fixed lookup table instead of a ring. By Envoy’s figures it builds its table about 10 times faster and picks hosts about 5 times faster than a large ring, at the cost of moving about twice as many keys when hosts change (Envoy load balancers, Maglev).
Hashing ignores load, so one hot key can overload its server. Bounded-load consistent hashing caps every server
at a multiple of the average load and, when a request’s first-choice server is full, picks another from the same
hash.
HAProxy’s hash-balance-factor does this: at 150, no server takes more than 1.5 times the average, and the manual
calls 125 to 200 reasonable values (HAProxy manual).
Client-side balancing
Everything above can run in a proxy or in the client itself. With client-side balancing each client gets the list of servers from service discovery and picks one per request, as gRPC clients can. gRPC’s guide weighs it against a proxy: no extra hop and better performance, but a more complex client, a separate implementation for every language, and clients that must be trusted, or a separate “lookaside” balancer that tells clients which servers to use (gRPC Load Balancing). Because every client sees only its own requests, client-side balancers are exactly the case where two random choices beats picking the single least loaded server.
JavaScript & TypeScript Online Compiler Paste the race above, change the fleet or the policies, and run it with your own numbers.Choosing a policy
Choosing a load-balancing policy in three questions
Text description of the diagram
The figure has three questions down the left, each leading either to the next question or to a policy on the right.
- Must a key keep its server, as with a cache, a session or a shard? If yes, hash the key with consistent or Maglev hashing, with a cap on any one server's load.
- If not: do request costs or server speeds differ? If no, use round robin, weighted if the servers differ in size.
- If they do: does one balancer see fresh counts of the requests in flight on every server? If yes, use least connections (least requests). If no, because there are many balancers or the counts are old or partial, use two random choices, comparing requests in flight or latency times load.
- Requests and servers alike, and a stable fleet: round robin, weighted if servers differ in size.
- Costs or speeds differ, one balancer with exact counts: least connections (least requests).
- Many balancers, client-side balancing or old counts: two random choices, on requests in flight, or on latency times load when servers differ in speed.
- The same key should reach the same server: consistent or Maglev hashing, with a load cap for hot keys.
- Whatever you pick: count errors as load, or a fast-failing server will attract traffic, and keep health checks and slow start from the previous lesson.
Interview questions
Warm-up (fresher to mid level): what is the power of two random choices in load balancing, and why is it popular? The balancer picks two servers at random and sends the request to the one with fewer requests in progress, or the lower latency-times-load score. In the balls-and-bins model it removes almost all of the imbalance that pure random choice leaves: the fullest of n bins drops from about log n / log log n balls to about log log n / log 2, and a third choice adds little. It is popular because it costs two lookups instead of a scan of every server, needs no global view, and stays good when load counts are old or each balancer sees only its own requests, where picking the single least loaded server sends everyone to the same place. Envoy, HAProxy, nginx and Finagle all offer it, and several use it by default.
Key takeaways
- Round robin and random give a slow server its full share: in the race its queue grew without limit and the p99 hit 30 s.
- Least connections reacts to slow servers through the requests they hold; with one balancer and exact counts it gave the best tail, 60.3 ms.
- Latency averages used alone herd traffic; combined with load and two random choices they adapt to faster servers.
- With old or partial counts, two random choices (p99 182 ms) beat least connections (471 ms).
- Two choices instead of one cut the fullest bin from 42 above the average to 2.
- Hash on a key only when locality matters, use consistent or Maglev hashing so scaling moves few keys, and cap the load so a hot key cannot take one server down.
Check yourself
6 questions about this lesson. Every answer and why it is right is on the page, behind “Show the answer”. Your score stays in this browser.
References
- The Power of Two Choices in Randomized Load Balancing (IEEE TPDS, 2001) (IEEE (author's copy, Michael Mitzenmacher))
- The Power of Two Random Choices: A Survey of Techniques and Results (Handbook of Randomized Computing (authors' copy, Mitzenmacher, Richa and Sitaraman))
- How Useful Is Old Information? (IEEE TPDS, 2000) (IEEE (author's copy, Michael Mitzenmacher))
- Supported load balancers (Envoy documentation) (Envoy Project)
- Maglev: A Fast and Reliable Software Network Load Balancer (Google Research (USENIX NSDI 2016))
- HAProxy 3.2 Configuration Manual (balance, hash-type, hash-balance-factor) (HAProxy Technologies)
- Module ngx_http_upstream_module (nginx)
- Using nginx as HTTP load balancer (nginx)
- Edit target group attributes for your Application Load Balancer (routing algorithms) (Amazon Web Services)
- Site Reliability Engineering, chapter 20: Load Balancing in the Datacenter (Google (O'Reilly Media))
- Load balancing (Linkerd documentation) (Linkerd Authors)
- Finagle clients (load balancing) (Finagle project)
- gRPC Load Balancing (gRPC Authors)
Related tools
Report a problem with this lesson
Kept only in this browser. Your Learn progress