HotShard
Traffic routing

Load Balancer

Spread requests across servers, and pick an algorithm that routes around a slow one.

A load balancer sits in front of a pool of identical servers and decides which one handles each request. The interesting part is how it decides. Simple rotation ignores load, so requests pile up behind a slow server. Load-aware policies route around it. And one trick, power-of-two-choices, gets most of that benefit for almost no extra work.

~5 min read

The problem: spread work, and route around a slow server#

TL;DRthe 30-second version
  • A load balancer fronts a pool of interchangeable backends and answers one question per request: which server handles this?
  • Round-robin is load-blind. It keeps handing a slow server its turn, so that server's queue grows while the others sit idle.
  • Least-connections and power-of-two-choices look at live load. Power-of-two-choices samples two backends at random and picks the less loaded one, which is nearly as good as checking all of them, at constant cost.
  • A real balancer also runs health checks, caps retries, and is itself made redundant so it isn't a single point of failure.

One server can only handle so much traffic. Past that ceiling, latency climbs and requests start failing. The fix is to run many identical servers and put a balancer in front to spread work across them. Now you grow capacity by adding machines.

But servers aren't always equal. One might be a bigger machine. One might be garbage-collecting. One might be stuck on a slow dependency, and one might have just crashed. A scheme that just rotates through them ignores all of that. It keeps handing requests to the struggling server while healthy ones sit idle, and keeps sending traffic to a dead one until something notices.

So the balancer has two jobs. Detect and avoid backends that are slow or dead. And for each request, pick the server that keeps load even, using only cheap information, because this decision runs on every single request.

How it works: one decision per request#

The balancer keeps a list of backends and, for load-aware policies, a little live state about each one: how many requests are in flight, recent latency, and whether it's healthy. Every request goes through the same loop.

  1. A request arrives at the balancer, which owns the service's public address.
  2. Drop any backend the health checks currently mark unhealthy.
  3. Run the selection algorithm over the healthy set.
  4. Forward the request and bump that backend's in-flight counter.
  5. On response or timeout, decrement the counter and record the latency. If the request is safe to repeat, it may be retried on another backend.
clientsmany, concurrent
requests
load balancerpick a backend per request
route by algorithm+ health checks
backend poolS1 ok Β· S2 ok Β· S3 ok Β· S4 slow
The balancer sits in front of the pool and chooses per request

The selection algorithm is where balancers differ. The options form a ladder, from load-blind to load-aware to hash-based.

  • Round-robin: hand each request to the next server in rotation. Simple and stateless, but load-blind. A slow server keeps getting its turn and backs up.
  • Weighted round-robin: rotate in proportion to each server's capacity, so a server twice as big gets twice the requests. Handles mixed hardware, still ignores live load.
  • Least-connections: send each request to the backend with the fewest in-flight requests. A slow server's count stays high, so it stops being picked. But you have to scan the whole pool on every pick.
  • Power-of-two-choices (P2C): pick two backends at random and send to the less loaded one. Almost as balanced as least-connections, but you only look at two servers.
  • Least-response-time (EWMA): keep a running average of each backend's latency and route to the fastest. Catches a server that answers, but slowly. More state and more tuning. Envoy and Finagle use this.
  • Consistent hashing (Maglev): hash a key such as the session id or cache key to a backend, so the same key always lands on the same server. That's what sticky sessions and cache affinity need. Adding or removing a backend remaps only a small share of keys. Maglev is Google's version, tuned for an even spread.
PredictYou run 4 servers. S4 is slow and its in-flight requests are piling up. Under plain round-robin, does S4 stop getting new requests until it catches up?

Hint: Round-robin counts turns, not load.

No. Round-robin is load-blind. It keeps handing S4 its turn no matter how backed up it is, so the queue keeps growing. Least-connections or power-of-two-choices look at live load and route around S4.

Why does looking at just two servers work so well? Model requests as 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. Now, for each ball, sample two bins and put it in the emptier one. The fullest bin drops to about log log n. That's an exponential improvement from one extra comparison. Sampling three or more instead of two only changes the constant, so real systems sample exactly two.

There's a second reason to prefer it over always picking the global minimum. At scale you have many balancer instances, each working from a slightly stale view of load. If they all pick "the least-loaded backend", they all pick the same one and stampede it. P2C's randomness breaks that synchronization while still steering away from hotspots.

Where the balancer runs is a separate choice. A dedicated box in the path (nginx, HAProxy, AWS ALB) is one place to configure and watch, but it's an extra hop and a thing to keep available. Client-side balancing (gRPC, Finagle) puts the pick inside each client using a service-discovery feed. No extra hop, but every client needs the logic and a fresh view of the pool.

Cost per decision#

The selection function runs on every request, so its cost is a real budget.

  • Round-robin and weighted round-robin: O(1) per pick, just advance a pointer. No per-backend state.
  • Least-connections: O(N) per pick to scan for the minimum, or O(log N) with a heap, plus an in-flight counter per backend. The scan and the shared counters get expensive with big pools and many balancer instances.
  • Power-of-two-choices: O(1) per pick, two random indexes and one comparison, with the same per-backend counter. Near least-connections balance at constant cost.
  • Least-response-time: O(1) to O(N) depending on the implementation, plus a latency average updated on every completion.
  • Consistent hashing: O(1) lookup into a precomputed table. The cost moves to rebuilding the table when the pool changes.
Trade-offs: balance vs state vs coordination

