Distributed KV Cache Methodology
This page explains how Metrum AI Bench measures a KV cache that is shared, offloaded or handed over across several machines. It covers what is measured, how a run works, and when a result counts as valid. For the steps to set up a run, see the Distributed KV Cache User Guide.
Use a distributed KV benchmark to answer questions like these:
- Does sharing KV cache between machines make the first token arrive sooner?
- How much of every prompt did the pool avoid computing again?
- Is it better to split prefill and decode across machines, or to keep them together?
- Which network path did the KV cache take, and how fast was it?
Key words
These terms appear throughout the page.
| Term | Meaning |
|---|---|
| Token | A piece of text the model reads or writes, about three quarters of a word. |
| Prompt | The text a user sends: the conversation so far plus the new question. |
| Prefill | The step where the model reads the whole prompt before it writes anything. It is the slow part of a long prompt. |
| Decode | The step where the model writes the answer, one token at a time. |
| KV cache | The keys and values the model computes for every prompt token during prefill. If they are kept, the next turn of the conversation does not compute them again. |
| Prefix | The start of a prompt that repeats, such as a shared system prompt or the earlier turns of a conversation. |
| Pool | The group of machines one run uses, with the engines, caches and routers it starts. |
| Tier | A place that holds KV cache: GPU memory, host memory (RAM), local disk, or a store shared between machines. |
| Stack | The software combination that runs a strategy, for example vLLM with LMCache and a Valkey store. |
| Time to first token (TTFT) | How long a user waits from sending a prompt to seeing the first word of the answer. |
The five strategies
A pool shares KV cache in one of five ways. Every run uses one of them.
- Each machine keeps its own cache. This is the aggregated baseline. Every machine reuses only the KV cache in its own GPU memory. Nothing moves between machines. Every other strategy is compared with this one.
- Spill to host memory and disk. This is offload per node. When GPU memory fills up, each machine moves older KV cache to its own RAM or disk and reads it back later. Nothing moves between machines.
- Share one cache pool across machines. This is the shared KV pool. Every machine can read KV cache that any machine wrote, through a shared store such as Valkey or Mooncake, or through a shared filesystem. A conversation can move to another machine without computing its history again.
- Split prefill and decode across machines. This is disaggregated prefill and decode. Some machines only read prompts. They hand the KV cache to other machines, which only write answers.
- Split and offload. Prefill and decode are split, and the KV cache is also kept in host memory, disk or a shared tier, so later turns can reuse it.
Where every prompt token comes from
For every prompt token, the run records where its KV cache came from. It either came from a cache tier, or prefill computed it again.
The sources are:
- GPU cache. The engine's own GPU memory already held it.
- Host memory. The machine's own RAM held it.
- Storage. The machine's local disk or a disk tier held it.
- Remote KV. Another machine or a shared store held it. This is the benefit only a distributed pool can give.
- Computed. No tier held it, so prefill computed it again.
Three numbers describe how well the pool reused the prompt:
- Reusable is the share of prompt tokens that an ideal cache of unlimited size could have served. It follows from the workload itself, before any engine runs.
- Reused is the share the pool actually served from any cache tier.
- Reuse captured is reused divided by reusable. A value of 100% means the pool lost no reuse to eviction, misses or routing.
For example, if 82% of the prompt tokens were reusable and the pool reused 78%, the pool captured 95% of the reuse it could have had.
The counts come from the engines' own counters on every machine, so a token is never counted as reused unless an engine reports it. A run whose engines do not report every source shows the missing share as not attributed, never as computed.
How one run works
A run moves through seven steps.
-
Plan. The platform chooses the machines, how many GPUs each engine uses, the ports and the size of every tier. It refuses a plan that cannot work, for example an engine that needs more GPUs than a machine has, and says why.
-
Preflight. Before anything starts, every machine proves it is ready. The checks include:
- The clocks agree.
- The GPUs belong to one family and answer their vendor tool.
- The RDMA ports the pool uses are active, and the network MTU matches.
- The host tier fits in the memory the machine can lock.
- No pool port is open to the internet, and no secret shows in a process list.
Preflight also measures the network line rate with
iperf3, and withib_write_bwandib_write_laton an RDMA pool. -
Launch. The head machine starts the shared store if the stack has one. Every machine then starts its cache server and its engines. A router starts in front of the engines.
-
Verify. The pool proves it works before it measures anything:
- A request through the router returns tokens.
- On a prefill and decode stack, a decode machine generates tokens from KV cache that a prefill machine computed.
- On a shared pool, one machine loads KV cache that another machine computed. The reader's own counters show it.
-
Benchmark. The workload runs. The load generator sends multi-turn conversations, either closed loop (each user sends the next turn when the last one ends) or open loop (requests arrive at a set rate).
-
Measure. Every machine samples its engine counters, its GPU and host telemetry, and its network counters while the workload runs.
-
Tear down. The pool flushes and removes the shared store, stops every process and deletes every disk tier directory. The next run starts from an empty cache.
If a preflight or verification check fails, the run stops before it measures anything, and the report names the check that failed.
The workload
The workload decides how much of every prompt can be reused. Four shapes are available:
- Agentic data set replay. Replays recorded coding-agent sessions. Each user sends the next round when the last one ends.
- Shared prompt, per-user history. Every user shares one system prompt and carries its own history.
- Synthetic sessions. Users come in groups that share a prefix. You set the prefix length, the number of turns, the input and output length, how much the context grows each turn, and the think time between turns. The reusable share follows from these settings.
- Public trace replay. Replays a published inference trace at its own timestamps, with prompts built so its reuse is reproduced.
The working set ratio compares the KV cache the workload needs with the KV cache the pool's GPUs can hold. A ratio above 1 means the conversations do not fit in GPU memory, so reuse depends on the tiers below the GPU. That is where offload and shared pools can help.
What is measured
Every result is reported per test point: one scenario at one concurrency, repeated as many times as the workload asks.
| Measure | What it tells you |
|---|---|
| Time to first token: mean, median (P50), P90, P99 | How long users wait before the answer starts. P90 means 9 in 10 requests were faster than this. |
| Time per output token | How fast the answer is written once it starts. |
| Output throughput, in total and per GPU | How many answer tokens the pool writes per second. Per GPU makes pools of different sizes comparable. |
| Goodput under an SLO | The rate of requests that met a latency target, such as an MLPerf or DistServe preset. |
| Prompt token sources | Where every prompt token came from, as described above. |
| KV transfer bandwidth | How fast KV cache moved between machines, compared with the measured line rate. |
| GPU KV usage and working set | How full the GPU KV cache was, and how the workload compares with its size. |
| Energy | Output tokens per joule, when the machines report their power. |
Each mean comes with a 95% confidence interval, computed with Student's t-distribution over the valid repeats. A test point with fewer than three valid repeats is marked preliminary.
How each measurement is taken
Every figure comes from a public, versioned tool wherever one publishes it. Each machine records, in the run's manifest, the tool and build that measured its counters, so a report names the source of every number.
| Measure | Source |
|---|---|
| Latency, throughput, goodput | The benchmark client's own record of every request: when it was sent, when its first token arrived and when it finished. |
| Engine and cache counters | Each engine's and LMCache server's own Prometheus endpoint. Every series name is checked on a live engine before a result that uses it can be published. |
| Router counters | vllm-router, the SGLang router, the Dynamo frontend or the llm-d endpoint picker, from their own metrics endpoints. |
| GPU utilisation and memory | all-smi. |
| GPU activity and PCIe traffic on NVIDIA | NVIDIA dcgm-exporter. |
| GPU power, temperature, memory activity, PCIe bandwidth and ECC errors on AMD | AMD Device Metrics Exporter v1.5.2. |
| NIC, RDMA and TCP counters | Prometheus node_exporter v1.12.1, its infiniband, ethtool, netdev and netstat collectors. |
| CPU, memory and drive I/O | node_exporter v1.12.1, its cpu, meminfo and diskstats collectors. |
| NVMe health | Prometheus community smartctl_exporter v0.14.0. |
| Node power | The machine's BMC, read with ipmitool dcmi or Redfish. |
Each exporter build is pinned by checksum or image digest, runs on the machine's loopback address only for the life of the pool, and is used only when it reports the pinned version. When an exporter cannot run, the agent's own reader takes its place and the manifest records that and why.
A few readings have no public exporter and stay on the agent's own readers, each named in the manifest:
- Per-process CPU time of the engines, from
/proc/<pid>/stat. - Memory-controller bandwidth, from
perf staton the uncore or UMC events, which node_exporter's perf collector cannot read. - RDMA counters node_exporter does not publish, for Broadcom NICs RoCE
discards, retries exhausted, and reads and writes sent, read from the
driver's
hw_counters. - BMC node power. An exporter exists (ipmi_exporter), but it needs FreeIPMI, which these machines do not carry.
A counter a machine's NIC driver does not expose is reported as not measured,
never as zero. The Broadcom bnxt_re 233.x driver, for example, has no PFC
pause duration counters.
Energy. Each machine counts once, at the best power it measured: its BMC node power, else the node power its telemetry sampled, else its GPUs. Each power series is integrated over its own samples, leaving out gaps longer than 30 seconds.
all-smi reports MI300X GPU power in kilowatts under a metric named in watts (0.148 for a GPU drawing 148 W, in versions 0.22.0 and 0.26.0), so AMD GPU power comes only from AMD's exporter. GPU power recorded before this change on MI300X machines was corrected by the same factor of 1,000.
Load duration, warm-up and the measured window
A measurement taken while the pool is still filling does not describe the pool at work. A job with a held load therefore has three phases.
- Warm-up. The first seconds of the load are discarded. The caches fill and the engines settle. The default is 300 s.
- Measured window. The rest of the load duration, from the end of the warm-up to the end of the load. Every published figure comes only from requests that start inside this window.
- Drain. A request that is still running when the load ends finishes and counts, because it started inside the window. No new request starts.
You set two values in the Measurement step of the workload form:
| Setting | Default | Range |
|---|---|---|
| Load duration | 1200 s | 120 to 14400 s |
| Warm-up | 300 s | 0 to 3600 s, and less than half the load duration |
With the defaults, 900 s (15 min) are measured after a 300 s (5 min) warm-up. The report states this under the headline, for example Measured over 15 min after a 5 min warm-up.
Only the Agentic data set replay workload shape holds its load for a set time. The synthetic, public trace and shared prompt shapes set their own length, so they have no load duration and no warm-up. Those results, and every result recorded before the warm-up setting existed, read Measured over the whole run.
Stability is reported, not graded
The report does not pass or fail a job on its stability. It states two facts about the measured windows:
- Variation is the coefficient of variation of output throughput across the measured windows: the standard deviation divided by the mean.
- Drift is the change in output throughput from the first to the last measured window, as a share of their mean, from a fitted straight line. A negative value means throughput fell.
Read them beside the figures. A high variation or a large drift tells you the load was not level, and you can decide what that means for your question.
When a result carries a warning
A result carries a warning in three cases only. The report shows the warning sentence beside the result.
- The load ran out of sessions. The data set had fewer sessions than the load duration needed, so users stopped before the load ended and the measured window is not full.
- More than 1% of the requests in the measured window failed.
- Less than 300 s was measured. A shorter window is too brief to average out normal variation.
A warning does not make a job invalid. It tells you how much weight the figures can carry.
Where the method comes from
The method follows how established benchmarks separate a settling period from the period they measure:
- MLPerf Inference requires a performance run of at least 600 s, so short runs do not stand in for a system at work. See the MLPerf Inference policies.
- NVIDIA AIPerf runs a warm-up phase that it discards, and then measures a fixed profiling window. Its metrics come from that window only. See AIPerf.
- SemiAnalysis InferenceX AgentX replays agentic coding sessions and measures a fixed window after the cache is primed. See InferenceX.
- NVIDIA Perf Analyzer repeats measurement windows until the last three agree within a stability percentage. Metrum AI Bench reports the same idea as variation and drift instead of retrying. See Perf Analyzer.
When a result counts
A job is valid only when every check that applies to it passed and nothing went wrong while it ran. A job is invalid when:
- a machine stopped responding during the job;
- a machine's clock drifted past the limit;
- a counter went backwards, which means a process restarted;
- the KV cache took a different network path than the stack intended;
- a reader used KV cache written in a different format;
- the cache held data before the job that its declared starting state says it should not;
- the stack compresses KV cache and the run has no accuracy result.
An invalid job still appears in the report, marked invalid, with each reason listed. Its numbers are shown as diagnostics, never as results.
A valid result can be published only when every machine of its pool recorded which tool measured its counters, and when no figure in it was read through an engine metric name that has not been confirmed on a live engine.
How KV cache travels between machines
A stack moves KV cache over one of two kinds of path.
- RDMA (remote direct memory access), over RoCE or InfiniBand. The network card reads and writes GPU or host memory directly, without the CPU. NIXL, Mooncake and AMD MoRI use this path.
- TCP. KV cache passes through the CPU and the operating system's network stack. LMCache with a Valkey store uses this path.
Preflight measures the line rate of the path before the benchmark, and the network counters sampled during the job show which path the bytes actually took. When a stack asks for RDMA and the bytes crossed TCP, the job is invalid.
When the shared store sits on a public address, the platform carries its traffic over a TLS link that accepts only the head machine's key.
NVIDIA and AMD
The method is the same on both vendors. These parts differ:
| Topic | NVIDIA | AMD |
|---|---|---|
| Engine images | vLLM, SGLang and TensorRT-LLM CUDA images | vLLM and SGLang ROCm images (for MI300X, the mi30x builds) |
| TensorRT-LLM | Supported | Not available |
| GPU tools read by preflight | nvidia-smi | rocm-smi or amd-smi |
| GPU telemetry | all-smi, with dcgm-exporter for activity and PCIe | all-smi for utilisation and memory, AMD Device Metrics Exporter for power, temperature and the rest |
| RDMA transfer libraries | NIXL, Mooncake | NIXL, Mooncake, AMD MoRI-IO |
| Multi-node NVLink checks | On GB200 and GB300 | Not applicable |
A pool uses one GPU family only. Machines with different GPU models, or different vendors, never share one pool.
Limits of the method
- A result describes the machines, network and software versions of that run. The report records all of them, so another run can reproduce it.
- The reusable share assumes an ideal cache. Real engines evict by their own rules, so reuse captured below 100% is normal.
- A stack with no validated run yet has no verdict. Check the validated configurations before you rely on a stack for a decision.