1. Problem → Why → Optimization#
PROBLEM Standard load balancing (round-robin, least-connections) performs
badly for LLM serving. Some replicas are overloaded while others idle,
and prefix cache hit rates collapse.
WHY Requests differ in cost by 1000x, are long-lived, and have per-replica
cache affinity. None of those hold for ordinary web traffic.
OPTIMIZE Load-aware and prefix-aware routing.2. Why the standard algorithms fail#
ROUND ROBIN
Assumes requests have similar cost. They don't.
A replica that receives three 100k-token requests in a row while another
receives three 50-token requests is 100x more loaded, and RR can't see it.
LEAST CONNECTIONS
Better — accounts for in-flight count. But a replica with 2 huge requests
is more loaded than one with 8 tiny ones, and this can't tell.
RANDOM
Same problems as round robin, plus variance.
LEAST RESPONSE TIME
Reacts to the symptom after it appears. For requests lasting minutes,
the signal is far too laggy.And all four destroy prefix caching, because consecutive turns of the same conversation land on different replicas with cold caches.
3. Simple analogy#
Assigning patients to doctors.
Round robin: next patient to the next doctor in rotation. Fine if every appointment is 10 minutes. Terrible when appointments range from 5 minutes to 4 hours.
Least connections: next patient to the doctor with fewest patients. Better, but a doctor with one 4-hour surgery is busier than one with three 10-minute consultations.
What you actually want: assign by estimated remaining work, and send returning patients back to the doctor who has their file (prefix affinity).
4. The routing signals#
SIGNAL QUALITY HOW TO GET IT
in-flight requests poor the LB counts
queue depth ok replica reports it
KV cache utilization % GOOD replica reports it ← the best load signal
running batch size good replica reports it
estimated queued tokens GOOD replica reports it
prefix cache hit probability GOOD hash the prompt prefix
p95 TTFT (recent) lagging LB measuresKV cache utilization is the single best load signal for LLM serving, because it directly measures the binding constraint. A replica at 95% KV utilization cannot accept more work regardless of how many requests it has in flight.
Modern engines expose these on /metrics:
vllm:gpu_cache_usage_perc
vllm:num_requests_running
vllm:num_requests_waiting
vllm:prefix_cache_hit_ratePoll them at 1-5 second intervals, or have replicas push them.
5. Routing algorithms#
Least-loaded (power of two choices)#
// Power of two choices: sample 2 replicas, not all N, and take the less loaded.
func route(replicas []*Replica) *Replica {
a, b := replicas[rand.Intn(len(replicas))], replicas[rand.Intn(len(replicas))]
if a.KVUtilization() < b.KVUtilization() {
return a
}
return b
}Power-of-two-choices is the right default. Sampling two and picking the better one gets almost all the benefit of checking all N, with none of the herding behavior that “always pick the least loaded” causes (where every LB instance sends everything to the same replica at once).
Prefix-aware routing#
type Router struct {
mu sync.Mutex
replicas []*Replica
prefixTable map[uint64][]*Replica // prefix hash -> replicas that have it cached
}
func (rt *Router) Route(tokenIDs []int) *Replica {
// hash the first N tokens (the likely shared prefix)
h := fnv.New64a()
for _, id := range tokenIDs[:min(512, len(tokenIDs))] {
binary.Write(h, binary.LittleEndian, int32(id))
}
key := h.Sum64()
rt.mu.Lock()
defer rt.mu.Unlock()
// among replicas that have this prefix cached, pick the least loaded
var best *Replica
for _, r := range rt.prefixTable[key] {
if r.Healthy() && r.KVUtilization() < 0.9 && (best == nil || r.KVUtilization() < best.KVUtilization()) {
best = r
}
}
if best != nil {
return best
}
// otherwise fall back to least-loaded-of-two, and remember the assignment
chosen := route(rt.replicas)
rt.prefixTable[key] = append(rt.prefixTable[key], chosen)
return chosen
}The tension: cache affinity vs load balance. Always honoring affinity overloads the replica
holding a popular prefix. The kv_utilization < 0.9 guard is the release valve.
A better formulation weights both:
score(replica) = w1 × predicted_cache_hit_rate(replica, request)
- w2 × replica.kv_utilization
- w3 × replica.queue_depthTune the weights by measuring end-to-end TTFT.
Session affinity#
For multi-turn conversations, route by conversation id:
Hash(conversation_id) → replica consistent hashing, so adding/removing
replicas moves few sessionsSimpler than prefix hashing and captures most of the benefit for chat, because a conversation’s turns share a prefix by construction.
Combine with a load release valve: if the affine replica is saturated, fall back and accept a cold cache.
Cost-aware routing#
Estimate the request's cost, and route expensive requests to a dedicated pool:
if request.prompt_tokens > 32000:
route to the long-context pool
elif request.max_tokens > 4000:
route to the long-generation pool
else:
route to the general poolSegregating long requests is one of the highest-value routing decisions. A single 128k-token request in a pool tuned for 4k requests consumes 32 sequences’ worth of memory and blocks prefill for everyone.
6. Under the hood — the LB’s own constraints#
Streaming responses hold connections for the whole generation:
- the LB must not buffer (Section II.08)
- connection limits per backend must be high
- idle timeouts must exceed max generation time
- draining a replica means waiting for in-flight generations
Health checking:
- use the readiness endpoint, not a real generation
- a replica at 100% KV utilization is BUSY, not unhealthy —
don't remove it from the pool, just stop preferring itThat last point matters: an LB that removes saturated replicas from the pool concentrates load on the remaining ones and cascades.
7. Performance#
Simulated, 8 replicas, heavy-tailed lengths, 60% prefix sharing:
Algorithm p50 TTFT p99 TTFT Prefix hit rate Throughput
Round robin 680 ms 8,400 ms 12% 1.00x
Least connections 520 ms 4,100 ms 12% 1.09x
Least KV utilization 410 ms 2,200 ms 13% 1.18x
+ session affinity 240 ms 2,400 ms 71% 1.42x
+ long-request segregation 230 ms 1,100 ms 71% 1.51xThe prefix-affinity row is the big jump, and it costs nothing but routing logic. The p99 improvement from segregating long requests is the second-largest effect.
8. Production implications#
- Do not use a stock HTTP load balancer’s default algorithm for LLM traffic. It will be round-robin or least-connections, and both are wrong.
- Expose load metrics from replicas and route on KV utilization.
- Use power-of-two-choices, not global-least-loaded (avoids herding).
- Implement session affinity for chat — it’s simple and captures most of the prefix benefit.
- Segregate long-context and long-generation requests into their own pools.
- Don’t remove saturated replicas from the pool. Deprioritize them.
- Disable buffering, raise connection limits, extend idle timeouts at the LB.
- The Kubernetes Gateway API Inference Extension standardizes model-aware routing. Its
InferencePoolAPI is stable (v1), and the endpoint picker that scores replicas on queue depth, KV-cache usage and prefix affinity is now developed in the llm-d project. If you are on Kubernetes, start from it rather than writing a balancer.
9. Common mistakes#
Round-robin with prefix caching enabled. You built the cache and then guaranteed it misses.
Global least-loaded across many LB instances. All of them pick the same replica simultaneously; that replica is now the most loaded; they all move together. Oscillation.
Removing busy replicas from the pool. Cascading failure.
Health checks that run real generations. Wastes GPU capacity.
Mixing long and short requests in one pool.
Buffering at the LB. Destroys streaming.
Session affinity without a release valve. One popular conversation overloads a replica.
10. Hands-on exercise#
A. Simulate the algorithms. Build a simulator with N replicas, heavy-tailed request costs, and a configurable prefix-sharing rate. Implement round-robin, least-connections, least-KV, and prefix-aware. Reproduce the table in section 7.
B. Power of two. Add “global least loaded” to your simulator with multiple concurrent LB instances. Observe the herding/oscillation. Compare to power-of-two-choices.
C. Affinity vs balance. Sweep the release-valve threshold (the KV utilization above which you abandon affinity) from 0.5 to 1.0. Plot prefix hit rate and p99 TTFT. Find the optimum.
D. Real routing. Put a simple custom router in front of 2-4 real vLLM replicas. Route on
vllm:gpu_cache_usage_perc scraped from /metrics. Compare TTFT distribution against
round-robin.
E. Segregation. Add a long-request pool to your simulator. Measure the p99 TTFT improvement for short requests.
11. Interview questions#
- Why does round-robin perform badly for LLM serving? Give two independent reasons.
- What is the best load signal for an LLM replica, and why?
- What is power-of-two-choices and what problem does it solve?
- How do you balance prefix cache affinity against load balance?
- Why shouldn’t you remove a saturated replica from the pool?
- Why segregate long-context requests, and how much does it help?
- What must a load balancer be configured for, to handle streaming responses?
12. Further reading#
- [FUNDAMENTAL] Mitzenmacher, “The Power of Two Choices in Randomized Load Balancing” (2001)
- [ESTABLISHED] Zheng et al., “SGLang” (2023) — cache-aware scheduling
- [ESTABLISHED] Kubernetes Gateway API Inference Extension (
InferencePoolv1) and llm-d’s endpoint picker - [REFERENCE] Envoy load balancing documentation (for the general algorithms)
- Next: 07 — Autoscaling and cold starts