The Power of Two Choices

How to spread work across servers, and why a coin flip helps.

Every time you load a video or run a search, your request lands on one of many identical servers. A load balancer decides which one, and a large service makes that decision millions of times a second.

The obvious rule is to send each request to the least busy server. The rule inside YouTube, NGINX and Envoy is stranger: pick two servers at random, and send the request to the less busy of the two.

Here are two clusters handling the same traffic. The left one sends each request to a random server. The right one picks the better of two.

Requests on the right finish about three times faster. That part is not too surprising: looking at two servers is better than looking at none. But why only two? To find out, we'll start with the simplest rule there is, picking at random.

1Picking at random

Let's look at one cluster up close. The balancer is at the top and ten servers are along the bottom. Each face is a request: the one sitting on a server is being handled, and the ones stacked above it are waiting. A line flashes each time the balancer sends a request, and faces get grumpier the longer they wait.

Each request takes about 100 ms, so the ten servers can handle about 100 requests a second. Here 90 arrive each second, keeping the servers busy 90% of the time.

This balancer sends every request to a random server. Watch for servers that turn orange. That's a server sitting idle while requests at other servers wait in line: there is work to do and a free server to do it, but they're in different places.

The simulation runs ten times slower than real life so you can follow individual requests. Use the speed buttons to fast-forward.

Random is fair on average: over time, every server gets the same share of requests. But in any moment some servers get unlucky streaks, and when the servers are busy, unlucky streaks are expensive. Try the slider. At 50 requests a second, a request spends about 200 ms in the system in total. At 90 a second it's about 1,000 ms, and at 95 a second about 2,000 ms. Those numbers come from a classic formula from queueing theory:

average time in the system = time per request1 − how busy the server is

At 90% busy, that's 100 ms ÷ 0.1 = 1,000 ms. As a server gets close to fully busy, the bottom of the fraction gets close to zero and the wait shoots up. A server that's busy almost all the time has little slack to catch up after a burst, so any line it builds takes a long time to drain. A little unfairness turns into a lot of waiting.

2Looking at every server

The obvious fix is to check every server before sending a request, and pick the one with the shortest line. Engineers call this least loaded, or "least connections." Here is the same traffic, 90 requests a second. Switch between the two rules.

With least loaded, the lines nearly vanish, the average drops from about 1,000 ms to under 200 ms, and no server ever sits idle while others have a line, because an idle server is always the first choice.

Learn more: what about round robin?

Most real load balancers have a simpler default: take turns. Send the first request to server 1, the next to server 2, and so on, then start over. This is called round robin. It needs no information about the servers at all, and it spreads requests perfectly evenly. Try it first with requests that are all the same size.

10 servers, 90 requests a second, one balancer.

When every request takes the same time, round robin is about as good as least loaded. Now switch to varied sizes, which is much closer to the real world, where one request might be a quick lookup and the next a big search. Round robin falls far behind. It can't tell that a server just drew a big request, so it keeps sending that server its turn while the line behind the big request grows. Least loaded sees the line and steers around it.

With one load balancer in front of ten servers, this is close to the best you can do. Real systems are not built that way, though.

3Many load balancers

A large service gets far too many requests for one load balancer to handle, and a single balancer would be a single point of failure. So there are many of them, dozens or hundreds, each taking a share of the incoming requests and deciding where to send them on its own.

Each balancer knows what it has sent. It doesn't know what the other balancers have sent. To find out how busy the servers really are, it would have to ask them, and asking every server before every request would add a network round-trip to every request. So a common design is for a monitor to collect every server's queue length and send a load report to all the balancers at once, here every 500 ms. Between reports, each balancer adds its own requests to the report's numbers.

In the picture, the gold mark in each server's line shows what the last report said, and the bar at the top shows how old the report is. The small bars under the servers count how many requests each server has received since the last report. Every balancer below uses least loaded. Start with one balancer, then add more.

With one balancer it works well. With twenty, it falls apart. Every balancer reads the same report, sees the same server as least loaded, and sends its next requests there. Each one counts its own requests, but none of them can see that the other nineteen are doing the same thing. Watch the bars under the servers: right after a report, nearly all the new requests go to one or two servers. By the next report those servers have long lines and different ones look emptiest, so the whole crowd moves there. With twenty balancers, least loaded does no better than picking at random.

This is called herd behavior. Michael Mitzenmacher, who analyzed it in 2000, compared it to a supermarket announcing that "Aisle 7 is now open": everyone rushes over, and Aisle 7 is soon the longest line (Mitzenmacher, 2000). Real systems run into it. The engineers who built TranSend, an early web proxy at Berkeley, "noticed rapid oscillations in queue lengths" because their balancers acted on stale load reports (Fox et al., 1997).

