Where the KV cache lives: storage tiering for LLM inference under real-world traffic
LLM serving runs out of GPU memory long before it runs out of compute. Using real conversation and production request traces, I measured what each storage tier below the GPU costs a live request and built a shared cache tier on GekkoFS that lets several vLLM instances reuse each other's work.
Reading an evicted prefix back from CPU memory took 73 ms on an A100, where recomputing it took 256 ms. The slower the GPU, the larger the saving, so cache tiering pays off most on cost-efficient hardware.
Controlled sweeps separated two properties that earlier evaluations change together. Serving survived a cut to a quarter of the disk bandwidth, but one millisecond of added per-operation latency multiplied time to first token several times over. Decode latency never moved.
With four vLLM instances on two nodes sharing one GekkoFS namespace, an instance reloads KV that another instance computed instead of recomputing it. After tuning, the shared tier tracks the private local-disk baseline.
Tuning the backend and raising the chunk size cut the reload of a 16k-token document to about 1.9 s and let every workload cell finish. With the default setup the largest cells could not complete.
The story
Every token an LLM generates depends on key and value vectors for all tokens before it. vLLM keeps these in a KV cache so it does not recompute them, and that cache grows with context length and concurrency until GPU memory is full. From there it spills to CPU DRAM and then to NVMe, and the storage path starts deciding how fast the first token comes back. Published systems show that offloading works. Few say what the storage underneath has to deliver.
I built the experiment platform on vLLM 0.19 with LMCache as the tiering layer, containerized with Apptainer on the Mogon cluster and Docker on a GPU workstation, and drove it with ShareGPT and BurstGPT traffic. Fourteen experiments moved from a single GPU up to production-style traffic on a 70B model with tensor parallelism. I isolated the storage variable with a FUSE throttling shim that caps bandwidth and injects per-operation latency independently, and cross-checked the ordering against a kernel-level dm-delay device.


What it means
The results locate the cost precisely. Writing the cache out is close to free in steady state and decode never touches storage, so the whole penalty lands on prefill. Within prefill, per-operation latency binds and bandwidth barely matters. That result reads as a requirements specification for any tier placed under the cache: a fraction of one NVMe device's bandwidth is enough, but every millisecond of round-trip time is paid in full.
Networked storage is exactly where latency is hard to keep low, so I integrated GekkoFS, an ad-hoc file system built at job start from the allocation's own node-local NVMe. I wrote SharedFSBackend, which plugs it into LMCache as a shared disk tier. On a single instance it lost to local NVMe, reading at about 0.7 GB/s against 6 GB/s with nothing to hide the round trip behind. Its value appeared once several instances served behind a load balancer. A private cache means each replica recomputes work its neighbors already did, whereas one shared namespace stores each prefix once and lets any replica read it.

Design rule
For an inference platform the design rule is concrete. Keep the hot tiers node-local and judge any shared tier by its per-operation latency before its throughput. A shared tier earns its place once the fleet runs more than one replica and users share prefixes. The 8B model used for most runs is the hardest case for a reload to win, so the parity measured here is a floor. At 70B, where recomputation costs far more per token, the same design should turn parity into a latency win.

Technology

Working on something similar?
Tell me what you run today and where it hurts. I will come back with how I would build it, and what I would leave alone.
Discuss your project

