Chapter 6 · Distributed inference
Production systems and serving benchmarks
6.4

Production systems and serving benchmarks

Everything so far in this chapter happens inside one deployment of one model. A production serving system is the layer above that: it decides which replica a request lands on, when it runs, and how the fleet scales, and its defining problem is the tail rather than the average. Clockwork, a model serving system that predates LLMs, made the case that this layer can be built on prediction instead of reaction: existing serving architectures used “well-known reactive techniques” against common-case latency but could not curtail the tail latency caused by unpredictable execution times, while DNN inference itself “has deterministic performance.” Built bottom-up from those predictable execution times, with a centralized scheduler in place of reactive workers, Clockwork supported thousands of models while meeting “100 ms latency targets for 99.997% of requests.” The lesson carries directly into LLM serving: latency targets are met by scheduling, not only by fast kernels.

llm-d is what that layer looks like for LLMs today. Engines like vLLM and SGLang handle running the model on accelerators; llm-d provides “orchestration and optimizations above model servers,” deployed on Kubernetes, and its feature list reads like this chapter turned into routing decisions: prefix-cache and load-aware request balancing, tiered KV cache offloading with global indexing of cache state, prefill/decode disaggregation and wide expert parallelism for the largest models, and SLO-aware autoscaling driven by real-time inference signals. Where a request lands now matters as much as how fast the engine runs it: llm-d reports 3x higher output throughput and 2x faster TTFT from prefix-cache-aware routing compared to round-robin, measured on Llama 3.1 70B across 4 AMD MI300X GPUs, a gain that comes from routing each request to a replica that already holds its prefix in cache.

That routing layer is also becoming a standard rather than a product. Gateway API Inference Extension, an official Kubernetes project, extends any gateway that supports the Gateway API and Envoy’s external processing protocol into an inference gateway: a load balancer coupled with an endpoint picker that chooses a replica using metrics and capabilities the model servers themselves report, such as prefix cache status or which LoRA adapters a replica has loaded. The same layer carries policy that a plain load balancer has no vocabulary for: serving priority, so a latency-sensitive chat model can outrank a latency-tolerant summarization model, and incremental rollouts of new model versions by splitting traffic on model names.

Scaling the fleet has a cold start problem of its own: a new replica must load its model’s weights from a checkpoint before it can serve a single token. ServerlessLLM, an OSDI 2024 system, attacks that startup path by harnessing “the substantial near-GPU storage and memory capacities of inference servers” to keep checkpoints close to the accelerators. Its three mechanisms are a loading-optimized checkpoint format with a multi-tier loading system that uses the full bandwidth of the storage hierarchy, live migration of running inference so a new request can claim a server that already holds its checkpoint, and a scheduler that places each model on the server that minimizes its time to start. The paper reports “reducing latency by 10 - 200X across various LLM inference workloads” against state-of-the-art serverless systems.

Measuring such a system takes more than a throughput number. The etalon benchmark framework works with the per-request latency metrics this chapter has already met: time to first token, defined as “the time taken between arrival and first output token,” which includes “both scheduling delay and prompt processing time”; time between tokens, the gap between two consecutive output tokens of a streaming response; and time per output token, total generation time divided by the number of output tokens. Etalon’s starting observation is that even these “fail to fully capture the nuances of LLM inference,” leaving an incomplete picture of user-facing performance in real-time applications like chat. What a production system optimizes is under : DistServe defines per-GPU goodput as “the maximum request rate that can be served adhering to the SLO attainment goal (say, 90%)” for each GPU provisioned. By these definitions a system can raise its throughput while lowering its goodput, simply by letting a slice of requests blow through their latency targets.

One request's lifetime, and the metrics cut from it. TTFT runs from arrival to the first output token and swallows both queueing and prompt processing. TBT is the gap between any two consecutive tokens; TPOT is the whole decode span divided by the token count. Goodput is not on this timeline at all: it counts how many such requests per second finish inside their latency targets. Definitions from the etalon documentation and DistServe.