4Two choices

Here is the trick. For each request, the balancer picks two servers at random, compares their lines in the report, and sends the request to the shorter one, the way you might glance at two checkout lines and join the shorter. The setup is the same as before: twenty balancers and a load report every 500 ms. Switch between the three rules, and use slow motion to watch each request check its two servers (dashed lines) and pick one (solid line).

With two choices, the average comes back down to about 450 ms, less than half of what least loaded manages with twenty balancers.

Look at the bars under the servers. Under least loaded, a few servers get almost everything after each report. Under two choices, new requests spread out over most of the servers. That spread comes from the randomness. Every balancer still reads the same report, but each request compares a different random pair of servers, so the balancers no longer all reach the same answer. The server that looks emptiest only gets a request when it happens to be in that request's pair, which is a few times its fair share instead of all of it.

Is checking more servers better?

It's natural to guess that if checking two servers helps, checking three helps more, and checking all ten helps most. But we've already seen that least loaded (the equivalent of checking all ten) falls apart when balancers share an old report. So where is the sweet spot? The chart shows the average response time for each number of servers checked, with twenty balancers. Checking one server is random, at the left end, and checking all ten is least loaded, at the right end. Start at 500 ms reports, then try the other report ages.

Each point is a long simulation: 10 servers, 90 requests a second, 20 balancers. Hover for the numbers.

With 500 ms reports, the best choice is two. Checking three is already worse, and checking all ten is no better than random. Now switch to live: with perfectly fresh information, checking more keeps helping, and the best choice is to check them all. In between, the sweet spot slides: about five servers at 100 ms, three at 250 ms, and two from 500 ms on.

Why would more information ever hurt? Each extra check gives a request more information, but it also makes it more likely that every balancer picks the same server. With fresh information that's fine. With a stale report, that agreement is what creates the herd. So there's a balance between picking the shortest line and picking a line nobody else is about to pick, and as the report ages, the balance tips toward fewer checks.

Learn more: is it really the old information?

Here is the herd setup again, twenty balancers using least loaded with 500 ms reports, with one change. Instead of one monitor sending the same report to every balancer at the same moment, each balancer gets its own report on its own schedule: still every 500 ms, still just as old, but at a different moment from the others. The small gold bar under each balancer shows how old its report is. With one shared report, all twenty fill up in lockstep. With staggered reports, they're out of step. Switch between the two.

10 servers, 90 requests a second, 20 balancers, a report every 500 ms. The bars under the servers count the requests each one received in the last 500 ms.

With staggered reports, the herd mostly disappears: least loaded drops from about 1,000 ms to about 400 ms, as good as two choices. Yet nothing got fresher; each report is still up to half a second old. What changed is that the balancers no longer read the same numbers at the same moment. When one balancer's report shows a server as empty, another balancer's report, taken a moment later, already shows requests piling up there, so they make different choices.

So the problem was never really the age of the information. It was many decision-makers acting on the same information at the same moment, without seeing each other's choices. Mitzenmacher saw this in 2000, when he gave each request information of a slightly different age:

"Tasks that enter at almost the same time may have different views of the system and, thus, make different choices. Hence, the 'herd behavior' is mitigated, improving the load balancing." (Mitzenmacher, 2000)

Now switch the rule to two choices and flip between shared and staggered. It barely changes, because two choices already makes the balancers disagree: every request compares its own random pair. Staggering adds randomness to when a balancer learns things; two choices adds randomness to which servers it compares. Either one breaks up the herd.

Average response time as the reports get older. Simulated: 10 servers, 90 requests a second, 20 balancers. Hover for the numbers.

So why not just stagger the reports? Because in real systems, synchronization creeps in by accident: a control plane pushes new numbers to every balancer at once, health checks run on the same timer, a deploy restarts everything together. And even perfectly fresh information fails if the decisions happen at the same instant. Mitzenmacher showed in his 1996 thesis that with a single round of simultaneous choices, the longest line grows just as fast as with random routing (Mitzenmacher, 1996). Two choices doesn't depend on any of that. Its randomness is built into every decision, which is a big part of why it's the default in real load balancers.

Two holds up across the whole range. The next section shows why: the first extra check is worth far more than all the others.

5But why two?

Let's go back to a single balancer with live information, so nothing goes stale. Picking at random averages about 1,000 ms. Checking every server averages about 190 ms. Checking just two averages about 300 ms: most of the way from no information to perfect information, from one extra look. Why is a single extra look worth so much?

Bad luck has to strike twice

