skip to content

Why does round-robin load balancing misbehave across LLM replicas?

level: middleimportance: should knowfreq 50%

answer

  1. requests are not equal in cost
  2. output length is invisible to the proxy
  3. counts balanced, work not balanced
  4. long streams make imbalance persist
  5. least outstanding, then metrics-aware

basics

~20 s

Round-robin equalizes request counts, but LLM requests differ in cost by orders of magnitude and run for seconds to minutes. One replica ends up holding all the long generations while another idles, so tail latency blows up even though every replica got the same number of requests.

solid answer

~50 s

Round-robin assumes requests are short and roughly equal — true for a web tier, false for inference. A 30-token classification and a 4,000-token summarization arrive as one request each, but the second occupies a slot in the running batch for a hundred times longer and consumes far more KV-cache. Dealing out requests in turn therefore equalizes counts while leaving work wildly uneven: one replica sits at its concurrency cap with a queue forming while its neighbour has free slots. Because generations are long-lived and streamed, the imbalance persists rather than averaging out. The minimum fix is to balance on outstanding requests instead — Envoy's LEAST_REQUEST policy, nginx's `least_conn` — which at least tracks what is still in flight. The better fix is an inference-aware router that reads each replica's real state from its `/metrics` (queued requests, running batch size, KV-cache occupancy) and sends the next request to the replica with genuine capacity.

code

nginx · 6 lines
nginx
upstream llm_backends {
    least_conn;
    server llm-0.svc:8000 max_conns=64;
    server llm-1.svc:8000 max_conns=64;
    server llm-2.svc:8000 max_conns=64;
}

go deeper

for a junior

Know that LLM requests vary hugely in how long they run, so handing them out in strict rotation is not the same as sharing the work evenly.

for a middle

Explain the two broken assumptions — uniform cost and short duration — and name the concrete improvement: balance on outstanding requests, such as Envoy LEAST_REQUEST or nginx least_conn.

for a senior

Demonstrate the diagnosis: per-replica running-batch and queue graphs showing one pod pinned, and the move to a router that reads engine metrics plus a shared queue and a per-replica in-flight cap.

for a principal

Own the routing layer as a platform decision — build versus adopt an inference-aware gateway, what it costs in operational surface, and how retry and load-shedding policy interacts with it during a fleet-wide overload.

## The assumption round-robin encodes Round-robin is optimal under one condition: request cost is roughly uniform and short relative to the arrival rate. Then dealing requests in turn approximates dealing work in turn, and any momentary imbalance washes out within a few hundred milliseconds. Every conventional HTTP load balancer defaults to it because that condition holds for typical web traffic. LLM traffic violates both halves of the condition. ## Cost variance Generation cost is dominated by output length, which the caller does not declare and the balancer cannot see. A yes/no classification emits perhaps 5 tokens; a document summary emits 2,000; an agent trace emits 10,000. Prompt length varies as widely and drives prefill cost and KV-cache footprint. So two requests that look identical to a proxy — same route, same method, same size class of body — can differ by two or three orders of magnitude in GPU seconds and in cache blocks held. ## Duration A web request finishes in milliseconds; an LLM generation occupies a slot in the running batch for seconds to minutes, streaming the whole time. That means imbalance accumulates instead of dissipating. If replica A happens to receive three long generations early, it holds them while the round-robin pointer keeps handing it new work at exactly the same rate as everyone else. Its running batch fills, its KV-cache occupancy climbs, and arrivals begin to queue — while replica B, which happened to draw short requests, has idle capacity the whole time. The user-visible result is a bimodal latency distribution. Median looks fine because most requests land on unloaded replicas; p95 and p99 are dreadful because the unlucky ones queue behind a saturated replica. Fleet-level dashboards make this hard to see: average GPU utilization looks healthy, average queue depth looks low, and only the per-replica breakdown reveals that one member is pinned while others coast. ## Least-outstanding-requests The cheap improvement is to balance on in-flight count rather than cumulative count. Envoy calls this LEAST_REQUEST (with the power-of-two-choices variant that samples two hosts and picks the less loaded, which is nearly as good and much cheaper than scanning every host); nginx offers `least_conn`. Because a long generation keeps its request or connection open for its entire duration, in-flight count is a far better proxy for occupied capacity than a round-robin counter. Two caveats. First, connection-based policies only approximate request count when connections are not multiplexed; with HTTP/2 many streams share one connection, so a connection-counting policy sees one unit where there are twenty streams — prefer a request-aware policy. Second, in-flight count still treats a 5-token request and a 5,000-token request as one unit each. It is a much better proxy, not a correct measure. ## Inference-aware routing The engines already publish the truth. vLLM exposes `vllm:num_requests_running`, `vllm:num_requests_waiting` and `vllm:kv_cache_usage_perc`; TGI and Triton publish their own queue and batch series. A router that consumes those — polling them, or receiving them pushed — can send the next request to the replica with actual free slots and cache headroom, and can refuse to pile onto a replica that is already preempting. This is what inference-aware gateways such as the Kubernetes Gateway API Inference Extension exist to do, and it is what any in-house router should approximate. A useful middle ground when you do not want a metrics-driven router: cap in-flight requests per replica at the client or proxy, so no replica can be handed more than it can hold, and let excess wait in one shared queue rather than in per-replica queues. A single shared queue is strictly better than N private ones because work goes to whichever replica frees a slot first. ## Where the balancer must not overreach One thing the balancer should not do is retry aggressively on slow responses. On a saturated fleet, a timeout-and-retry policy duplicates the most expensive requests at exactly the moment capacity is scarcest, and the duplicate usually lands on another already-loaded replica. Retries here need a strict budget, and streaming responses complicate them further because a retry after the first token has already shipped bytes to the client. ## What to watch Always graph the per-replica distribution, not the fleet average: running batch size, queued requests and KV-cache occupancy, per pod. A healthy fleet shows those curves tracking one another. A round-robin fleet under mixed traffic shows one line pinned to the ceiling and the rest well below it — that single picture is the whole argument.

  • Why is least-outstanding-requests still only an approximation?
    Because it counts requests, not work. A replica holding three 5,000-token generations and one holding three 20-token replies look identical to it, even though their remaining GPU work differs by orders of magnitude. It fixes the worst of round-robin's blindness — persistent occupancy — without seeing cost, which is why reading the engine's queue and KV-cache metrics is the real answer.
  • Does a connection-based policy like nginx least_conn work if clients use HTTP/2?
    Poorly. HTTP/2 multiplexes many streams over one connection, so a connection-counting policy sees a single unit where there may be dozens of concurrent generations. Either terminate at a proxy that balances per request, or use a request-aware policy such as Envoy's LEAST_REQUEST, which counts active requests rather than sockets.
  • How should the router behave when every replica is at its cap?
    Hold the excess in one shared queue and apply a bounded wait rather than pushing it into per-replica queues, so the next freed slot anywhere takes the oldest request. Past the bound, shed load explicitly with a retryable status and a hint of when to come back. Silently queueing without a limit converts an overload into a timeout storm.

Round-robin is a supermarket policy of sending every fifth customer to lane three regardless of trolley size; least-outstanding-requests at least looks at how many people are already in each lane.

saying these in an interview costs you the question

  • Assumes LLM requests cost roughly the same
  • Uses fleet-average metrics to judge balance
  • Counts connections under HTTP/2 as concurrency
  • Retries slow generations aggressively during overload
  • Gives each replica its own deep private queue

context