Skip to main content

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.

TermMeaning
TokenA piece of text the model reads or writes, about three quarters of a word.
PromptThe text a user sends: the conversation so far plus the new question.
PrefillThe step where the model reads the whole prompt before it writes anything. It is the slow part of a long prompt.
DecodeThe step where the model writes the answer, one token at a time.
KV cacheThe 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.
PrefixThe start of a prompt that repeats, such as a shared system prompt or the earlier turns of a conversation.
PoolThe group of machines one run uses, with the engines, caches and routers it starts.
TierA place that holds KV cache: GPU memory, host memory (RAM), local disk, or a store shared between machines.
StackThe 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.

The five ways a pool shares KV cache

  1. 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.
  2. 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.
  3. 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.
  4. 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.
  5. 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.

Where every prompt token comes from

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.

The lifecycle of a distributed KV run

  1. 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.

  2. 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 with ib_write_bw and ib_write_lat on an RDMA pool.

  3. 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.

  4. 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.
  5. 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).

  6. Measure. Every machine samples its engine counters, its GPU and host telemetry, and its network counters while the workload runs.

  7. 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.

MeasureWhat it tells you
Time to first token: mean, median (P50), P90, P99How long users wait before the answer starts. P90 means 9 in 10 requests were faster than this.
Time per output tokenHow fast the answer is written once it starts.
Output throughput, in total and per GPUHow many answer tokens the pool writes per second. Per GPU makes pools of different sizes comparable.
Goodput under an SLOThe rate of requests that met a latency target, such as an MLPerf or DistServe preset.
Prompt token sourcesWhere every prompt token came from, as described above.
KV transfer bandwidthHow fast KV cache moved between machines, compared with the measured line rate.
GPU KV usage and working setHow full the GPU KV cache was, and how the workload compares with its size.
EnergyOutput tokens per joule, when the machines report their power.
Switch countersWhat the switches between the machines saw: congestion marks, pause frames, drops, buffer peaks, FEC errors and optical power, per port and per queue.
Energy with networkOutput tokens per joule with the pool's share of its switches' energy added, shown beside the machines-only figure.

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.

MeasureSource
Latency, throughput, goodputThe benchmark client's own record of every request: when it was sent, when its first token arrived and when it finished.
Engine and cache countersEach 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 countersvllm-router, the SGLang router, the Dynamo frontend or the llm-d endpoint picker, from their own metrics endpoints.
GPU utilisation and memoryall-smi.
GPU activity and PCIe traffic on NVIDIANVIDIA dcgm-exporter.
GPU power, temperature, memory activity, PCIe bandwidth and ECC errors on AMDAMD Device Metrics Exporter v1.5.2.
NIC, RDMA and TCP countersPrometheus node_exporter v1.12.1, its infiniband, ethtool, netdev and netstat collectors.
CPU, memory and drive I/Onode_exporter v1.12.1, its cpu, meminfo and diskstats collectors.
NVMe healthPrometheus community smartctl_exporter v0.14.0.
Node powerThe machine's BMC, read with ipmitool dcmi or Redfish.
Switch counters and switch powerThe switch's own counters, read by the pool's head machine through gNMI (gnmic v0.49.0), RESTCONF, the SONiC counters database, or the Dell Enterprise SONiC command line. See Switch counters.

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 stat on 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.

Switch counters​

A team connects its switches once, with a read-only login. See Network Switches. The pool's head machine then reads every connected switch that a pool machine's RDMA network card is cabled to, and the uplinks between them when the pool spans more than one switch, from the start of the pool to its end.

  • Which reader. When a switch is connected, a machine cabled to it tests every reader once, and each counter goes to the first reader that read it: gNMI, SONiC gNMI, RESTCONF, the SONiC counters database over SSH, then the command line. gNMI and RESTCONF carry the switch's own timestamps.
  • How often. Every 10 seconds, or every 60 seconds through the command line, where one full reading takes about 34 seconds on a Dell Z9864F-O64. The report names the reader and the interval of every switch.
  • How a job value is made. Each counter has a reading kind. A running total reports its change over the job. A level, such as optical power, reports its last reading. A peak the switch keeps until it is cleared, such as a queue watermark, reports its last reading only if it rose during the job, since a peak that did not rise may be older than the job. A last-event time, such as the last link-down, reports how many times it changed. A running total that went down during a job was cleared or restarted, and is reported as a reset, never as a value.
  • Clocks. The head records the switch clock's offset from its own, so the timeline places every reading on the job's clock.
  • Agreement. For every network card, the report compares what the card sent with what its switch port received, and the other way round. They agree when both directions are within 2 percent. Below 1 MB/s both ways, the link counts as idle. Some drivers count RDMA traffic outside the card's ordinary byte counters, so the comparison uses whichever host total matches.
  • Network energy. A switch's energy in a job is its PSU input power, averaged over the job's readings, times the job's length. The pool's share is the ports it used out of the ports that were up. Tokens per joule with network adds that share to the machines' energy. It is shown only when every machine of the job reported node power, because a switch's wall power added to GPU-only energy would mix two different measures. When a switch does not report its ports, none of its energy is attributed, the job gets no network figure, and the report says so.

A counter a switch does not expose is reported as not measured, never as zero. A counter read through a field name not yet confirmed on a live switch of that operating system is shown, marked as unconfirmed.

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.

How a distributed KV run is measured

  1. Warm-up. The first seconds of the load are discarded. The caches fill and the engines settle. The default is 300 s.
  2. 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.
  3. 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:

SettingDefaultRange
Load duration1200 s120 to 14400 s
Warm-up300 s0 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, when no figure in it was read through an engine metric name that has not been confirmed on a live engine, when every switch its pool's RDMA network cards are cabled to was read, and when no switch figure was read through an unconfirmed field name. Switches the team ignored and switches Metrum AI Bench cannot read are left out of that rule. A switch first found after the pool ended does not hold back that pool's result.

How KV cache travels between machines​

A stack moves KV cache over one of two kinds of path.

How KV cache travels between machines

  • 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:

TopicNVIDIAAMD
Engine imagesvLLM, SGLang and TensorRT-LLM CUDA imagesvLLM and SGLang ROCm images (for MI300X, the mi30x builds)
TensorRT-LLMSupportedNot available
GPU tools read by preflightnvidia-smirocm-smi or amd-smi
GPU telemetryall-smi, with dcgm-exporter for activity and PCIeall-smi for utilisation and memory, AMD Device Metrics Exporter for power, temperature and the rest
RDMA transfer librariesNIXL, MooncakeNIXL, Mooncake, AMD MoRI-IO
Multi-node NVLink checksOn GB200 and GB300Not 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.