Hi everyone, thanks for joining our talk.
I'm Lu and this is my colleague Chen Ru. So today we're going to talk about routing LLM inference in production. Specifically how our system evolved from routing based on feedback loops driven by engine signals to a more explicit and predictable policy, which is still informed by engine signals. However, the way we use it is different. So for the agenda today, we're going to begin by introducing the inference load balancer, what it is, what it does, and how it has evolved.
And then Chen Ru will walk us through the newer control plane and data plane driven architecture, what are the responsibilities of each, and followed by a concrete case study of how we reduce the global network overhead. And in the end, I return to discuss the protection mechanisms that help keep the system stable under production level stress. So to begin with, what is the inference load balancer? And where does it sit? So this is a very high level diagram of the system we are talking about. On the left hand side are the front end clusters. Those are the CPU clusters that act as gateways into our system.
And they receive user requests, then prepare them into the inference request that can be processed by the inference engines. And on the right hand side are the engine clusters, which are usually GPU clusters, and each hosting multiple inference engines. So that's why they got the name of engine clusters. And as you may have already heard, nowadays the GPUs are popular and expensive. So sitting in the middle, it is the ILB or inference load balancer. It actually runs on the front end clusters, but it's also a bridge into our inference stack. It has two main responsibilities: select an engine and approximate the request.
For this talk, we are going to focus on the engine selection parts. So in some ways, ILB resembles a very traditional load balancer, because a request usually targets a model, and a model is backed by multiple engines. They may live on different clusters in different regions, or even across continents, because that gives us good resiliency towards localized degradation or cluster failures. However, the inference stack, or the uniqueness of inference, introduces a lot of nuances.
Like it has to consider a bunch of signals reported in real time, like the well-known time-to-first token, TTFT, time between output tokens, also known as token throughput, or time between tokens, and other healthiness and utilization signals. Besides, there's an important concept of a KV cache, which is also well-known. But for example, when the conversation already has a lot of useful context cached in one engine, sending the follow-up terms of the same conversation back to the same engine, we are avoiding recomputation, improve efficiency, and reduce latency.
So the combination of performance, reliability, locality, cache awareness is what makes it such an interesting problem. So how we attempted in the problem, let's take a look at the early days. And to be honest, early days in this industry sounds a lot more historic than it really is. And the routing process at that time began with a filter of each request may not be served by all the engines because of constraints such as capabilities, or geo-restrictions due to compute or data residency. And among the remaining engines, ILB used a weighted consistent hashing to select the best destination engine for a request for a certain user.
Then the important question becomes, where do the weights come from? So they were generated by a periodic feedback loop. The inference engines, as mentioned earlier, report all kinds of the signals we care about. And the controller will periodically smooth out those signals and compute performance score. The performance score then will be compared against the fleet average. Then the weight will be adjusted for each engine, if the weight goes up, if the performance is better, or it goes down when the performance is worse than the fleet average. And this generated weight will impact the routing, and then it's basically a control loop.
Conceptually, it's very similar to the PID controller. And no, this PID controller will not help you cure a Linux process, but instead it's a classic control theory technique that continuously steers the system towards its desired state. And we just borrowed this important concept, the proportional part of it, and applied it into our load balancer. So it has a lot of nice properties. For example, it could combine the useful signals we care about into the single routing decision.
And because it adapts to the observed performance, as what we mentioned earlier, there's a lot of constraints, and those constraints might have some engines busier because they can serve more requests, more kinds of requests than the remaining. But those busier signals will be fit into the next loop, and resulting in the less constrained request that can go to more of those kinds of engines. So basically, they're self-balanced out. And to some extent, this just means we don't need to do a lot of manual intervention, and it just works. However, that kind of adaptability comes with big trade-offs.
Because of the same reason that it combines so many signals, it's also very hard to reason about a particular routing decision, or why some engine gets a higher weight than we expect. And every time we want to fine-tune towards some aspect, it's almost impossible to not impact something else. And the load is not always very evenly distributed, because sometimes a model is served by engines on different GPU SKUs, and they have different characteristics. Then the problem becomes a lot trickier.
And the feedback loop sometimes creates bad oscillations, because when you shift an engine away some traffic, the engine turns a bit cooler, and this signal gets fit to the controller. The controller now thinks, hey, this engine can take a lot more traffic. Then some traffic is going to be shifted back and forth between a few engines, and disrupting the KV cache utilization. So all those limitations motivated us to rethink about the architecture, and see if we have new ways to address the problem. So I'm going to hand over to Chen Ru to deep dive into the new architecture we tried out. Yeah, thank you, Lu.
So I'm going to talk about the architecture of the load balancer, and how do we reduce the overall overhead with our routing algorithm. So the load balancer answers one question: for each request from a CPU cluster, which engine should serve it? A most naive baseline might be round robin, which sends requests across engines evenly. But if you think a little bit more, that doesn't make sense. Because engines are not homogeneous, they can have different hardware and capacity, different health, and also different distance from CPU cluster. Also, round robin could break cache locality. Related requests that could reuse the same engine cache might be sent to different engines.
A probably better solution might be for each CPU cluster, it choose the best engine from its own local view. But that's not enough either. Think about one extreme case. Multiple CPU clusters route traffic to the same engines independently, which could overload that engine while leaving other engines underutilized. So what we need is a globally optimized solution. A control plane that has a global view for all the CPU clusters and GPU engines, and could compute a globally optimized routing answer. And the data plane can make a routing decision quickly based on the answer pulled from the control plane. Now, let's look inside the control plane and data plane.
In the data plane, there is an engine selector, which selects engine for each request. It reads the local routing state, which includes the candidate engines and the routing weights for each candidate engine. Both of them are refreshed asynchronously in the background. So we don't need to ask the control plane before we make a routing decision for each request.
Also, the data plane collects real-time engine signal, such as number of ready replicas, engine health, etc., to service fast local guardrails. In the control plane, the data loader combines those live engine signals and network overhead. And with offline regressions of capacity, TTFT and TBOT, the optimizer could turn those data into routing weights. And the control plane will publish the routing weights for each data plane to pull.
In this way, no request needs to wait on the data plane. The control plane continuously computes the next globally optimized routing weights snapshot, while the data plane makes a routing decision based on the latest snapshot already installed locally. In summary, there are three important paths through the system. The first path is the inference request path. The request arrives to the CPU cluster and the data plane inside the CPU cluster will select engine for that request based on the local routing state and forward the request to the selected engines. The second path is the engine signal path.
The system continuously collects real-time engine signal, such as TTFT, TBOT, number of ready replicas, and engine health, etc. Both planes need those real-time engine signals. The control plane needs them to compute a globally optimized routing weights, while the data plane needs them to service fast local guardrails.
And the third path is the routing weights path. The control plane computes and publishes the routing weights, and the data plane pulls the updates to its local cache. So only the first path is synchronous, but it's fast and only local inside the data plane of the CPU cluster. The other two loops are asynchronous loops, and they are to improve future routing decisions. So that's pretty much of the architecture part, but that still leaves one question: how do we compute those routing weights? But before answering that question, let's answer another question first.
Why not just send a request to the nearest engine? That's because the traffic demand and GPU capacity are not geographically balanced. For example, in region one, CPU cluster A sends 90 RPS, and the nearby engine A can serve 100 RPS. So in this case, nearest engine only is fine. While in region two, CPU cluster B sends 120 RPS, and the nearby engine B could only serve 100 RPS. So in this case, if we insist on keeping everything local, the extra 20 RPS needs to wait on an overloaded engine B.
While in region three, we are only using 40 RPS of an 80 RPS engine C. That still leaves 40 RPS spare. So if we send the extra 20 RPS from cluster B to engine C, that will add network distance. So in this case, a further engine C is faster end to end. That's why we need something better than the nearest only routing. Now let's open the black box of the optimizer. The optimizer accepts four types of input. The request from each CPU cluster, the network latency to each engine, the available engine capacity and health, and also the TTFT, TBOT latency profiles that tell us how the engine side latency changes as the load increases.
And with those inputs, the optimizer turns the input to the output routing weights. The routing weights say for each CPU cluster, what fraction of its traffic should go to each GPU engine. And the optimization goal is straightforward: it's to minimize the expected end-to-end latency across all routed traffic. The important part is that the end-to-end latency includes both the network distance and the engine side latency. That means a nearby engine might be attractive when it still has room to serve traffic. While a further engine might be better if all the nearby engines are close to full. And the optimizer also needs to respect several hard constraints.
First, it needs to route all the traffic demand. Second, it needs to ensure all the engines stay within the effective capacity. Third, it needs to keep the routing weights non-negative. And that's pretty much my part. And Lu will continue to talk about the protection mechanisms in the system. Thanks, Chen Ru. So, as AI engineers, we all know that production in many cases is not behaving in the most ideal case. So, clusters can fail, GPUs or individual nodes can degrade, and networking can just get all kinds of mysterious issues. So, how do we keep our production system healthy as much as possible under heavy load? The first thing we have is penalties.
Basically, when an engine is an outlier, we detect the abnormality and try to reduce the routing weight to that engine. In that way, we give it a chance to either recover by themselves if there's some transient issue, or we can have a human intervene to rotate it out or replace the faulty hardware. And secondly, the retries, which is a very common technique used to mitigate problems. However, in some cases, it actually could make them even worse. Like when the system is very close to a tipping point or very heavily utilized. Retries send more load. And this more load causes more failures and causes more retries, which is the infamous retry storm.
So, we implemented caps or budgets to constrain retries into an acceptable region. And this actually needs to be dynamic because in the happy time or in the normal time, we can tolerate a lot more retries than when the system is heavily utilized. And finally, we have load shedding, which is our last resort when the production capacity couldn't meet the increasing amount of inference demands. So, we instead will try to have the system fail gracefully. We basically proactively load shed a portion of the traffic to have the system degrade gracefully. So, that pretty much concludes our talk today.
And thanks for joining us. Both of us will be around in our booth area this afternoon. So, if you have further questions, feel free to walk to the area and chat with us. Thank you. that can go to more of those kind of engines. So basically, they're self-balanced out. And to some extent, this just means we don't need to do a lot of manual intervention, and it just works. However, that kind of adaptability comes with big trade-offs. Because of the same reason that it combines so many signals, it's also very hard to reason about a particular routing decision, or like why search engine gets a higher weight than we expect.
And every time we want to fine-tune towards some aspect, it's almost impossible to not impact something else. And the load is not always very evenly distributed, because sometimes a model is served by engines on different GPU skills, and they have different characteristics. Then the problem becomes a lot more trickier. And the feedback loop sometimes creates bad oscillations, because when you shift an engine away some traffic, the engine turns a bit cooler, and this signal gets fit to the controller. The controller now thinks, hey, this engine can take a lot more traffic. Then some traffic is going to be shifted back and forth between a few engines,
and disrupting the Kiwi cache utilization. So all those limitations motivated us to rethink about the architecture, and see if we have new ways to address the problem. So I'm going to hand over to Qian Ru to deep dive into the new architecture we tried out.
Yeah, thank you, Lu. So I'm going to talk about the architecture of the load balancer, and how do we reduce the overall overhead with our routing algorithm. So the load balancer answers one question. For each request from a CPU cluster, which engine should serve it? Well, most naive baseline might be round robin, which send requests across engine evenly. But if you think a little bit more, that doesn't make sense. Because engines are not homogeneous, they can have different hardware and capacity, different health, and also different distance from CPU cluster. Also, round robin could break cache locality.
Related requests that could reuse the same engine cache might be sent to different engines. A probably better solution might be for each CPU cluster, it choose the best engine from its own local view. But that's not enough either. Think about one extreme case. Multiple CPU cluster route traffic to the same engines independently, which could overload that engine while leaving other engines underutilized. So what we need is a globally optimized solution. A control plan that has a global view for all the CPU cluster and GPU engines, and could compute a globally optimized routing answer. And the data plane can make a routing decision quickly based on the answer
pulled from the control plane.
Now, let's look inside the control plane and data plane. In the data plane, there is an engine selector, which selects engine for each request. It reads the local routing state, which includes the candidate engines and the routing weights for each candidate engine. Both of them are refreshed asynchronously in the background. So we don't need to ask the control plane before we make a routing decision for each request. Also, the data plane collects real-time engine signal, such as number of ready replica, engine house, etc., to service fast local guardrails. In the control plane, the data loader combines those live engine signals
and never overhead. And with offline regressions of capacity, TDFT and TBOT, the optimizer could turn those data into routing weights. And the control plane will publish the routing way for each data plane to pull. the control plane is a control. In this way, no request need to wait on the data plane. The control plane continuously computes the next globally optimized routing way snapshot, while the data plane makes a routing decision based on the latest snapshot already installed locally.
In summary, there are three important paths through the system. The first path is the inference request path. The request arrives to the CPU cluster and the data plane inside the CPU cluster will select engine for that request based on the local routing state and forward the request to the selected engines. The second path is the engine signal path. The system continuously collects real-time engine signal, such as TTFT, TBOT, number of ready replica, and engine house, etc. Both planes need those real-time engine signals. The control plane need them to compute a globally optimized routing way, while the data plane need them to service fast local guardrails.
And the third path is the routing way path. The control plane computes and publishes the routing way, and the data plane for the updates to its local cache. So only the first path is synchronous, but it's fast and only local inside the data plane of the CPU cluster. The other two loops are asynchronous loop, and they are to improve future routing decision. So that's pretty much of the architecture part, but that still leaves one question. How do we compute those routing weights? But before answering that question, let's answer another question first.
Why not just send a request to the nearest engine? That's because the traffic demand and GPU capacity are not geographically balanced. For example, in region one, CPU cluster A sends 90 RPS, and the nearby engine A can serve 100 RPS. So in this case, nearest release only is fine. While in region two, CPU cluster B sends 120 RPS, and the nearby engine B could only serve 100 RPS. So in this case, if we insist on keeping everything local, the extra 20 RPS needs to wait on an overloaded engine B.
While in region three, we are only using 40 RPS of an 80 RPS engine C. That still leaves 40 RPS spare. So if we send the extra 20 RPS from cluster B to engine C, that will add network distance.
So in this case, a further engine C is a faster end to end. That's why we need something better than the nearest only routing.
Now let's open the black box of the optimizer. Let's open the optimizer. The optimizer accepts four types of input. The request from each CPU cluster, the network latency to each engine, the available engine capacity and health, and also the TTFT, TBOT latency profiles. That tell us how the engine side latency change as the low increases. And with those inputs, the optimizer, turn the input to the output routing weights. The routing weights say for each CPU cluster, what fraction of its traffic should go to each GPU engine. And the optimization goal is straightforward. It's to minimize the expected end-to-end latency across all routed traffic.
The important part is that the end-to-end latency includes both the network distance and the engine side latency. That means a nearby engine might be attractive when it still has room to serve traffic. While a further engine might be better if all the nearby engines are close to full. And the optimizer also need to respect several hard constraints. First, it needs to route all the traffic demand. Second, it needs to ensure all the engines stay within the effective capacity. Third, it needs to keep the routing weights non-negative. Third, it needs to ensure that the control plane is the routing weights and the optimizer.
And the data plane pulls them and uses them to make a globally optimized routing decision. And that's pretty much my part. And Lu will continue to talk about the protection mechanisms in the system. Lu Wang Huang Thanks, Chenru. Lu Wang Huang So, as AI engineers, we all kind of know that production in many cases are not behaving in the most ideal case. So, clusters can fill, GPUs or individual nodes can degrade, and networking can just get to all kind of mysterious issues. So, how do we keep our production system healthy as much as possible under the heavy load? Lu Wang Huang The first thing we have is penalties.
Basically, when an engine is an outlier, we detect the abnormally and try to reduce the routing weight to that engine. In that way, we give it a chance to either recover by themselves if there's some transient issue, or we can have a human intervene to rotate it out or replace the faulty hardware. Lu Wang Huang And secondly, the retries, which is a very common technique used to mitigate problems. However, during some cases, it actually could make them even worse. Like when the system is very close to like a tip over or very heavily utilized. Retries, we are sending more load.
And this more load, we are causing more failures and causing more retries, which is the infamous retries storm. So, we implemented caps or budgets to constantly retries into an acceptable region. And this is actually even need to be dynamic because in the happy time or in the normal time, we can tolerate a lot more retries than when the system are heavily utilized. And finally, we have the load shedding, which is our last result when the production capacity couldn't meet the increasing amount of inference demands. So, we instead will try to have all the system fill. We basically proactively load shed a portion of the traffic to have the system degrade gracefully.
So, that pretty much concludes our talk today. And thanks for joining us. Both of us will be around in our booth area this afternoon. So, if you have further questions, feel free to walk to the area and chat with us. Thank you.
Thank you.