PidokuInfra

Distributed KV Cache

Advanced 1h Difficulty 4/5 Topic 11 of 12

Prerequisites V.05, 03, 08


1. What is it?#

Where the KV cache lives when a model instance spans multiple GPUs — and, increasingly, when you deliberately move KV between GPUs, nodes, or storage tiers.

Three distinct topics:
  1. How KV is SPLIT by the parallelism strategy (automatic; understand it)
  2. How KV is SHARED across replicas for prefix caching (an optimization)
  3. How KV is MOVED between machines (disaggregation, Section XIII.06)

2. How each parallelism strategy splits the KV cache#

TENSOR PARALLELISM (TP=N)
  KV heads are split across GPUs, along with the attention heads.
  Each GPU stores KV for n_kv_heads/N heads, for ALL sequences.
  
  Llama-3-70B, n_kv_heads=8, TP=8:
      each GPU holds 1 KV head's worth: 320 KiB / 8 = 40 KiB per token
  
  → total KV capacity scales with N. Good.
  → constraint: TP must divide n_kv_heads, or KV heads are REPLICATED
    (e.g. TP=16 with 8 KV heads → each head on 2 GPUs → 2x KV memory waste)

PIPELINE PARALLELISM (PP=P)
  Each stage holds KV for ITS layers only, for all sequences.
  
  Llama-3-70B, 80 layers, PP=2:
      stage 0 holds KV for layers 0-39
      stage 1 holds KV for layers 40-79
  
  → total capacity scales with P. Good.
  → but the KV for one sequence is spread across stages, so a sequence
    cannot migrate between stages.

DATA PARALLELISM (DP=N)
  Each replica has its OWN complete KV cache for its OWN sequences.
  
  → no sharing. A sequence lives entirely on one replica.
  → total capacity scales with N, but per-sequence capacity does not.
  → each replica has its own prefix cache — hence prefix-aware routing
    (Section VIII.06).

EXPERT PARALLELISM
  KV is in the ATTENTION layers, which are typically TP-parallel, not
  EP-parallel. So EP doesn't affect KV placement.

SEQUENCE PARALLELISM
  KV is split by TOKEN POSITION across GPUs.
  → this is the only strategy that lets ONE sequence's KV exceed one GPU.

The important asymmetry: TP and PP increase the capacity available to a single sequence; DP does not. If one request needs 100 GB of KV, DP doesn’t help at all.


3. Simple analogy#

A patient’s medical file across a hospital.

  • TP: the file is split by specialty — cardiology holds the heart notes, neurology the brain notes. Every department sees every patient, but only their own sections.
  • PP: the file is split by time — the first ward holds the admission notes, the second the treatment notes.
  • DP: each hospital has complete files for its own patients. No sharing. A patient must go to the hospital that has their file.
  • SP: one patient’s file is so large it’s split across buildings by chapter.

4. Prefix cache sharing across replicas#

With DP, each replica has its own prefix cache. For a workload with heavy sharing (a common system prompt, a shared document), this means N copies of the same KV.

8 replicas, 2,000-token shared system prompt, Llama-3-8B:
  KV for the prompt: 2000 × 128 KiB = 256 MB
  × 8 replicas = 2 GB of duplicated KV

Not catastrophic here. But for a 100k-token shared document:
  12.8 GB × 8 = 102 GB duplicated.

Three approaches:

1. PREFIX-AWARE ROUTING (Section VIII.06)
   Route requests with the same prefix to the same replica.
   ✓ simple, no new infrastructure
   ✓ each prefix is cached once (on its assigned replica)
   ✗ load imbalance if one prefix is very popular
   → THE STANDARD ANSWER

2. SHARED KV STORE  [EMERGING]
   A external tier (CPU memory, NVMe, or a dedicated KV service) holding
   computed prefix KV. Replicas fetch on demand.
   ✓ one copy globally; any replica can serve any prefix
   ✗ fetching costs bandwidth: 256 MB over 50 GB/s IB = 5 ms
     (vs recomputing the prefill: often comparable or cheaper!)
   → viable when the prefix is LONG (recompute is expensive) and the
     interconnect is FAST

3. RECOMPUTE
   Just prefill it again on whichever replica gets the request.
   ✓ zero infrastructure
   ✗ pays full prefill cost
   → the default when routing can't help

The decision hinges on: is fetching the KV cheaper than recomputing it?

Recompute cost:  prefill_FLOPs / achieved_FLOPs
Fetch cost:      KV_bytes / interconnect_bandwidth

Llama-3-8B, 10,000-token prefix:
  recompute: 2 × 8e9 × 10000 = 160 TFLOP / 600 TFLOP/s = 267 ms
  fetch:     10000 × 128 KiB = 1.28 GB / 50 GB/s = 26 ms      ← 10x cheaper

Llama-3-8B, 500-token prefix:
  recompute: 8 TFLOP / 600 = 13 ms
  fetch:     64 MB / 50 GB/s = 1.3 ms                          ← still cheaper
             but the fixed overhead (RPC, lookup) may dominate

Fetching KV is generally cheaper than recomputing it, by roughly the ratio of arithmetic_intensity_of_prefill to interconnect_bandwidth/FLOPs. This is the fundamental insight behind disaggregated serving (Section XIII.06) and shared KV stores.

The catch: it requires a fast interconnect and careful engineering. And note that this same arithmetic is why preemption by recomputation beats swapping over PCIe (Section V.09) — there the interconnect is much slower, flipping the answer.


5. KV cache tiering#

