Load Balancing: The Intuition Behind the Power of Two Random Choices
betterprogramming.pub
betterprogramming.pub
In the balls and bucket scenario, ideally you'd want to scan all the buckets and put the ball in the least filled one.
By picking two random buckets and putting the ball in the least filled of the two, you approximate the ideal algorithm.
The chance of picking two relatively filled buckets is inversely proportional to how "bad" it is. If there are just two buckets which have more balls than the others then it's relatively bad to pick those two, but chances are low. And vice versa.
Still, hadn't thought about this before, interesting trick indeed!
Feels like picking an adjacent bucket would work just as well and increase locality.
There's a difference, but it's minimal. Only the initial choice has to be random. As long as you have more than one bin to pick from, you are good regardless of your strategy (random vs adjacent).
https://medium.com/@mihsathe/load-balancing-a-very-counterin...
- It's O(1), rather than O(buckets).
- It's robust to stale data and concurrency. One way systems try to avoid the O(buckets) work is by caching, or by scanning in parallel, either of which cause stale load estimates. "Pick the best" tends to lead to significant over-filling http://www.eecs.harvard.edu/~michaelm/postscripts/handbook20... or https://brooker.co.za/blog/2012/01/17/two-random.html).
- Best-of-k (rather than best-of-two) provides a parameter `k` that allows the system to tradeoff between optimality (higher k is more optimal), and robustness to stale data and concurrent accesses (lower k is more robust). Mitzenmacher's core observation is that the step up from k=1 to k=2 is a massive improvement in optimality with only a small loss in robustness (and much more than most people would expect).
With random placement, as was said in the article, each flip is independent. With PoTRC it’s using knowledge of current queue depth.
The author didn’t explain the challenges around tracking queue depth, or why PoTRC is better than say keeping an ordered list of queues depth and picking the smallest one.
I’m assuming the trade off is about cost of managing the sorted list of queues depths vs just O(1) look ups of the random nodes’ depths.
It just would’ve been nice to cover that in the article.
There's some alternative hashing solutions (rendezvous hash, jump hash, maglev) that improve load distribution without the memory overhead, but still have similar challenges with node addition/removal overhead. Also, integrating state into any of these solutions is incredibly difficult and really antithetical to their point (being deterministic). This is fine when you're assigning uniform workloads and preferred when these assignments need to be long-lived, but for request load balancing you're dealing with uneven workloads and short-lived assignments.
This is where power of two load balancing shines, as it allows you to perform stateful load balancing (using queued requests, cpu utilizations, host perf) efficiently with low memory overhead and minimal effort to add/remove nodes. And if you lose a node with state its no big deal because the lifetime of a request is so short.
The proposed algorithm only needs to consider two, and from what I can gather should work in parallel with multiple load balancers without coordination.
The other answers are saying for algo reasons. Sure, but, pick two is even more better for uptime reasons.
In serving at scale, quite often the "least occupied" server of all of them is one with an error causing it to have least load/connections/whatever.
By starting with a random two, you pick among probably healthy servers, and then the least load.
> good improvements with nginx load balancing when changing its configuration to least_conn, that chooses the server with least connections.
nginx doesn't know why your app is failing to keep connections open long enough to serve them, least_conn between two random results in far fewer 100% "outages" when a single box or subset of boxes among 10s or 1000s is bad.
Sure with 10 or so hosts it’s trivial to select the least busy one but at the scale of 1000s of hosts you begin to run into a ton overhead just to decide where to forward a request.
As many others have said, O(N) for each request vs O(1) while avoiding hot spots.
Short version: it's robust to staleness and concurrency, which is important in distributed settings ("find the best" is hard to do when "best" is a moving target).
A server is taking requests 30 at a time in separate threads. To add a sequential step for picking which target is next, you need to add a synchronous piece of code that all of the threads depend on- a bottleneck. That's doable, but it adds complexity.
And then you have (let's say) 30 servers making calls. Do they each use round robin load balancing separately, or do they try to coordinate? That's difficult.
"Pick two randomly" requires no state. Every thread and server can run the same algorithm without regard for each other.
That’s not quite true. Picking between two “buckets” still requires knowing how many “balls” are in each which is state. That state can be local to each server or global, that state can be accessed concurrently or synchronously, but you still have the same problem to solve.
The proposed algorithm should work sufficiently well for multiple load balancers without coordination, assuming the number of load balancers is small compared to the number of servers processing the load.
In distributed systems with large numbers of producers (like one of the systems Mihir has worked, AWS Lambda) round robin essentially degrades into random placement as the number of producers increases.
We had an API that was a kitchen sink of different tools, including reporting. Eventually we split reporting off into its own service/LB group because you can get calls delayed that are returning 100 bytes of text being delayed by a few megabytes of data being assembled.