Say 1 server in 10 has a long line right now. A request sent to a random server lands in one of those lines 1 time in 10. A request that checks two servers lands in one only if both of the servers it checked had long lines, since otherwise it takes the shorter. That happens 1 time in 10 × 10, or 1 in 100. One extra look turns a 1-in-10 risk into a 1-in-100 risk.

And it compounds

The same thing happens at every line length. For a line to grow to 3, a request had to find two servers that both already had at least 2. Lines of 2 were already fairly rare, so lines of 3 are much rarer, and lines of 4 rarer still. At each step up, the share of servers with a line that long is roughly squared, and squaring a fraction makes it smaller, fast: a tenth becomes a hundredth, then a ten-thousandth. Here is how many servers out of a million have a line at least a given length, at 90% busy:

Line lengthRandomTwo choices
1 or more900,000900,000
2 or more810,000729,000
3 or more729,000478,297
4 or more656,100205,891
5 or more590,49038,152
6 or more531,4411,310
7 or more478,297about 2

From the formula Mitzenmacher (1996) and Vvedenskaya et al. (1996) derived for a very large cluster; the simulations on this page match it.

With random routing, almost half of all servers have a line of 7 or more. With two choices, about two servers in a million do. Random routing trims the count by the same 10% at every step. Two choices squares it.

How long can a line get?

Squaring also explains the most famous form of this result. Take n servers and drop n requests onto them, one after another, with nothing finishing in between. Then look at the longest line. Here it is for clusters from 10 servers to 10 million:

Average longest line when n requests are dropped onto n servers. Simulated; hover for the numbers.

With random routing, the longest line keeps growing as the cluster grows: about 5 at a thousand servers, about 9 at a million. With two choices it barely moves: 3 at a thousand, 4 at a million. That pattern holds in general. With two choices, you have to square the number of servers to make the longest line one request longer: a thousand servers to a million, a million to a trillion.

The share of servers with longer and longer lines goes roughly ½, ¼, 1/16, 1/256, 1/65,536: the number of digits roughly doubles at every step. Computer scientists write the length of the longest line as log log n, which is 5 or less for any cluster that has ever existed. With random routing it grows like log n / log log n, which keeps climbing. In Azar, Broder, Karlin and Upfal's own experiments with about 16 million servers, the longest random line was usually 10; with two choices, it was 4 in every run.

"This apparently minor change in the random allocation process results in an exponential decrease in the maximum occupancy per location." (Azar, Broder, Karlin & Upfal, 1994)

Why not three?

Three checks would cube instead of square, and the chart shows that it helps: the longest line drops from 4 to 3. But that is small next to the jump from one check to two. The first extra look changes how long lines shrink, from a little at each step to squaring at each step. Each look after that raises the power by one, to cubing and then to the fourth power, but by then lines are already disappearing so fast that each extra look saves only a step here and there. And each extra look also makes the balancers more likely to agree on the same server, which is what sets off a herd when their information is stale. Two gets almost all of the benefit, at the smallest possible cost.

"The gain from each ball having two choices is dramatic, but more choices prove much less significant: the law of diminishing returns appears remarkably strong in this scenario!" (Mitzenmacher, 1996)

The formulas behind these results show the same pattern far beyond any real cluster:

The leading terms of the theorems, plotted out to 10100 servers (Azar, Broder, Karlin & Upfal, 1994). They show how fast each grows, not exact line lengths: at real cluster sizes, the simulated lines above are longer.

6Where this is today

The power of two random choices was worked out in the 1990s. Azar, Broder, Karlin and Upfal proved it for balls thrown into bins (1994), and Mitzenmacher (1996) and, independently, Vvedenskaya, Dobrushin and Karpelevich (1996) worked it out for servers with lines like the ones above. As early as 1986, Eager, Lazowska and Zahorjan had found by simulation that "extremely simple" policies using "very small amounts" of information did nearly as well as complicated ones (1986).

Today the idea is built into NGINX, HAProxy and Envoy, three of the most widely used load balancers. In 2024 Google described Prequal, the load balancer behind most of YouTube's serving stack: it checks a few random servers for each request and picks the best one, and it cut tail latency in half (Wydrowski et al., 2024).

Real systems raise harder questions:

The lesson reaches well beyond load balancers. When many decision-makers act on the same information at the same moment, more information can make things worse. A little information plus a little randomness often does better.

Below is a sandbox with everything from this page. Or skip ahead to Part 2: The real world, where we take on the harder questions that real systems raise.

7Sandbox

Here is everything from this page in one place, plus a few things we haven't covered. Change the rule, the number of servers and balancers, how the balancers get their information, and what the requests look like, and watch what happens.

Some things to try:

Continue to Part 2: The real world →