Load Balancing
One address, many servers: why one server falls over, how a balancer chooses between several, and what it does when one of them is slow, broken or dead.
An interactive System Design lesson: 24 steps, about 30 minutes, on a live simulation in your browser.
The shop's API runs on one machine, web-1. Customers reach it through a front door, lb, which for now has exactly one place to send anything.
Follow one request from client-1. It spends 10 ms crossing the internet to lb, 1 ms to web-1, 20 ms of work on web-1, then the same 11 ms back: 42 ms in all.
What you will learn
One server falls over
- One request, one server: A request's latency is the sum of its hops plus the work at the far end. Only the work part grows when the server is busy.
- Capacity is a number you can compute: A server's capacity is its concurrency divided by the time one request takes: 4 slots at 20 ms is 200 requests a second, and not one more.
- Past capacity: Past capacity, latency climbs until the queue is full and then errors start. A queue buys time for a burst and nothing for a sustained overload.
Many servers, one address
- Three servers behind one address: A load balancer turns many servers into one address. You get their combined capacity, and one of them failing is no longer an outage.
- Round-robin: take turns: Round-robin counts turns, not work. It is fair exactly when requests cost the same and servers are the same.
Who gets the next request
- Break it: one slow server: A slow server under round-robin gets the same share as a fast one, so its queue overflows while the others idle. One sick server sets the fleet's p99.
- Least-connections: Least-connections measures the one thing round-robin ignores: how long each server is taking. Slow servers hold connections longer and are chosen less.
- Drill: least connections in nginx
- Weighted round-robin for unequal servers: A weight is a server's share of the turns. Give servers weights in proportion to their capacity and round-robin is fair again.
Sticky sessions
- ip-hash: the same client, the same server: Stickiness trades balance for locality. Hashing clients spreads load only as evenly as the clients themselves are spread, and a big client lands whole on one server.
- The server holding the carts dies
- Hash the key instead: Hash the client and you get client affinity; hash the key and you get data affinity. Consistent hashing is the version that moves only the departed server's share when the pool changes.
When a server dies
- A server crashes under load: Detection is not instant. Between a crash and the health check that notices it, the dead server's share of traffic fails, and each of those failures costs a full timeout.
- Back up is not back in rotation: A recovered server earns its way back: it re-enters rotation only after several health checks in a row pass.
- Break it: up, but broken: A shallow health check proves the process is reachable, not that it works. A server can pass every check and fail every request.
- A deep health check: A deep check asks "can you do your job?" instead of "are you there?". Make it deep enough to catch a bad deploy and not so deep that one shared dependency takes out every server at once.
- Drill: the health check thresholds
Retries, ejection and the balancer itself
- Retries hide a failure, for a price: A retry turns a failure into latency: the client gets an answer, and pays for the failed attempt in time.
- What about a write?: Retry only what is idempotent. A write that failed may still have happened; make it idempotent with a request id before any layer is allowed to repeat it.
- Ejection: stop asking the broken one: Passive checks judge a server by its real answers and eject it after consecutive errors. They catch the failures that a shallow active check cannot see.
- Break it: retries into an overloaded fleet: Retries help when failures are random and capacity is spare. Under overload they multiply the load that caused the failures.
- Break it: the balancer is one box too: A balancer removes the servers as single points of failure and becomes one itself. It needs its own redundancy: a floating IP, DNS, or anycast.
Recap & playground
- Cheat sheet
- Playground