Load Balancing Algorithms Explained: Round Robin to Two Choices
Every system design diagram has a box labelled "load balancer" and most candidates are never asked to defend what is inside it. When an interviewer does ask — "which algorithm, and what happens when one backend gets slow?" — round robin is the reflex answer. Round robin is also the specific algorithm that turns one slow backend into a user-visible outage, which makes it an unusually bad default to have memorised.
This post covers what the balancer can actually see at layer 4 versus layer 7, the five algorithms worth knowing with the arithmetic for when each one breaks, why I default to power-of-two-choices least-request, how health checking causes the cascading failures it was meant to prevent, and the problems no load balancer solves.
What the balancer is actually deciding
One choice — this request goes to that backend — made under three constraints that explain every algorithm below.
The decision budget is microseconds. A proxy handling 100,000 requests per second has about 10 microseconds of CPU per request for everything, routing included. Anything requiring a lock across all backends, a sort, or a network call is out.
It only knows what it can see. The balancer does not know a backend's CPU, its heap pressure, or that it is three seconds into a stop-the-world pause. It knows what it sent and what came back.
There is never one balancer. Production runs M of them behind DNS or anycast, each with its own private view of backend state and no idea what the others just did. An algorithm that is optimal for one balancer can be actively harmful when twelve of them run it simultaneously.
Layer 4 and layer 7: what each one can see
A layer 4 balancer forwards packets. It hashes the connection's 5-tuple (source IP, source port, destination IP, destination port, protocol) to pick a backend and sends every packet of that connection to the same place. It never sees a URL. The payoff is enormous throughput — Google's Maglev paper reports saturating 10 Gbps line rate on commodity hardware, and AWS NLB advertises millions of requests per second with single-digit-millisecond added latency.
A layer 7 proxy terminates the client connection, parses the request, and chooses a backend per request. That costs a TLS handshake, a parse, and typically a few hundred microseconds of added latency, and it makes the proxy a real piece of software you have to version and roll. In exchange it can route by path, retry an idempotent request on another backend, rewrite headers, and — the part that matters most for load balancing — observe the response status and latency of every request it forwards.
| Layer 4 | Layer 7 | |
|---|---|---|
| Decides once per | connection | request |
| Can route on | IP, port | path, header, cookie, method |
| Health signal | TCP connect succeeds | HTTP status, latency, per-route |
| Added latency | tens of microseconds | hundreds of microseconds to ~1 ms |
| Typical | AWS NLB, Maglev, ECMP | Envoy, NGINX, ALB, HAProxy |
The trap worth knowing: layer 4 balancing in front of HTTP/2 or gRPC is close to no balancing at all. gRPC clients open one long-lived connection and multiplex thousands of requests over it. An L4 balancer pins that connection to one backend for its lifetime, so ten clients talking to fifty backends will use ten of them. This is a real production failure, and it is a strong thing to volunteer in an interview. The fixes are an L7 proxy, a client-side balancer, or forcing periodic connection recycling (max_connection_age in gRPC servers).
My default: L7 for service-to-service and user-facing HTTP, L4 only for non-HTTP protocols or when you genuinely need packets-per-second that software proxies cannot reach.
Round robin, and the case it cannot handle
Round robin hands out requests in rotation: backend 1, 2, 3, 1, 2, 3. It is stateless, costs one increment, and distributes request counts perfectly.
Equal counts are not equal work. Round robin is fine when three things hold: request costs have low variance, backends are identically sized, and the per-backend request rate is high enough that variance averages out over a short window. Violate any of them and the algorithm has no feedback loop with which to notice.
Here is the arithmetic. Ten backends, 500 requests per second, a typical request taking 20 ms. Round robin gives each backend 50 rps, and by Little's law each carries 50 x 0.02 = 1 request in flight on average. Now one backend stalls for two seconds — a stop-the-world GC, a lock contention spike, a slow query holding a connection.
Round robin keeps metering 50 rps into that backend for the whole two seconds. It accumulates 100 queued requests, none of which can complete until the stall clears. Over the same window the cluster served 500 x 2 = 1,000 requests, so 10% of all traffic in that window lands in a queue behind a stalled process. Those are your p99 and p99.9.
Weighted round robin fixes only the second of the three conditions. Give a 16-core instance twice the weight of an 8-core one and counts are distributed proportionally. The weights are static configuration, so they are correct on the day you set them and progressively wrong afterwards.
Least outstanding requests
Route to the backend with the fewest requests currently in flight. This is the single most valuable upgrade over round robin, and the reason is that in-flight count is an observed quantity: the balancer increments on dispatch and decrements on response, needing no cooperation from the backend at all.
A rising in-flight count is the one signal that captures every reason a backend might be slow — GC, a cold cache, a noisy neighbour on the host, a slow downstream dependency, a thread pool exhausted by something unrelated. You do not have to enumerate the causes.
Take the same stall. The moment S2's in-flight count exceeds everyone else's, the balancer stops choosing it. It receives perhaps one or two more requests, not 100. The other nine backends now carry 500/9 = 55.6 rps each, an in-flight average of 55.6 x 0.02 = 1.11. That is an 11% load increase on the healthy backends, which no user will notice, versus 100 requests parked behind a dead process.
The subtlety to name before an interviewer does: with M balancers, each one counts only the requests it dispatched. With twelve balancers and ten backends, each balancer's view of a backend's load is roughly one twelfth of reality, and at low per-balancer request rates those counts are mostly noise. Exact least-connections also has a sharper failure, which is the next section.
Random, and why two choices beats one
Pure random is a legitimate baseline: no state, no coordination, O(1), and it degrades gracefully because no two balancers correlate. Its weakness is the tail. Throwing n requests into n backends uniformly at random leaves the busiest backend with roughly log n / log log n in flight. For n = 1,000 that is about 3.6, when the average is 1.
Now sample two backends at random and send the request to whichever has fewer outstanding. The busiest backend's load drops to about log₂ log n, which for n = 1,000 is about 2.8 and, more importantly, grows doubly-logarithmically rather than logarithmically. The constants at realistic cluster sizes are modest; the asymptotics are the famous part, but the operational reason to use it is different and better.
That is the real argument. Exact least-connections has every balancer computing the same global minimum from the same stale snapshot, so they all pick the same backend in the same instant — the idlest host becomes the busiest, and the system oscillates. Sampling two at random decorrelates the decisions while keeping almost all of the benefit. It also costs two lookups instead of a scan of the whole backend set, which matters at a thousand backends.
This is not a research curiosity. Envoy's LEAST_REQUEST policy samples two hosts by default (choice_count: 2). HAProxy ships balance random(2). Finagle has used two-choices with a peak-EWMA load metric for years.
import random
def pick(backends, inflight):
"""Power of two choices: sample two, take the one with fewer in-flight."""
if len(backends) < 2:
return backends[0]
a, b = random.sample(backends, 2)
return a if inflight[a] <= inflight[b] else b
# dispatch
host = pick(healthy_backends, inflight)
inflight[host] += 1
try:
response = send(host, request)
finally:
inflight[host] -= 1 # must run on timeout and error too
The finally is the part that bites people. A decrement skipped on a timeout path leaks in-flight count, the backend's apparent load ratchets upward forever, and the balancer quietly stops routing to a perfectly healthy host.
If you take one position from this post: weighted least-request with two-choices sampling is the right default for HTTP and gRPC service traffic. It needs no backend cooperation, reacts to stalls within a request or two, and does not herd.
Latency-aware routing, and the black hole
Least response time goes one step further: track an exponentially weighted moving average of each backend's latency and prefer the fast ones. Finagle's peak-EWMA variant weights the worst recent observation so a single bad spike is not immediately forgotten.
It has one spectacular failure mode. A backend that fails fast looks fast. A host whose dependency is misconfigured and returns HTTP 500 in 2 ms has the best latency in the fleet, so a latency-aware balancer routes everything to it. This is the black hole, and it is a genuine outage pattern, not a thought experiment.
The fix is that latency must never be the only input. Combine the latency metric with ejection on error rate, and count a 5xx as infinitely slow rather than as 2 ms.
Hash-based routing
Sometimes you do not want even distribution — you want the same request to reach the same backend, so a local cache stays warm or a session stays put. Hash the client IP, a cookie, or a key from the path, and map it onto the backend set.
Plain modulo over the backend count remaps almost everything whenever a backend joins or leaves, which is why production implementations use a ring or a lookup table: Envoy offers RING_HASH and MAGLEV, NGINX has hash ... consistent. The mechanism, the arithmetic for why modulo collapses, and the virtual-node fix are covered in detail in consistent hashing explained.
The cost is that you have given up load balancing in exchange for affinity. Load now follows the key distribution, so one hot tenant or one viral object concentrates on one backend and no algorithm at the balancer can help. Use it when cache locality is worth more than evenness, and say so out loud when you do.
Health checks are the other half of the answer
An algorithm choosing between backends is only as good as the set it chooses from, and this is where candidates who have operated systems separate themselves from candidates who have read about them.
Active checks probe each backend on a timer. They are predictable and they lie: a /healthz that returns 200 whenever the process is up tells you nothing about the dependency it cannot reach. They are also slow — a 5-second interval with a 3-failure threshold means up to 15 seconds of serving errors to real users before ejection.
Passive checks (Envoy calls this outlier detection) use real traffic as the probe. Five consecutive 5xx responses and the host is ejected for 30 seconds, then tentatively returned, with the ejection time growing on repeat offences. This reacts in milliseconds and tests the path users actually take. It is the better primary signal; keep active checks to detect when an ejected host has genuinely recovered.
Two failure modes to name:
Ejection cascades. Load shifts off every ejected host onto the survivors, which makes the survivors slower, which gets them ejected. The cluster unravels from the top. Envoy's defence is a panic threshold: when fewer than 50% of hosts are healthy, it ignores health status entirely and balances across all of them, on the reasoning that a degraded cluster beats a stampede onto the last two machines. Knowing this exists is a strong signal in an interview.
New hosts get flooded. A freshly started backend has an empty in-flight count, so least-request sends it everything — precisely when its cache is cold, its JIT has not warmed, and its connection pool is empty. It immediately becomes the slowest host in the fleet. The fix is a slow-start window that ramps a new host's weight from near zero over 30 to 60 seconds, which both Envoy and NGINX Plus implement. Deploying with least-request and no slow start is a reliable way to make every rolling restart produce a latency spike.
What load balancing does not solve
It does not add capacity. If aggregate demand exceeds aggregate capacity, every algorithm is only choosing whose requests are slow. You need autoscaling, load shedding, or admission control — see design a rate limiter for the mechanics of turning work away deliberately.
It does not stop retry storms. When backends slow down, clients retry, and the retries arrive exactly when there is least capacity to serve them. A balancer faithfully distributing a 3x amplified load is doing its job and making things worse. You need retry budgets (cap retries at a small percentage of request volume), circuit breakers, and jittered exponential backoff.
It does not balance state. If your service is sharded and one shard is hot, nothing the balancer does helps — the requests have to go where the data is. That is a sharding problem, not a routing one.
It does not give fairness between tenants. One client sending 90% of your traffic looks to the balancer exactly like 900 clients sending 0.1% each. Per-tenant quotas live above the routing decision.
It cannot see the future. Every algorithm here is reactive. A request dispatched to a backend one millisecond before its GC pause begins was a perfect decision that produced a terrible outcome. This is why timeouts and hedged requests exist: routing accepts that some choices will be wrong and recovers from them.
Alternatives worth naming
Client-side load balancing. Give the client the backend list from service discovery and let it run two-choices itself. Removes a network hop and an entire tier of infrastructure, and it is what gRPC-LB, Finagle and service meshes do. The cost is that every client language needs a correct implementation of the policy — which is the problem a sidecar proxy exists to solve, at the price of a hop back.
Pull instead of push. Workers take work from a queue when they are free. This is the only scheme that is perfectly balanced by construction, because a worker that is busy does not pull. It applies only to work that tolerates queueing rather than request-response latency, and it comes with its own set of problems around ordering and redelivery — message queues explained covers those. If a system design question is about background jobs, propose this rather than a balancer.
DNS and anycast. Coarse, cheap, global. DNS distributes by handing out different records, and resolvers ignore your TTLs, so convergence after a failure is measured in minutes. Anycast routes to the nearest BGP-announced point of presence and shifts traffic when a site withdraws. Both are the right tool for getting users to a region and the wrong tool for choosing a process.
| Approach | Reacts to a stalled backend in | Main cost |
|---|---|---|
| Round robin | never | no feedback at all |
| Least request + P2C | 1–2 requests | needs in-flight bookkeeping |
| Peak EWMA | a few hundred ms | fast failures look fast |
| Consistent hash | never (by design) | load follows key skew |
| Queue / pull | immediately | only for queueable work |
| DNS | minutes | resolver caching |
Saying this in an interview
When the interviewer points at your load balancer box, the answer that scores is four sentences long:
- Name the layer and why. "L7, because this is HTTP and I want per-request routing and real health signals from status codes. If it were raw TCP throughput I'd use L4."
- Name the algorithm and the failure it avoids. "Least outstanding requests with two-choices sampling. Round robin keeps feeding a stalled backend; in-flight count notices within a request or two."
- Volunteer the distributed caveat. "Each balancer only sees its own in-flight counts, and exact least-connections would make them all herd onto the same idle host, which is why the sampling matters."
- Mention health checking as a separate decision. "Passive outlier detection on consecutive 5xx as the primary signal, active checks to detect recovery, and slow start so a newly deployed host doesn't get flooded while it's cold."
Then stop. The follow-up is almost always "what if one backend starts failing fast?" — the black hole — or "what about sticky sessions?", and having set up the first three points you already have the vocabulary to answer both.
If you want to find out whether that actually comes out cleanly under time pressure, which is a different skill from knowing it, you can run a system design mock interview by voice on Whitepad.
Practice this out loud
Reading is the easy part. Sit across from a senior AI interviewer that talks, watches your whiteboard, and scores you like the real thing — your first mock is free.
Start a free mock →