The load those metrics are measured under matters as much as the metrics themselves. Clockwork was evaluated “using production trace workloads,” and BurstGPT makes the same possible for LLM serving research: a public trace of real GPT-3.5 and GPT-4 traffic served on Azure, covering 121 consecutive days and roughly 5.29 million requests, each entry recording request and response token counts and whether the call came through a conversation or the API. The trace’s own overview plots weekly and daily periodicity in request arrivals, structure that a uniform synthetic load has none of, and its usage notes suggest scaling the trace’s average request rate to the evaluation setup rather than discarding the arrival pattern.

ServeGen carries that argument further. Built on a characterization of workloads “collected from our worldwide cloud inference serving service” at Alibaba, covering language, multimodal, and reasoning models, it argues that prior analyses were too limited in scale and scope to capture how real traffic behaves, and generates realistic workloads “by composing them on a per-client basis” rather than drawing from one aggregate distribution. Its production use case makes the stakes concrete: benchmarking with ServeGen “avoids 50% under-provisioning compared to naive workload generation.” A capacity plan validated against uniform synthetic load can simply be wrong.

Comparing systems across vendors needed a referee, and MLPerf became it. The MLPerf Inference paper (ISCA 2020) describes a field where over 100 organizations were building inference chips and existing systems spanned “at least three orders of magnitude in power consumption and five orders of magnitude in performance,” and answers with rules and best practices “to ensure comparability across systems with wildly differing architectures”; its first call for submissions drew more than 600 reproducible measurements from 14 organizations, of which audit tests cleared 595 as valid, spanning four orders of magnitude of performance. What makes it a benchmark rather than a leaderboard is that a submission is a scenario, not a number. MLPerf Inference defines four of them. Single-stream injects one query at a time and sends the next only when the last completes, scored on the query stream’s 90th-percentile latency. Multistream sends a new query of N samples at a fixed interval, 50 to 100 milliseconds depending on the task, and is scored on how many streams the system sustains inside that interval. Server is treated below. Offline sends a single query containing every sample, lets the system process them in any order, and is scored in samples per second. The paper’s own finding is that “performance can vary drastically under these scenarios and their corresponding metrics,” which is the reason a serving number is worth little until you know which scenario produced it.

The server scenario is the one this chapter has been describing without using the name. Its queries carry one sample each and arrive “in accordance with a Poisson distribution”; the system must answer each inside a benchmark-specific latency bound between 15 and 250 milliseconds; no more than 1% of queries may exceed it for the vision tasks and no more than 3% for translation; and the metric is “the Poisson parameter that indicates the queries-per-second (QPS) achievable while meeting the QoS requirement.” That is goodput, written down in 2020 for image classifiers. What makes it enforceable is a component MLPerf ships and the submitter is not allowed to touch: the Load Generator, “a traffic generator for MLPerf Inference that loads the SUT and measures performance,” which produces query traffic according to the scenario’s rules, records the queries and responses, and decides at the end whether the run was valid at all. The paper calls that ability to simulate the realistic behavior of the system under test “unique among AI benchmarks.” Submissions then land in a closed division under strict rules, or an open division that lets the submitter change the model. In the next subsection, vLLM’s benchmark client will synthesize arrival times from a Poisson process and report a 99th-percentile TTFT beside its throughput. That is the server scenario, run by hand, on one deployment.

MLPerf Endpoints is the same idea aimed at LLM serving: a submitted system is a measured curve rather than a single number, each point one real test at a fixed concurrency, relating system throughput, interactivity in tokens per second per user, and 95th-percentile TTFT. As MLCommons puts it, as load rises “total throughput goes up and per-user speed comes down. That is the core tradeoff in serving AI,” which is goodput’s tradeoff restated as a purchasing decision.

Running the measurement#

Those definitions turn operational the moment something computes them. vLLM ships a benchmark client that drives a running server and reports exactly these numbers; the documentation’s own example points it at a ShareGPT dump and ten prompts.

driving a running server, from vLLM's benchmark CLI documentation
# download dataset
# wget https://huggingface.co/datasets/anon8231489123/ShareGPT_Vicuna_unfiltered/resolve/main/ShareGPT_V3_unfiltered_cleaned_split.json
vllm bench serve \
  --backend vllm \
  --model NousResearch/Hermes-3-Llama-3.1-8B \
  --endpoint /v1/completions \
  --dataset-name sharegpt \
  --dataset-path <your data path>/ShareGPT_V3_unfiltered_cleaned_split.json \
  --num-prompts 10