Every step up the ladder buys better balance with more state and more coordination. Round-robin needs nothing shared. Least-connections needs an accurate count of in-flight requests per backend, which is hard to keep consistent when ten balancer instances are each making picks. P2C is popular because approximate, per-instance state is good enough for it.

  • Smarter balancing costs state: load-aware policies update per-backend metrics on every completion. That's memory, contention on shared counters, and a feedback loop that can oscillate if it reacts too fast.
  • Global optimum invites herding: picking the single least-loaded backend looks ideal but synchronizes many balancers onto the same target. Randomized sampling trades a little optimality for stability.
  • Affinity fights even load: hashing a key to a backend means you no longer get to balance freely. A few hot keys can overload their backend no matter how idle the rest of the pool is.
  • L4 vs L7: a layer-4 balancer routes by IP and port. It's fast and protocol-agnostic, but it balances whole connections and can't see URLs. A layer-7 balancer understands HTTP, so it can route by path, header, or cookie, retry single requests, and terminate TLS, at a bit more CPU per request. Large deployments often stack both: L4 in front for raw connection spreading, L7 behind it for smart routing.
Algorithms compared
AlgorithmBalance qualityPer-pick costState neededAffinity
Round-robinPoor (load-blind)O(1)None (a pointer)No
Weighted round-robinFair (static weights)O(1)Per-backend weightNo
Least-connectionsExcellentO(N) scanIn-flight count per backendNo
Power-of-two-choicesNear-excellentO(1)In-flight count per backendNo
Least-response-time (EWMA)Excellent (latency-aware)O(1)–O(N)Latency average per backendNo
Consistent hashing / MaglevFair (key-dependent)O(1) lookupPrecomputed hash tableYes
Where load balancers run in the wild
  • nginx: the L7 reverse proxy most engineers meet first. Round-robin, least_conn, ip_hash, and in Nginx Plus a P2C 'random two' method.
  • HAProxy: L4 and L7. roundrobin, leastconn, source and URI hashing, random with draws=2 for P2C, plus detailed health checks.
  • Envoy: the L7 proxy behind service meshes like Istio. Its default is weighted least-request, an O(1) P2C policy. Also ring-hash, Maglev, and EWMA, with active health checks and passive outlier detection built in.
  • AWS ALB and NLB: managed balancers. ALB is L7 (path and host routing, TLS, least-outstanding-requests). NLB is L4 (very high throughput, static IPs).
  • Google Maglev: a software L4 balancer using consistent hashing for connection affinity with minimal disruption on changes. It underpins Google Cloud's global load balancing behind a single anycast address.
  • Client-side: gRPC, Finagle, and Netflix Ribbon push the pick into the client, using a service-discovery feed instead of a middlebox.
Failure modes and how to handle them
  • The balancer is a single point of failure. Every request flows through it, so if it dies the service dies. Run several balancer instances and fail over between them: an active/passive pair sharing a virtual IP (VRRP, keepalived), an anycast IP advertised from many machines, or DNS returning several balancer addresses.
  • Health-check flapping. A backend hovering at the edge of its check threshold flips healthy, unhealthy, healthy, and the pool churns with every flip. Require several consecutive failures to eject and several successes to re-add. Ramp a re-added backend up slowly so it isn't flooded at once.
  • Retry storms. When a backend slows, retries pile on top of the original load and can knock over the rest of the pool. Cap retries with a retry budget (retries as a small share of total traffic), add jittered backoff, and shed load instead of queueing it without bound.
  • Sticky-session imbalance. Pinning clients to backends means load follows the keys, not the algorithm. A few heavy sessions overload one backend, and removing it drops all its sessions at once. Prefer stateless backends with shared session storage. If you must pin, use consistent hashing so removals remap as little as possible.
  • Stale pool view. A balancer with an out-of-date view sends traffic to a dead backend, or herds onto one with many instances. Propagate health fast, use a randomized policy, and add outlier detection that ejects a backend returning errors even if it passes the basic health check.
Active vs passive health checksActive checks probe each backend on a schedule, say GET /healthz every few seconds. Clear signal, but extra traffic and a few seconds of detection lag. Passive checks (outlier detection) watch real request results and eject a backend that starts erroring or timing out. Instant and free, but they only see backends currently receiving traffic. Production systems run both.
In an interview

Name the ladder and one weakness per rung: round-robin is load-blind, weighted handles mixed capacity, least-connections is load-aware but scans the pool, P2C gets nearly the same balance at O(1), EWMA sees slow-but-alive servers, consistent hashing is for affinity. If you land one crisp fact, make it P2C: two random samples cut the worst-case load from log n / log log n to log log n, and the randomness avoids the herd. Then expect the operational follow-ups below.

What does power-of-two-choices actually buy you?

Sampling two random backends and picking the less loaded one drops the worst-case server from about log n / log log n requests to about log log n, for one extra comparison. That's near least-connections balance at O(1) cost. And because it's randomized, many balancers don't all stampede the same 'least-loaded' server.

Why not always send to the single least-loaded server?

Two reasons. Finding the global minimum means scanning every backend and needing an accurate, shared view of all their loads, which is hard with many balancer instances. And it synchronizes them: each acts on a slightly stale snapshot, they all pick the same server, and they stampede it.

How is the load balancer itself made highly available?

Make the balancing tier redundant. An active/passive pair sharing a virtual IP that fails over (VRRP, keepalived), an anycast IP advertised from many machines so routing picks a live one, or DNS handing out several balancer addresses.

What do L4 and L7 mean?

Layers of the network stack. An L4 balancer routes by IP address and port. It's fast and protocol-agnostic, and it balances whole connections. An L7 balancer understands HTTP, so it can route by path, header, or cookie, retry individual requests, and terminate TLS, for a bit more work per request.

Do sticky sessions break load balancing?

They constrain it. Pinning a client to one backend means load follows the keys, so a few heavy sessions can overload their server while others idle, and losing that backend drops all its sessions. Prefer stateless backends with shared session storage. If you must pin, use consistent hashing so backend changes remap as few sessions as possible.

References & further reading
References

Feedback on this topic β†’