Imagine a single web server handling every request to your app. It works fine on launch day. Then you get popular, traffic doubles, and that one server starts to sweat — requests queue up, latency climbs, and eventually it falls over. You could buy a beefier machine (vertical scaling), but there's always a ceiling, and a single machine is a single point of failure.
Load balancing takes the other path. Instead of one big server, you run several identical ones and put a traffic cop in front of them.
The problem: one server can't do it all
When all traffic funnels into a single server, two things go wrong as you grow:
- Capacity — one machine can only handle so many concurrent requests before it saturates CPU, memory, or network. And it gets slow well before it's full: once a server is busy most of the time, new requests usually find others ahead of them and wait in line.
- Availability — if that machine restarts or crashes, your entire app is down.
Both problems have the same shape: you've put all your eggs in one basket. Step through a year of growth on one server below. Before traffic passes what the server can finish, predict what happens to page times.
How it works
A load balancer sits between clients and your servers. Clients connect to it — they never address the servers directly. For each incoming request, the balancer picks one server from a pool of identical ones and forwards the request there.
Because the servers are interchangeable, it doesn't matter which one handles any given request. Add three more servers and you've roughly tripled your capacity. This is horizontal scaling: you grow by adding machines, not by enlarging one.
The simplest way to pick is round robin: take turns, 1, 2, 3, 1, 2, 3. That works while every server is equally fast, but servers have bad days: a noisy neighbour on the same host, a long garbage-collection pause, a slow disk. Step through what happens when one of three servers slows down. You'll predict where round robin sends the next request, then flip to least connections, which sends each request to the server with the fewest requests in progress.
Identical servers are the key precondition. Load balancing only works cleanly when any server can handle any request. That usually means servers are stateless — they keep no per-user data in local memory. Anything that must persist (sessions, uploads) lives in a shared store like a database or cache.
One of your four servers has a failing disk, and requests it handles take five times longer. Which balancing rule limits the damage on its own?
Routing around failure
The load balancer continuously runs health checks — small periodic requests to each server ("are you alive?"). When a server fails enough checks in a row, the balancer marks it unhealthy and simply stops sending it traffic. From then on, requests quietly flow to the healthy servers.
Detection isn't instant, though. A balancer that checks every 5 seconds and waits for two misses (so one blip doesn't eject a healthy server) keeps sending traffic to a dead server for up to about 10 seconds. Shorter intervals find failures faster but make more noise; it's a tuning trade-off.
When the sick server recovers and starts passing health checks again, it's added back to the rotation. This is what turns a pile of servers into a resilient system. Predict how long a crashed server keeps getting traffic, then switch the checks off to see the same crash without them.
How does it choose a server?
The routing algorithm decides which server gets each request:
- Round-robin — hand requests out in a cycle: 1, 2, 3, 1, 2, 3… Simple and even when requests cost about the same.
- Least connections — send the next request to whichever server is currently handling the fewest. Better when some requests are much heavier than others, or one server is slower.
- Hashing — derive the server from something stable, like the client's IP, a user ID or a URL. The same input always maps to the same server, which is useful for cache locality: that server's local cache already holds the key's data.
Hashing has a catch. The obvious rule, hash % number of servers, depends on how many servers there are, so changing the count changes the answer for almost every key. Consistent hashing puts keys and servers on a ring instead, so adding or removing a server only moves the keys next to it. Predict how many of 12 keys move when a fourth server joins, first with % and then on a ring.
Sticky sessions pin a given user to the same server for their whole session. It's a quick fix when servers do hold local state — but it undermines even distribution and makes failures more disruptive (a dead server takes its users' sessions with it). Prefer stateless servers with a shared session store instead.
Ten servers each cache data for the users routed to them by hash(user) % 10. During a traffic spike you add an eleventh. Roughly what share of users now land on a server whose cache has never seen them?
When to reach for it
Reach for a load balancer when you need to scale past one machine or you need redundancy so a single failure doesn't take you offline — which, in practice, is almost any production web service. It's one of the most common building blocks in system design, and it pairs naturally with techniques like caching (to cut work per request) and circuit breakers (to handle downstream failures gracefully).
It isn't free. Every request now makes one extra network hop (usually a small fraction of a millisecond inside a data centre), and the balancer itself must not become the new single point of failure, which is why managed load balancers run as several instances behind one address. You also have to keep spare capacity: when one of three servers drops out, the other two must absorb its share.