Why is load balancing gRPC tricky?
majidfn.com
majidfn.com
This is good practice anyway. Having long-lived servers causes all kinds of issues. In the old days, if you had a server that was up for years, you had no idea whether you could reboot the server safely or not. Additionally, from a security perspective, the longer-lived a server is, the more likely that an attacker installed some manner of backdoor on the server that's just waiting for the attacker's opportune moment to come back and exploit it.
Much better to make chaos an expected part of your operations - make sure your servers are cattle (not pets), give them maximum lifetimes, and terminate them on a regular basis. The more frequent the termination, the better prepared you are for single-node overload, and the less likely that it happens in your use case in the first place.
If you want to ignore most of the picture and focus on the part you like, you're like a horse with eye covers. Only what you look at exists, and where you look determines where you want to go. It's not pretty.
I feel like that's not quite true. The server could notice it's over-worked and start telling clients to back off or otherwise to try another system. At this point the clients can make the call to either back off or terminate the connection. The server can increase its "back off for N ms" until it finally has to terminate connections.
It requires a bit of extra coordination, but not much.
Also, even without such a system, forcing a reconnect when a new instance starts up/ periodically still reduces connection load a ton. If you force reconnects every minute, and you have 10/rps, that's 1 reconnect per 600 calls.
I guess in theory a load balancer is the most information to make these decisions, so really you'd want it to be aware of 'back off' requests from the server. One way of doing this might be something with health checks? Dunno.
Let me describe a scenario my team ran into while using option #3 i.e., Look-Aside. We used etcd as service discovery store. Servers as part of their startup would register with etcd. Clients would query etcd and randomly pick a server to talk to. Now, for some reason the entire service fleet got isolated. A couple of hours later the network partition got back just fine but all the servers had shut down.
The real pain started when we began bringing up servers. As soon as a server was brought up it would register with etcd, immediately get bombarded by hundreds of clients and then would promptly die. No matter what permutation we tried servers would just refuse to come up because of thundering herd. The only option was to shut down the client app, which was another fiasco because we hadn't built a clean way to do it, and bring up all the servers and then gingerly bring up clients hoping they won't kill the servers.
The downtime lasted almost for a day. And this is a largish app based ride hailing business I'm talking about so you could imagine the shitstorm it created among customers and investors alike.
A key learning for us was to isolate service startup from load-balancer registration. A service should not be responsible for registering itself with a load balancer that should be for someone else to do.
While gRPC does have lots of positives it still has some way to go before it reaches the operational maturity. People will discover these deficiencies in a painful manner.
I don't know what the state of the art around tooling to manage this is. For some reason, I suspect that the service meshes punt on this, because N * M isn't a problem in the demo environments where these systems spend most of their time. Meanwhile, the big companies that hit the scalability limit of N * M connections across upstream/downstream pairs wrote their own service discovery stuff decades ago.
I like gRPC overall, but it's definitely got some rough spots like this. I'm fortunate enough to be using it at a small enough scale that they don't really hinder me.
My sense, based on the multitude of GitHub issues like this, is that the project itself is thrashing. They've got a relatively small core team that seems to be micro-managing all the official implementations, largely out of a desire to maintain cross-language compatibility. But the cross-language compatibility is actually rather poorer than you'd expect, because they're trying to maintain it by directly micro-managing the implementations, which turns every feature design project into a complicated n-body problem.
It would be nice if they could just publish a comprehensive formal spec, and let the language implementations follow it. Then each (official) language would have only one external thing to track, instead of ten.
This is what I meant by "operational maturity". A big factor for a new tech to gain adoption is the ease with which it fits in the existing ecosystem. For gRPC it's about load-balancer, reverse proxies, deployments, and so fourth. I'm sure gRPC will get there, depending on how fast it's adopted in big tech companies but I suspect it's not there yet.
http://cbonte.github.io/haproxy-dconv/2.4/configuration.html...
Something that I implemented at my last job and others rediscovered as part of my current job is to implement an "administratively up/down" API as part of the control plane and only have the server announced if it was "up." Decoupling the announcement from process start/initialization complete allowed us to roll out new versions of software in a disabled fashion and then "flip the switch" (red/black deployments). It also enabled us to take individual instances out of service without killing them, enabling developers to debug issues/anomalies more easily.
Load shedding/backpressure/rate limiting at various layers is also extremely helpful, whether at the load balancer/API gateway or at individual servers. That has saved our bacon numerous times.
You can also load-balance on the client side. Clients connect to a management server, and the server tells them what percentage of connections to send to which upstream. I don't think anyone has written any publicly available software to manage this, but it can use the same underlying protocol as Envoy's configuration servers (which are also rare in public; I wrote one, but it sends clients Endpoints and Clusters and the gRPC implementation happens to want Listeners). The xDS implementation is a newer version of grpclb, which maybe had more support but I've never heard of anyone using it. (Google internally uses the client side load balancing. Worked great when I used it, so was always surprised it wasn't more popular in the real world.)
The article is really negative about client-side load balancing, but you can't complain about the latency properties -- no proxy between you and the other end. I always wished the Internet standards did a better job here; you can serve multiple A records for a DNS lookup, and browsers kind of sort of try to load balance (by selecting one of the endpoints randomly), but it doesn't work too well if one of the backends is broken. If you could just supply a pre-defined health-check function and maybe select a simple balancing algorithm, HTTP would be just a little bit smoother for everyone. (If you put topology information in the DNS lookup, you could even balance geographically without having to run special DNS servers that serve different replies to different users. A lot of upstream DNS caching hacks could then go away. It would be great! People are onto the value of this; Envoy Mobile lets you serve this information to your mobile app, so it can easily survive regional failures in your server-side infrastructure. Probably not a lot of adoption, but something that browsers should start doing.)
I believe this is exactly the way that Hashicorp's Consul Connect service mesh is implemented for L7 routing.
That seems strange to me. Without end-to-end http/2, ALBs would load balance individual requests, not the entire TCP session, even with keepalives. And it sounds like individual messages on the same connection can be routed to different target groups, which would make it even more stange not to balance individual requests.
>Either of these approaches defeats a fundamental gRPC benefit: reusable connections.
Let's say the client reconnects every 50 requests. You have reduced the connection overhead by 98%. Are those last 2% really going to ruin your day?
Some hybrid of count and time might do, e.g. `50 requests AND 1+ minute since establishment`. But it's nuanced - surprisingly hard to find logic that works well in all cases.
Also, does anybody know how gRPC is used at Google? Why do they support its development (other than GCP clients)? I'm under the impression Stubby is a completely different codebase and is far superior and fully integrated in their stack.
As for Stubby it is a completely different codebase and is tightly integrated with a lot of internal tooling and features.
One other thing I miss from Stubby C++ compared to gRPC is that async stubby is super opinionated about it's threading model and came with batteries included, where I find the completion queue approach very hard to use.
Lastly, as far as load balancing goes. I think some load balancers with native gRPC support allows you to let different servers get different individual requests instead of requiring all requests go to the same server. Of course this only really works for unary RPC calls. But load balancing HTTP/2 isn't unique to gRPC nor keep alive http/1.1 connections. (Outside of the streaming RPCs of course, but there is an equivalent of websockets)
And yes, the completion queue approach is pretty cumbersome, would be very interesting to see how stubby went about it :P
Likely to maintain less code long term. It's nice to have a primative like this open source so you can use it everywhere (even on mobile). I'm not sure about the claim that gRPC isn't customizable enough, I think there are a lot of hooks in place. For example Open Census does a lot of the instrumentation that you get for free at Google with Stubby, but for gRPC: https://opencensus.io/zpages/
> How stubby went about it: Stubby is a lot closer to the experimental callback API talked about here: https://github.com/grpc/proposal/pull/180
E.g. you have 200 clients connected to 100 servers, your service discovery sends each client a list of 10 servers, and 2000 connections are established instead of 20000.
"Partial mesh" is a good description, but it's mostly used to describe routing, not load balancing.
Meaning that a rough lower bound on the memory usage under this scenario would be something like ~8GiB per million open connections, and ~8TiB per billion open connections.
You're right that it comes down to the extent to which that's a significant fraction of total system cost. I think I can say with reasonable confidence, though, that this would be a significant fraction of total system cost for the vast majority of systems.
A much more likely scenario is each of your 10100 things has at least 2-4GB of memory per CPU core, meaning at least 20-40TiB of memory in your system, meaning that the million connections cost you 0.02% of your system memory, or essentially nothing.
That and we've violated the laws of physics. Because even by the numbers that you agreed are grossly optimistic (which was deliberate), we're looking at that consuming an order of magnitude or two more RAM than we have across all our servers. So that would admittedly be cool. But still.
This is nearly two orders of magnitude off (250-500MB using 4-16 CPU cores) for some of our most connection-intense services. And as you said 8KiB is extremely low, in practice more like 2-10x that once all buffering is taken into account. >5% of my memory spent on connection handling that (by the original assumption) I don't need, yes, that worries me.
Because the service has no substantial local state so doesn't need any more than that? Should we invent unnecessary local state just so we can use more memory? That sounds considerably more surreal to me.
Or maybe instead we should use all the memory we saved to keep our 200-300GB working set for a different service in RAM colocated, which is what we do...
I found http://liblb.com/learn.html useful for understanding how different load balancing algorithms work.
> The official implementation uses a Per Call Load Balancing and not a per connection one. So each call is being load balanced separately. This is the ideal and the desired case and it will avoid having heavy sticky connections.
I wonder how much of this is just leaking from how stubby & GSLB is done internally at Google. I'll admit, I don't know much about gRPC.
Separate point: if you control both the clients and the servers completely, you can have servers respond with load shedding signals if they start to get overloaded, and use that as a hint for the client to go look elsewhere. This doesn't solve all problems, but it does make some cases easier to handle. Also it is obviously not a viable solution if you don't fully control all clients - some of them will just not respect your load shedding and continue to beat an individual server into the ground.
Any sane HTTP/2 client making N queries to the same host/IP ideally should use 1 connection since h2 supports multiplexing (which is bad for load balancing if you have an L4 LB). So even if you don't use gRPC, consider putting a per-call (L7) LB in front of your app if you use HTTP/2 protocol between your apps.
https://linkerd.io/2018/11/14/grpc-load-balancing-on-kuberne...
I was making this argument for years. WebSockets was designed for bidirectional data exchange between client and server. On the other hand, HTTP2 was designed for pushing static resources to the client to speed up loading of front ends (it was meant to make bundling front-end scripts redundant). HTTP2 was not designed for bidirectional data exchange over a stateful channel.
If it was properly stateful, then sticky sessions would not be required. With WebSockets, sticky sessions are not required when load balancing.
It already worked fine way before October with envoy or haproxy for example.
I don't see anything difficult here...