The Power of Two Choices
Part 2: The real world
What happens to two choices when it meets real servers, real requests, and real failures.
In Part 1, twenty load balancers sending requests to ten servers taught us that checking two random servers beats checking all of them. All along, though, every rule judged servers only by the length of their line and assumed it accurately told us how long a new request would wait.
Jeff Dean and Luiz Barroso, who built much of Google's infrastructure, explained why that trust only goes so far. Checking a few servers and picking the least loaded, they wrote, falls short because
"load levels can change between probe and request time; request service times can be difficult to estimate due to underlying system and hardware variability; and clients can create temporary hot spots by all clients picking the same (least-loaded) server at the same time." (Dean & Barroso, 2013)
So far we've only dealt with the last of these, the herd. Here we'll take on the rest, and then go further: what real systems do when servers, requests, and fixes stop behaving the way we assumed.
1Load is not what you should balance
Real servers have bad moments. One might stall for a few seconds while it cleans up memory, or slow down because another program on the same machine is hogging it. Engineers at Linkerd and Twitter tested what load balancers do when that happens, by making one of eleven servers take 2 seconds per request for 30 seconds (Linkerd, 2016).
Here is a version of their test: one balancer using two choices, live information, and requests that each take a steady 100 ms.
10 servers, 80 requests a second. A stalled server takes 2 seconds per request.
While the balancer counts requests, the stalled server keeps getting work. Its line never looks long: it usually holds one or two requests, about the same as everyone else. But each of them is stuck for seconds. Counting tells the balancer how many requests are waiting, not how fast they're moving.
Watching response times fixes it. Every response comes back through the balancer, so it knows how long each server has been taking. It scores each of its two servers by the requests ahead, times that server's recent time per request (the number under each server). One slow response is enough: the stalled server's score jumps, and the balancer leaves it alone, apart from an occasional request to check whether it has recovered.
Linkerd's real test came out the same way. With a 1-second timeout, sending requests to servers in turn failed about 1 request in 20; counting requests, about 1 in 100; watching response times, about 1 in 1,000.
Learn more: why 1 in 100 vs. 1 in 1,000 matters so much
One slow request in a hundred sounds like a rounding error. But when you load a page, your request rarely touches just one server. Behind the scenes it fans out to dozens or hundreds of them, for your feed, each recommendation, the ads and so on, and the page is only as fast as the slowest of them. Drag the slider to see what that does.
A page is fast only if every server it waits on is fast:
chance a page is slow = 1 − (chance one server is fast)number of servers
With 100 servers that are each slow 1 time in 100, that's 1 − 0.99100 ≈ 1 − 0.37, so about 63% of page loads are slow. At 1 in 1,000, it's about 10%. Even at 1 in 10,000, a page that touches 2,000 servers is slow almost 1 time in 5. That's why the slowest requests, not the average, decide how a large service feels to the people using it. As Jeff Dean and Luiz Barroso put it:
"Even rare performance hiccups affect a significant fraction of all requests in large-scale distributed systems." (Dean & Barroso, 2013)
Google learned the same lesson at YouTube. Its old balancer kept CPU usage almost perfectly even across servers, yet under heavy traffic more than a quarter of requests failed, because some servers shared their machines with other busy programs. Prequal, which replaced it, routes by requests in flight and recent response times, and in the same test nothing failed.
"The real goal of a load balancer is not to balance load, it is to direct load where capacity is available." (Wydrowski, Kleinberg, Rumble & Archer, 2024)
One earlier lesson still applies. If every balancer chases the server with the best response times, they can herd onto it, so real designs like Prequal and C3 deliberately avoid being too greedy.
All of this works because every request here takes about the same 100 ms, so a server's recent response times predict the next one. When some requests take ten times longer than others, they stop predicting much at all. That's where the next chapter begins.
2Stop guessing
Now make the requests uneven, the way they are in most real services: most take about 50 ms, but 1 in 10 takes about half a second, because it's a bigger search or a bigger page. We're back in the realistic setup from the end of Part 1: twenty balancers using two choices, with a load report every 500 ms. Huge requests have a dark ring.
10 servers, 80 requests a second, 20 balancers, a load report every 500 ms. Most requests take about 50 ms; 1 in 10 takes about 500 ms.
Two choices now walks into a trap it can't see. A line with one request in it looks short, but if that request is huge, everyone behind it waits half a second. Nothing in the load report says which lines have a huge request at the front. Watching response times doesn't help either: a server's last few responses say nothing about the request it's working on right now. In our simulation, watching response times actually makes things worse here. As the designers of Sparrow, a scheduler built at Berkeley, put it:
"Queue length provides only a coarse prediction of wait time." (Ousterhout, Wendell, Zaharia & Stoica, 2013)
Join both lines
So stop predicting. With join both lines, a balancer picks two servers at random and puts the request in both lines. Whichever server gets to it first takes it and tells the other server to drop its copy. A request waiting in two lines has a pink ring; when one copy starts, the other disappears. No one has to guess which line is faster, because the lines themselves find out. Here is one request, slowed down:
Requests taking over a second drop from about 17% to about 5%, and the average falls by almost half. And this rule needs no load information at all: there's no report to go stale and nothing for the balancers to herd onto. Google calls these tied requests. In its storage system, they cut the slowest 0.1% of reads by about 40% while adding less than 1% extra work, because the second copy is cancelled before it starts (Dean & Barroso, 2013).
Other systems use the same move in different forms. Sparrow has servers hold a reservation for a task and ask for the task only when they're ready to run it, which brought it within 5% of a perfect scheduler. Microsoft's Join-Idle-Queue has idle servers announce themselves to the balancers, so requests go to servers that are known to be free (Lu et al., 2011). In each case, the decision moves to the last possible moment, when the answer is known instead of guessed.
3Not every server is the same
So far, any server could handle any request equally well. Real servers remember. A server that just fetched someone's profile keeps a copy in its cache, and serving it again is about four times faster than going back to the database. Here each request asks for one of 200 items, a few far more popular than the rest, and each server caches the 10 items it used most recently. Requests for the most popular item have a purple ring.
10 servers, 140 requests a second, one balancer with live information. A cached item takes 25 ms; anything else takes 100 ms. The percentage under each server is how often its cache already had what was asked for.
With two choices, the cluster can't keep up. Spreading requests evenly means every server sees every item, so caches keep getting overwritten and most requests make the slow trip to the database. The load is perfectly balanced, and the lines grow anyway.
The opposite rule sends every request for an item to the same server, its home. The standard way to pick homes is consistent hashing, invented in 1997 for caching web pages and the basis of Akamai (Karger et al., 1997): servers and items sit at random points on a circle, and an item's home is the next server clockwise. Now the caches work beautifully. But follow the purple rings. Every request for the most popular item lands on one server, and the same happens with the next few. Their homes can't keep up, while other servers sit idle. It's Aisle 7 again, except this time the crowd never moves on.
Home, unless busy
The fix, from researchers at Google, is to go home unless home is already busy (consistent hashing with bounded loads). If the home server is holding more than 1.25 times the average number of requests, the request skips it and goes to the next server around the circle instead (Mirrokni, Thorup & Zadimoghaddam, 2018). Popular items spill onto a neighbor or two, which start caching them too. It's the only one of the three rules that keeps up here. Vimeo uses it, with that same 1.25, for its video servers.
The same tension now sits at the center of serving AI models, where a GPU that has already read a long prompt can reuse that work instead of redoing it (Srivatsa et al., 2024). Spread those requests evenly and the work is thrown away; send them all to the GPU that already has the prompt and it becomes the next Aisle 7. Mitzenmacher saw a version of this coming. He ended his 2001 paper on two choices by looking past a world where every server is interchangeable:
"Perhaps the most interesting open question is to include locality in this model." (Mitzenmacher, 2001)
4When the cure becomes the disease
So far, every problem has started with the servers. This one starts with the clients. When a request takes too long, a client usually gives up and sends it again. That's a sensible habit: most slow requests are bad luck, and a second try usually lands somewhere faster. Here, each client gives up after 1 second and retries once. Here is what that looks like for one request:
The servers are only 70% busy, with plenty of room to spare. But what happens when there's a hiccup?
10 servers, 70 requests a second, one balancer using two choices. A hiccup slows every server to a quarter of its speed for 3 seconds.
During the hiccup, lines grow, requests start taking over a second, and clients send them again. The retries pile more work onto servers that are already behind, so more requests time out, and more retries go out. Worse, the servers still finish every request whose client has already given up, the faded faces, doing work no one is waiting for. When the hiccup ends, the servers are back at full speed, and it doesn't matter. Twice as many requests arrive as before, nearly all of them time out, and the cluster stays down.
Nathan Bronson and his colleagues named this a metastable failure: the system stays broken after whatever broke it is gone. Their example is nearly this one, a database whose clients retry once after 1 second, knocked over by a 10-second network outage:
"The system is now in a metastable failure state. It will remain there until the load is significantly reduced or the retry policy is changed." (Bronson, Aghayev, Charapko & Zhu, 2021)
It isn't rare. A follow-up study found that at least 4 of the 15 major outages at Amazon Web Services over a decade were metastable failures (Huang et al., 2022). The hiccup was only the trigger. What keeps the system down is the retries, the very feature meant to hide small failures.
Retry within a budget
So one fix is to limit them. With a retry budget, a client only retries while retries make up less than a tenth of what it sends; past that, it reports the failure instead. Google's guidance for its own services recommends the same kind of cap (Google SRE, 2016). Some requests fail during the hiccup, but afterward the cluster recovers within seconds, because retries can never add more than a sliver of extra work.
The same guidance adds one more rule, and it's an old friend. Clients that all retry on the same schedule turn into a herd, all rushing the servers at once, and the fix is the same as before: randomness.
"If retries aren't randomly distributed over the retry window, a small perturbation (e.g., a network blip) can cause retry ripples to schedule at the same time, which can then amplify themselves." (Google SRE, 2016)
5Where things stand
Every fix in this part works, and each one runs in real systems today. None of them is the last word.
There's still no universal measure of "busy." Each system tunes its own signals, and a balancer that chases the fastest server can start a herd of its own (Suresh et al., 2015). Joining two lines costs extra work, which is cheap only while there's capacity to spare, and extra work is exactly what turned a hiccup into an outage. Bounded loads handles caches, but AI serving raises the stakes: reusing a GPU's earlier work saves far more than a cache hit does, and how to weigh that against load is still an open question. And retry budgets defuse the failure loops engineers already know about. As Bronson and his colleagues put it:
"A systematic approach for building systems that are robust against unknown metastable failures remains an open problem." (Bronson et al., 2021)
Learn more: a stutter in AI serving
Reusing prompts isn't AI serving's only new problem. Serving a chatbot takes two kinds of work. Reading a new prompt is one big burst. Writing answers is many small steps: each step adds one word to every answer the GPU is working on. Here is a pair of GPUs that are already writing six answers when long new prompts arrive, handled three ways:
Mixing gets new chats started fastest, but every answer on that GPU freezes while a prompt is read, and the freezes stack up when prompts arrive together. Cutting prompts into chunks keeps answers flowing, but new chats wait longer for their first word. Separate GPUs avoid both, but the reading GPU sits idle between prompts, and when many arrive at once it becomes the bottleneck. The field is split. Some systems give prompt-reading its own GPUs (Zhong et al., 2024), as Moonshot AI does in production for its Kimi chatbot (Qin et al., 2024); others read prompts in chunks between answer words (Agrawal et al., 2024). The papers explicitly leave the head-to-head comparison open, and it remains an active area of research.
Under all of these sits one tradeoff every large service has to make: how close to full to run its servers. Busy servers are cheaper, but everything on these pages gets harder as servers fill up. Lines grow faster, herds hurt more, and retries tip systems over. Many services run near the edge anyway, because the savings are too large to pass up:
"Many production systems choose to run in the vulnerable state all the time because it has much higher efficiency than the stable state." (Bronson et al., 2021)
Better load balancing mostly buys the right to run closer to that edge. Prequal's biggest win wasn't just faster responses. It let YouTube and other services at Google run their servers much busier than before (Wydrowski et al., 2024).
This series has followed one question: which server should get the next request? Real load balancers juggle much more: checking that servers are still alive, steering around ones that fail, spreading traffic across data centers, deciding when to add servers. Two choices is one small part of a much bigger field.
Below is a sandbox with everything from this part, to try these problems together. Or skip ahead to the Epilogue: Beyond servers, where we follow the same idea into hash tables, routers, and caches.
6Sandbox
Here is everything from this part in one place. Each chapter changed one thing at a time; here you can combine them. Stall a server while clients retry, turn on caches and cause a hiccup, or put twenty balancers on stale reports and make some requests huge.
Some things to try:
- Why retry at all? Stall server 4 with two choices and clients that don't retry: about 1 request in 100 fails. Switch to retry within a budget and failures nearly vanish, because the second try lands on a healthy server. That's what retries are for. Now cause a hiccup with retry right away, and you'll see what they cost.
- Find the edge. With retry right away, lower the traffic to 45 requests a second and cause a hiccup. The cluster recovers: even with every request sent twice, the servers can keep up. At 60 requests a second, the same hiccup takes it down for good. Bronson and his colleagues call the region in between vulnerable: perfectly healthy, until something goes wrong.
- Retries make hot spots worse. Turn on caches, choose always home, and raise traffic to 120 requests a second. The most popular item's server falls behind, its requests time out, and with retry right away every timeout sends it more work. Switch to home, unless busy and the spillover keeps things moving.