9 ms·
Deterministic Aperture: A distributed, load balancing algorithm
- web007 6y ago[2019] please - I can't believe I didn't see this before now! I don't see any mention of the "load" feedback mechanism to let clients discover a server's load, but I assume that comes as part of their framework.
- jeffbee 6y agoIsn't the squid-like graph in the middle of the article caused simply by using tiny subsets over tiny populations? Would you really bother with subsetting if each participant had fewer than, say, 100 connections? Put a different way, is the reduction in connections from 280k to 25k at the end of the article a thing worth achieving? Any single Linux process can deal with 280k connections, so having that be the aggregate number of connections in a distributed service architecture strikes me as a non-problem.
- scottlamb 6y agoInteresting to compare with Google's approach to the same problems: https://sre.google/sre-book/load-balancing-datacenter/ https://sre.google/sre-book/load-balancing-datacenter/ In my experience Google's deterministic subsetting usually ensures backends have even connection counts. There's one significant corner case though: a large backend service that has a large number of small clients, such that client_task_count * subset_size < backend_task_count. Then it basically degrades to random selection. This can happen to infrastructure that supports all of Google's products, many of which are relatively tiny. Seems like their "Continuous Ring Coordinates" section onward is designed to address this same problem. I got a bit lost following it though; I might need to read through it sometime my kids are not also wanting my attention...
- fastest963 6y agoHow do they decide how big to make the slice for each service? Is that still manual?
- nednar 6y agoIf you have three clients, each gets a third of the ring's space. If you have 200 clients each gets a 200th of the ring's space.
- scottlamb 6y agoUnfortunately this wouldn't be practical when the client's task count changes frequently. [1] You wouldn't want all the client tasks to recalculate their subset every time, so anything that requires each client task to know the total client task count is a dead end. Likewise when the backend task count changes, client tasks shouldn't change whether backend tasks that existed before and still exist are in their chosen subset. [1] Autopilot: Workload Autoscaling at Google Scale, https://research.google/pubs/pub49174/ https://research.google/pubs/pub49174/
- nvartolomei 6y agoI have implemented an interactive demo of this algorithm which can be found at https://nvartolomei.com/weighted-deterministic-aperture/ https://nvartolomei.com/weighted-deterministic-aperture/. Make sure to click "Verbose state" spoiler. Context: Earlier this year I explored using deterministic aperture for balancing writes/appends based on storage usage and implemented a visualization to help explain the algorithm to my colleagues.
- siscia 6y agoThis kind of reports are always interesting, but they ALWAYS miss why they are not following standard control-theory principles. At this point it is a long time since I finished studying control-theory so I may be missing something obvious.
- federico_c 6y agoI might be spilling some secret sauce here but... Control theory is mostly based on physical devices and analog variables, so it does not often fit with software systems. However you are right nonetheless. Large networks contain a lot of caching of different kinds and other forms of data replication. Caches warm up by transferring data and the available bandwidth to do so is never infinite. Especially when flipping traffic between whole datacenters. Most load balancing systems are simply unaware of this.
- smadge 6y agoFor an unfamiliar audience, can you describe which standard control theory principles these kinds of reports omit?
- tadkar 6y agoI am going to massively over simplify, but here goes. The really big idea from control theory is the idea of negative feedback. This basically boils down to measuring the output of your system and making the input to the system some function of the difference between the input signal and the output. The parent post refers to PID control. This refers to the three types of commonly used things to do with the difference (or error) signal. P is for proportional - where you just multiply the error signal with a constant. What this does is to encourage the system to track the level of the input signal (but potentially with some lag) I is for integral. This is where you integrate the difference signal over time. What this does is to reduce the lag between the input and output signal D is for derivative. This is where you feed back in the derivative of the error signal What this does is to damp down the swings in the system especially those that come from being too aggressive with the above two knobs. Good controller design often comes down to picking the right weights for each of the three types of feedback functions you can input into the system. So in this example, it might be that you distribute your requests to servers based on how over or under loaded the servers are...
- bullen 6y agoI use DNS round-robin for my setup, that way it's distributed geographically: http://host.rupy.se?dark http://host.rupy.se?dark
- twic 6y agoWhat if the backends announced their load via multicast? Then every client could know the load of every backend, and could continue to use P2C or something like that, without needing to maintain a socket for every backend. There's still a problem to solve, because you don't want to be opening and closing connections to backends all the time - you want to reuse a small number of connections. So maybe you keep a set of connections open, distribute requests between them using P2C or a weighted random choice, but also periodically update the set according to load statistics.
- jeffbee 6y agoThis is a bit like what Google does, but not multicast. Inactive RPC channels are downgraded from TCP to UDP and the health/load packets are sent less often. A client can maintain a smaller active subset of channels but still have complete load data for all its peers.
- scottlamb 6y agoGoogle doesn't do the UDP stuff anymore. IIRC, they never did it for channels not in the subset, and I don't think it'd make sense to given the goal of having a stable subset choice. And they haven't been able to actually send RPCs over a channel in UDP mode since LOAS IIRC, so the UDP stuff was useless and eventually removed. I'm not sure it was ever that useful anyway... They do still deactivate inactive channels, just totally rather than downgraded to UDP. Annoyingly, they only compute this inactivity for channels to individual tasks, not the greater load-balancing channels, and freshly started backend tasks always come up as active. So if you have a lot of inactive clients, when your tasks restart the inactive clients all rush to connect to it and your task sees noticeably higher health-checking load until the clients hit the inactivity timeout and disconnect again. To more directly answer twic's question: I can think of a few reasons multicasting all the load reports probably doesn't make sense: * That sounds like an awful lot of multicast groups to manage and/or a fair bit of multicast traffic. Let's say you have a cluster of 10,000 machines running 10,000 services averaging 100 tasks per service (and thus also averaging 100 tasks per machine). (Some services have far more than 100 tasks; some are tiny.) Each service might be a client of several other services and a backend to several other services. Do you have a multicast group for each service, and have each client tasks leave and join it when interested? I'm not a multicast expert but I don't think switches can handle that. More realistically you'd have your own application-level distributor to manage that, which is a new point of failure. Maybe you'd have all the backend tasks on a machine report their load to a machine-level aggregator (rather than a cluster-wide one), which broadcasts it to all the other machines in the cluster, and then fans out from there to all the interested clients on that machine. That might be workable (not a cluster-wide SPOF, each machine only handles 10,000 inbound messages per load reporting period, and each aggregator's total subscription count is at most some reasonable multiple of its 100 tasks) but adds moving pieces and state that I'd avoid without a good reason. edit: ...also, you'd need to do something different for inter-cluster traffic... * They mention using power of two choices to decrease the amount of state ("contentious data structures") a client has to deal with. I think mostly they mean not doing O(backend_tasks) operations on the RPC hot path or contending on a lock ever taken by operations that take O(backend_tasks), but even for less frequent load updates I'm not sure they want to be touching this state at all, or doing it in a lockless way, and ideally not maintaining it in RAM at all. * The biggest reason: session establishment is expensive in terms of latency (particularly for inter-cluster stuff where you don't want multiple round trips) and CPU (public-key crypto, and I don't think an equivalent of TLS session resumption to avoid this would help too often). That's why they talk about "minimal disruption" being valuable. So if you had perfect information for the load of all the servers, not just the ones in your current subset, what would you do with it anyway?