The dataset name is drawn from a list that includes BurstGPT, so the trace above is not a separate exercise, it is an argument. The arrival process is a separate concern from the prompt count. --request-rate is requests per second, and its default of infinity fires every request at time zero, while any finite rate synthesizes arrival times from a Poisson process; --burstiness replaces that Poisson process with a gamma distribution, values below one producing burstier arrivals and values above one more uniform ones. A finite request rate read against a tail-latency bound is MLPerf’s server scenario rebuilt on one deployment; what this client does not bring is the load generator’s verdict on whether the run counted. --max-concurrency is a different lever again, capping how many requests are allowed to execute at once in order to simulate an upstream component enforcing a limit, and the documentation is careful to note that combining it with a request rate can push the achieved rate below the requested one when the server cannot keep up. Load is a shape, not a number.

What comes back has a fixed form: run totals and three throughput lines covering the whole benchmark, then per-request latency split into three blocks, each reported as a mean, a median, and a 99th percentile. Ten requests over 5.78 seconds is far too small a sample to conclude anything from, which is exactly why it is worth reading as a form rather than as a result.

the serving benchmark report, from vLLM's benchmark CLI documentation
============ Serving Benchmark Result ============
Successful requests:                     10
Benchmark duration (s):                  5.78
Total input tokens:                      1369
Total generated tokens:                  2212
Request throughput (req/s):              1.73
Output token throughput (tok/s):         382.89
Total token throughput (tok/s):          619.85
---------------Time to First Token----------------
Mean TTFT (ms):                          71.54
Median TTFT (ms):                        73.88
P99 TTFT (ms):                           79.49
-----Time per Output Token (excl. 1st token)------
Mean TPOT (ms):                          7.91
Median TPOT (ms):                        7.96
P99 TPOT (ms):                           8.03
---------------Inter-token Latency----------------
Mean ITL (ms):                           7.74
Median ITL (ms):                         7.70
P99 ITL (ms):                            8.39
==================================================

The last two blocks look redundant and are not. vLLM measures at the benchmark client and defines each block by where it takes its timestamps: time to first token is the interval from sending a request to receiving its first streamed output; inter-token latency records the gap between consecutive streamed outputs, aggregated across every successful request; and time per output token is computed once per request, as the end-to-end latency minus the TTFT divided by one less than the number of output tokens, and only then aggregated. With ordinary decoding a streamed output carries one token and the two blocks nearly agree, which is what the numbers above show. With speculative decoding one output can carry several accepted tokens at once and they part company: in the documentation’s worked example the mean inter-token latency is 40 ms, because the three tokens bundled into a single streamed output add no gaps of their own, while the time per output token is 20 ms, because that metric spreads the 80 ms between the first token and the end of the request across the four tokens that follow it. The engine did not change between those two numbers. The measurement did, which is why vLLM warns that metric terminology is not standardized across tools and asks readers to compare measurement points and formulas rather than names.

Goodput is a flag on the same command. --goodput takes service level objectives as key-value pairs, where the key is one of ttft, tpot, or e2el and the value is a bound in milliseconds, and the help text sends the reader to the DistServe paper for what the word means. That flag is the whole distance between a target written on a slide and a target the measurement enforces: a run with no objective reports throughput, and a run with one reports how much of that throughput was worth anything.

One run is still one point. vllm bench sweep serve starts a server and iterates the benchmark across parameter combinations, three runs each by default; its workload variant explores levels of request rate or concurrency specifically to find the tradeoff between latency and throughput, working outward from serial inference at one extreme to every request at once at the other; and a Pareto plot over the results “helps pick configurations that balance per-user and per-GPU throughput,” with the frontier showing the best achievable pairs across the runs. That is the MLPerf Endpoints curve, drawn on one deployment instead of across vendors. vLLM is candid that its own client is aimed at evaluating specific features and at regression testing, and points at GuideLLM for benchmarking production servers, but the shape of the answer does not change: a serving system is a frontier, and a single throughput number is a point somebody chose from it.