TIER            CAPACITY    BANDWIDTH      LATENCY     USE
GPU HBM         10-100 GB   3,350 GB/s     ~0.5 µs     active sequences
CPU DRAM        0.5-2 TB    ~50 GB/s (PCIe) ~10 µs     recently-used prefixes
                            900 GB/s (C2C)
Local NVMe      1-100 TB    3-14 GB/s      ~100 µs     long-tail prefixes
Remote store    ∞           1-50 GB/s      ~1 ms       shared across the fleet

A tiered KV cache promotes/demotes blocks based on access:

Hot prefix (used every second)      → GPU
Warm prefix (used every minute)     → CPU DRAM
Cold prefix (used every hour)       → NVMe or remote

[EMERGING] — systems like LMCache and vLLM’s experimental connectors implement this. The engineering is nontrivial (async prefetch, eviction policy, consistency) and the benefit depends heavily on your prefix reuse distribution.

Grace-Hopper changes the calculus: NVLink-C2C gives 900 GB/s CPU↔GPU, making the CPU tier genuinely fast rather than a fallback.


6. Under the hood — what moving KV actually requires#

To transfer a sequence's KV from GPU A to GPU B:

1. Identify the blocks (from A's block table)
2. Ensure B has free blocks
3. Copy: A's HBM → B's HBM
     same node, NVLink:  P2P copy, ~370 GB/s
     across nodes:       GPUDirect RDMA, ~50 GB/s
4. Build B's block table pointing at the new blocks
5. Free A's blocks

Complications:
  - the layout must match (same TP degree, same head split)
  - RoPE positions are baked in, so the sequence must resume at the same
    position (Section V.05)
  - the transfer must not block ongoing compute → separate stream
  - a partially-transferred sequence is in an inconsistent state

Point 1 in the complications list is a hard constraint: KV computed under TP=8 has a different head layout than KV computed under TP=4. Systems that move KV between differently- configured workers must reshape it, which costs time and complexity. Disaggregated serving designs usually require matching TP degrees on both sides for exactly this reason.


7. Performance#

Operation                                    Time (Llama-3-8B, 10k tokens)
Recompute prefill                            267 ms
Copy KV GPU→GPU same node (NVLink)           1.28 GB / 370 GB/s = 3.5 ms
Copy KV GPU→GPU across node (IB NDR)         1.28 GB / 50 GB/s = 26 ms
Copy KV GPU→CPU (PCIe Gen5)                  1.28 GB / 55 GB/s = 23 ms
Copy KV GPU→CPU (NVLink-C2C)                 1.28 GB / 900 GB/s = 1.4 ms
Load KV from local NVMe                      1.28 GB / 7 GB/s = 183 ms
Load KV from remote object store             1.28 GB / 2 GB/s = 640 ms

Read the last two rows against the first. Loading KV from NVMe (183 ms) is slower than recomputing it (267 ms)? No — it’s faster, but not by much, and the object store is far slower. The tier only helps if it’s faster than recompute, which rules out slow storage for short-to-medium prefixes.


8. Production implications#

  • Understand how your parallelism splits KV. It determines per-GPU capacity, which determines concurrency.
  • TP degree must divide n_kv_heads, or you waste memory replicating KV heads.
  • Prefix-aware routing is the standard answer for cross-replica sharing. Do that before considering a KV store.
  • Compute fetch-vs-recompute for your prefix lengths. For short prefixes, recompute wins on simplicity.
  • KV tiering is [EMERGING]. Evaluate it if you have very long shared prefixes and a fast interconnect; otherwise it’s complexity without payoff.
  • KV layout compatibility constrains architecture. Workers exchanging KV must have matching TP degrees and layouts.

9. Common mistakes#

TP degree not dividing n_kv_heads. Silent memory waste.

Assuming DP increases per-sequence KV capacity. It doesn’t.

Building a KV store before implementing prefix-aware routing. Much more complexity for overlapping benefit.

Ignoring the RoPE position constraint when moving or reusing KV.

Using slow storage for a KV tier. If it’s slower than recompute, it’s worse than nothing.

Mismatched TP degrees between KV producers and consumers.


10. Hands-on exercise#

A. Compute the splits. For a model at TP=8 and PP=2, compute per-GPU KV bytes per token. Verify against what the engine reports.

B. Fetch vs recompute. For your model and interconnect, compute the crossover prefix length at which fetching KV beats recomputing it. Plot both curves.

C. The TP=16 penalty. For a model with 8 KV heads, compute the memory waste at TP=16 (KV head replication). Is TP=16 worth it?

D. Measure a KV transfer. Implement a GPU→GPU KV block copy over NVLink and over host memory. Measure both. Compare to the recompute time.

E. Duplication cost. For your workload and replica count, compute how much KV is duplicated across replicas due to prefix caching. Would routing fix it?


11. Interview questions#

  1. How does each parallelism strategy split the KV cache?
  2. Why does TP increase per-sequence KV capacity but DP doesn’t?
  3. What happens when the TP degree doesn’t divide the number of KV heads?
  4. When is fetching KV cheaper than recomputing it? Derive the crossover.
  5. Why does RoPE constrain KV cache movement?
  6. What must match between two workers that exchange KV cache?
  7. When would a tiered KV cache be worth the complexity?

12. Further reading#

  • [ESTABLISHED] Pope et al., “Efficiently Scaling Transformer Inference” (2022)
  • [EMERGING] Zhong et al., “DistServe” (2024); Qin et al., “Mooncake” (2024) — KV-centric disaggregated architectures
  • [EMERGING] LMCache and vLLM KV connector documentation
  • Next: 12 — When more GPUs make things worse

↑↓ navigate↵ openesc close