Random Load Balancing Is Unevenly Distributed
evanjones.ca
evanjones.ca
In general, LOC is more forgiving for uneven workloads, especially for Ruby services. If a slow request hits, that machine gets less traffic than the others and gets some room to recover and burn off its queue.
Where we found weird problems was at the extremes: very low concurrency and exceptionally high concurrency.
Low concurrency meant that one host consistently had 2 connections and another consistently had 3. That's 50% more load on one host. I believe AWS has fixed this and now tie-breakers are resolved randomly, but it was kinda painful to deal with. Our solution was to scale vertically and run fewer hosts. 19 on one host vs 20 on another was less jarring.
The other extreme was 1k+ websocket connections per host. When we autoscaled, it would add a new host. That host would have 0 connections and so the next 1k connections would all go to the new host. For that system, we changed it back to round robin.
I assume you're saying that this happened over a period much longer than the lifetime of a connection? If all of the connections are super long lived, I'm not really sure what else the load balancer could do, other than forcibly close one and make the client swap to the other server, but naively that doesn't really seem like something I'd want a load balancer to do.
Something like that.
This seems like a non-problem. What would you wish an idealized load balancer to do in this case and how would that improve any relevant metric?
This worked well on a fleet of a couple thousand servers handling thousands of requests per second each. I can't recall why, but lower volume and few servers caused uneven load, and we went back to least connections for those clusters.
Some load balancers support readiness checks. When a backend is over-worked it can fail readiness to excuse itself from more work. We implement one variation of this by looking at load metrics (cpu, disk, mem, net) [0]. readiness is binary (ok / !ok) and not weighted (ex: no. of conns). The backend, however, is free to load shed (reject already sent requests) as it sees fit.
[0] https://www.brendangregg.com/blog/2017-08-08/linux-load-aver...
https://github.com/torvalds/linux/blob/b97d64c722598ffed42ec...
The parent's concept of reporting "hotness" is really powerful. It lets a central authority (the load balancer in this scenario) make an informed decision about where to send traffic: the coolest backend. In the case where all systems are actually quite hot, it probably makes sense in most scenarios to continue sending traffic as best-effort.
In the case of readiness checks, nacking a check removes the backend from the pool of eligible new-traffic recipients. At least this is true for systems I've encountered. The solution you propose is problematic as it's possible for all systems to decline traffic resulting in no available backends -- rarely a desirable state -- which the load balancer usually can't do much about.
A good system will allow you to decouple the signals (loadavg, connections, latency, CPU, etc.) and the health decision (readiness checks) and the interpretation of that decision (load balancer policy) to provide the best of both these worlds.
Now, there are many scenarios where simply excusing yourself works fine -- processing expensive batch jobs, for example -- so this can work nicely, but for typical production traffic scenarios I'd advise against this approach.
I don't see why this is a problem? This is when we should start rejecting requests on the frontend
Let's say each node can handle 5k rps and you have 10 nodes. You can handle 50k rps. If you are receiving 40k rps, a good strategy will put each node at 80% capacity. A bad strategy will knock a node out, reducing your total capacity, putting extra pressure on the rest of the system, causing more failures, and more pressure. This is called a thundering herd.
At some point, your only option is load shedding. But with a bad LB strategy, you start load shedding much much earlier than you should. This is a bad experience for customers that is avoidable with good LB strategies.
We run load balancers in a fail open mode. As in, if every backend is excusing itself, then none are excused.
But as you point out, load balancing is a hairy beast unto itself.
We don't use loadavg. The linked article in my parent comment makes the same point as you are (:
The more traffic you’re handling, the more the central limit theorem applies as you are summing the behavior of lots of random events drawn from various random distributions, and the more the system behavior regresses to the mean.
I think of it as analogous to how when an engine is idling it tends to pop and judder, but as you rev it up the smoother it sounds until at sustained high revs it can make almost a pure note. All the variations are smoothed out and the system runs much more evenly and predictably, but when it drops to idle speed the subtle variation in the fuel air mix and the pressure in the exhaust and inlet every stroke makes it chug and sputter chaotically.
This included regular stratified sampling[2] as well as using quasi-random low-discrepancy sequences[3] like the Sobol sequence[4], which would "space out" the random samples.
Probably not terribly relevant for load balancing, but I did find the quasi-random stuff fascinating. How to make something "random enough" without being "too random".
edit: Though now I got a shower thought moment, could stratified sampling be used as a variant of the "two choices", even if each load-balancer worked independently? That is, each load-balancer stratifies the resources, and keeps a separate index to the current bin. Then for each request pick a random/suitable resource from the current bin and increment the bin index, wrapping when needed. Hmm...
[1]: https://en.wikipedia.org/wiki/Path_tracing
[2]: https://en.wikipedia.org/wiki/Stratified_sampling
^^ This is worth highlighting.
"Randomness is clumpy"
-- Attributed to various authors
The load balancer could keep track of how long requests for a specific path tend to take, and what servers are currently handling and use that information to be smarter about routing requests. I wouldn’t be surprised if that already exists in things like haproxy or nginx.
There’s issues where people stuff data into the URL, like ids or keywords for SEO. But that’s possible to workaround.
The article points out random load balancing scales sub-linearly with increase in capacity?
If we distribute a set of items randomly across a set of servers (e.g. by hashing, or by randomly selecting a server), the average number of items on each server is num_items / num_servers. It is easy to assume this means each server has close to the same number of items. However, since we are selecting servers at random, they will have different numbers of items, and the imbalance can be important.
...
This is a classic balls in bins problem... the summary is that the imbalance depends on the expected number of items per server (that is, num_items / num_servers). This means workload is more balanced with fewer servers, or with more items. This means that dividing a set of items over more servers makes the distribution more unfair, which is a reason we can get worse than linear scaling of a distributed system.*There are other techniques that are optimal like Parallel Depth First scheduling from Blelloch.
Since then there has been a lot of work for distributed work-stealing.
Boggles my mind that at "cloud-scale" people still use round robin and assume uniform workloads and that at the server level there is uniform load as well.
Seems many load balancer providers missed that memo.
And there is no standard way I’m aware of to usefully report back or communicate this data to the LB.
So what else are you going to do?
Given tasks do not run in identical time, there's going to be variance come what may.
I would suggest a pool of 3 should tend to be more stable. But, your platform people may be recommending that 3 and below need to be sized to work as 1, only 4 and above can start to assume a minimum cohort (2?) exist.
Sharding is not LB. It's just a hash approximation to grouping things. Maybe you're telling a story about using a random() to emulate how far out of 50/50 a 2 part sharding story can get?
Routers doing LB may well do simple prefix triage. You may be stuck on one side of the LB because of who you are, and nothings going to change it.
(I don't do this for a living, wiser people may tell you I'm wrong)
For a truly random source, the odds of getting the same number in the next roll are identical to getting it on any prior roll. For example, flipping a coin twice is as likely to get you heads twice as it is to get you heads/tails.
eventually things even out. But if it’s truly random, when that is going to be the case should also be completely unpredictable. It might end up being after the service is deprecated, for instance. You might end up never getting a specific ‘desired’ answer ever. For instance, if the test is flipping a coin and we do it 10 times, there is a small chance we never get tails and get 10 heads in a row from a truly random source.
That is working as expected.
Usually (at least with load balancing) they instead want a smoother distribution, where it’s less likely to get the same answers they just got, when they just got them.
So for example, they’d prefer that if they just got a heads they’re more likely to get a tails next (but it won’t always happen).
Notably, this isn’t really random. It’s also a really hard thing to do in practice without also getting all nodes in lockstep, being predictable to an outside attacker, or having internode communication. Maybe impossible, I don’t remember.
https://techcommunity.microsoft.com/t5/core-infrastructure-a...
AWS ALB (and others I'm sure) can balance by "Least outstanding requests" [1] which means a server with the least in-flight HTTP requests (not keep-alive network connections!) to an app server will be chosen.
If the balancer operates on the network level (eg NLB) and it maintains keep-alive connections to servers, the balancing won't be as even from a HTTP request perspective because a keep-alive connection may or may not be processing a request right now and so the request will be routed based on number of TCP connections to app servers, not current HTTP request activity to them.
[1] https://docs.aws.amazon.com/elasticloadbalancing/latest/appl...
[Also, to be pedantic, it should really be called "fewest requests" - but I've never seen a product call it that!]
I think what the article misses highlighting is that this is about stateless load balancing which improves throughput at the LB, at the cost of potentially suboptimal distribution.
The other nit I have is calling hashing random, it's not, and a lot goes into selecting hash algorithms for LB that don't result in mass movement of requests every time a server is added or removed from the pool.
If the author wants stateless balanced load and assumes all requests are equal, why not use an approach like round robin?
Beyond that, you can get into more stateful methods that look at CPU use on servers, or any other metric you like to load balance on the resource you care about, ie memory or CPU vs number of open sockets.
Developers have a nasty habit of assuming a glitch or someone else changed something and if a reload fixes it they don’t note the pattern.
In this scenario, unless each request is exactly the same you will end up with uneven load as some sessions persist longer than others.
This can still be stateless with use of a cookie for protocols that support it but is not the idealized "round robin and every request takes same time and resources".
The other alternative is a stateless app that is architected for round robin lb.
It’s a good trade off when throughput is low and the impact of lumpy traffic is high. At higher volumes, random allows stateless load balancing.
Use the tool that makes sense for your project. Even at the same system I’ll have some services that operate at low volumes with high latency and others with high volume and low latency.
Just like JSON (and js) hashmaps should really be treemaps and retain order.
Last but not least you should be able to make a TCP stream loose order. So you get UDP speed without changing pipes/ports and all that jazz.
It's all about doing that job when you design a protocol.
That only works if the clients choose the first IP address in the response to connect to. But IPv6, IIUC, specifies that clients instead always use the IPv6 address it has the longest stretch of bits in common with, compared to its own IPv6 address.