Skip to content
4 changes: 4 additions & 0 deletions docs/features/kv_offloading_usage.md
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,10 @@ flowchart LR
CPU <--> SN["..."]
```

## Terminology: Chunks

The unit of operation is a **chunk** — a fixed-size piece of KV data covering a group of tokens. By default, a chunk maps to a single accelerator block. A configurable `blocks_per_chunk` parameter allows larger chunks, yielding larger I/Os to the host and secondary tiers.

## Single-Tier Setup (CPU Only)

```bash
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -704,7 +704,7 @@ def test_prom_metrics_registers_tiering_metrics_from_spec():
)

metric = prom_metrics._offloading_metric_defs[
TieringOffloadingMetrics.BLOCK_QUERIES
TieringOffloadingMetrics.CHUNK_QUERIES
]
assert metric.kwargs["labelnames"] == ["model_name", "engine", "tier"]

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -96,7 +96,7 @@ def test_swa_offload_window_covers_unaligned_hit(boundary, eagle, left_state):
groups.append(
KVCacheGroupSpec([f"layer{i}"], kv_spec, is_eagle_group=eagle and i == 1)
)
manager = CPUOffloadingManager(num_blocks=100)
manager = CPUOffloadingManager(num_chunks=100)
spec = SimpleNamespace(
tokens_per_block=(256, 64, 8),
tokens_per_hash=8,
Expand Down Expand Up @@ -2888,7 +2888,7 @@ def _make_req_status(
req=req,
req_context=ReqContext(req_id="test-req"),
offloading_context=RequestOffloadingContext(
policy=OffloadPolicy.BLOCK_LEVEL
policy=OffloadPolicy.CHUNK_LEVEL
),
num_locally_computed_tokens=num_computed_tokens,
)
Expand Down
12 changes: 6 additions & 6 deletions tests/v1/kv_offload/cpu/policies/test_factory.py
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,7 @@
from vllm.v1.kv_offload.base import OffloadKey, ReqContext
from vllm.v1.kv_offload.cpu.manager import CPUOffloadingManager
from vllm.v1.kv_offload.cpu.policies.arc import ARCCachePolicy
from vllm.v1.kv_offload.cpu.policies.base import BlockStatus, CachePolicy
from vllm.v1.kv_offload.cpu.policies.base import CachePolicy, ChunkStatus
from vllm.v1.kv_offload.cpu.policies.factory import CachePolicyFactory
from vllm.v1.kv_offload.cpu.policies.lru import LRUCachePolicy

Expand All @@ -20,10 +20,10 @@ class _DummyCachePolicy(CachePolicy):
def __init__(self, cache_capacity: int) -> None:
self.cache_capacity = cache_capacity

def get(self, key: OffloadKey) -> BlockStatus | None:
def get(self, key: OffloadKey) -> ChunkStatus | None:
return None

def insert(self, key: OffloadKey, block: BlockStatus) -> None:
def insert(self, key: OffloadKey, chunk: ChunkStatus) -> None:
pass

def remove(self, key: OffloadKey) -> None:
Expand All @@ -34,7 +34,7 @@ def touch(self, keys: Iterable[OffloadKey], req_context: ReqContext) -> None:

def evict(
self, n: int, protected: set[OffloadKey]
) -> list[tuple[OffloadKey, BlockStatus]] | None:
) -> list[tuple[OffloadKey, ChunkStatus]] | None:
return None

def clear(self) -> None:
Expand Down Expand Up @@ -72,7 +72,7 @@ def test_register_and_resolve_custom_policy(self):
policy_cls = CachePolicyFactory.get_cache_policy_cls("dummy")
assert policy_cls is _DummyCachePolicy

manager = CPUOffloadingManager(num_blocks=4, cache_policy="dummy")
manager = CPUOffloadingManager(num_chunks=4, cache_policy="dummy")
assert isinstance(manager._policy, _DummyCachePolicy)

def test_unregistered_policy_raises(self):
Expand All @@ -98,7 +98,7 @@ def test_manager_resolves_policy_via_module_path(self):
"""End-to-end: CPUOffloadingManager resolves an unregistered policy
purely from cache_policy_module_path."""
manager = CPUOffloadingManager(
num_blocks=4,
num_chunks=4,
cache_policy="_DummyCachePolicy",
cache_policy_module_path="tests.v1.kv_offload.cpu.policies.test_factory",
)
Expand Down
4 changes: 2 additions & 2 deletions tests/v1/kv_offload/cpu/test_canonical_layout.py
Original file line number Diff line number Diff line change
Expand Up @@ -232,9 +232,9 @@ def test_cross_topology_roundtrip(writer_tp: int, reader_tp: int):
def canonical_view(rank: int, world_size: int) -> torch.Tensor:
region = SharedOffloadRegion(
engine_id=engine_id,
num_blocks=num_blocks,
num_chunks=num_blocks,
rank=rank,
kv_bytes_per_block=row_stride,
kv_bytes_per_chunk=row_stride,
cpu_page_size=row_stride // world_size,
)
regions.append(region)
Expand Down
92 changes: 46 additions & 46 deletions tests/v1/kv_offload/cpu/test_gpu_worker.py
Original file line number Diff line number Diff line change
Expand Up @@ -26,7 +26,7 @@
from vllm.v1.kv_offload.cpu.shared_offload_region import SharedOffloadRegion

NUM_GPU_BLOCKS = [64]
NUM_CPU_BLOCKS = [256]
NUM_CPU_CHUNKS = [256]
GPU_PAGE_SIZES = [512, 1024]
BLOCKS_PER_CHUNK_VALUES = [1, 3]
NUM_TENSORS = [4]
Expand Down Expand Up @@ -225,7 +225,7 @@ def record_region_cleanup() -> None:
@pytest.mark.parametrize("gpu_page_size_bytes", GPU_PAGE_SIZES)
@pytest.mark.parametrize("blocks_per_chunk", BLOCKS_PER_CHUNK_VALUES)
@pytest.mark.parametrize("num_gpu_blocks", NUM_GPU_BLOCKS)
@pytest.mark.parametrize("num_cpu_blocks", NUM_CPU_BLOCKS)
@pytest.mark.parametrize("num_cpu_chunks", NUM_CPU_CHUNKS)
@pytest.mark.parametrize("num_tensors", NUM_TENSORS)
@pytest.mark.parametrize("seed", SEEDS)
@pytest.mark.parametrize("device", DEVICES)
Expand All @@ -241,7 +241,7 @@ def test_transfer(
gpu_page_size_bytes: int,
blocks_per_chunk: int,
num_gpu_blocks: int,
num_cpu_blocks: int,
num_cpu_chunks: int,
num_tensors: int,
seed: int,
device: str,
Expand Down Expand Up @@ -288,58 +288,58 @@ def test_transfer(
SharedOffloadRegion.BLOCK_SIZE_ALIGNMENT,
)
simulated_world_size = 2
kv_bytes_per_block = (
kv_bytes_per_chunk = (
cpu_page_size if replicated_layout else cpu_page_size * simulated_world_size
)
mmap_region = SharedOffloadRegion(
engine_id=str(uuid.uuid4()),
num_blocks=num_cpu_blocks,
num_chunks=num_cpu_chunks,
rank=0,
kv_bytes_per_block=kv_bytes_per_block,
kv_bytes_per_chunk=kv_bytes_per_chunk,
cpu_page_size=cpu_page_size,
)

worker = CPUOffloadingWorker(
kv_caches=kv_caches,
blocks_per_chunk=blocks_per_chunk,
num_cpu_blocks=num_cpu_blocks,
num_cpu_chunks=num_cpu_chunks,
mmap_region=mmap_region,
)

# select block mappings
gpu_blocks = random.sample(range(num_gpu_blocks), num_mappings * blocks_per_chunk)
cpu_blocks = random.sample(range(num_cpu_blocks), num_mappings)
cpu_chunks = random.sample(range(num_cpu_chunks), num_mappings)

# expand cpu blocks to gpu-page granularity for uniform comparison:
# each cpu block maps to blocks_per_chunk consecutive sub-blocks
cpu_blocks_expanded = [
cpu_block * blocks_per_chunk + j
for cpu_block in cpu_blocks
# expand cpu chunks to gpu-page granularity for uniform comparison:
# each cpu chunk maps to blocks_per_chunk consecutive sub-blocks
cpu_chunks_expanded = [
cpu_chunk * blocks_per_chunk + j
for cpu_chunk in cpu_chunks
for j in range(blocks_per_chunk)
]

# maybe skip some GPU blocks to test reading/writing from the middle of a CPU block
# maybe skip some GPU blocks to test reading/writing from the middle of a CPU chunk
blocks_to_skip = blocks_per_chunk - 1
if blocks_to_skip > 0:
gpu_blocks = gpu_blocks[blocks_to_skip:]
cpu_blocks_expanded = cpu_blocks_expanded[blocks_to_skip:]
cpu_chunks_expanded = cpu_chunks_expanded[blocks_to_skip:]

# set transfer direction
if gpu_to_cpu:
handler = worker._store_handler
src_spec = GPULoadStoreSpec(
gpu_blocks, group_sizes=(len(gpu_blocks),), block_indices=(blocks_to_skip,)
)
dst_spec = CPULoadStoreSpec(cpu_blocks)
dst_to_src = dict(zip(cpu_blocks_expanded, gpu_blocks))
dst_spec = CPULoadStoreSpec(cpu_chunks)
dst_to_src = dict(zip(cpu_chunks_expanded, gpu_blocks))
num_dst_sub_blocks = num_gpu_blocks
else:
handler = worker._load_handler
src_spec = CPULoadStoreSpec(cpu_blocks)
src_spec = CPULoadStoreSpec(cpu_chunks)
dst_spec = GPULoadStoreSpec(
gpu_blocks, group_sizes=(len(gpu_blocks),), block_indices=(blocks_to_skip,)
)
dst_to_src = dict(zip(gpu_blocks, cpu_blocks_expanded))
dst_to_src = dict(zip(gpu_blocks, cpu_chunks_expanded))
num_dst_sub_blocks = num_gpu_blocks

# randomize src and dst tensors before transfer
Expand Down Expand Up @@ -410,7 +410,7 @@ def test_transfer(
@pytest.mark.parametrize("gpu_page_size_bytes", GPU_PAGE_SIZES)
@pytest.mark.parametrize("blocks_per_chunk", BLOCKS_PER_CHUNK_VALUES)
@pytest.mark.parametrize("num_gpu_blocks", NUM_GPU_BLOCKS)
@pytest.mark.parametrize("num_cpu_blocks", NUM_CPU_BLOCKS)
@pytest.mark.parametrize("num_cpu_chunks", NUM_CPU_CHUNKS)
@pytest.mark.parametrize("seed", SEEDS)
@pytest.mark.parametrize("device", DEVICES)
@torch.inference_mode()
Expand All @@ -421,7 +421,7 @@ def test_transfer_multi_group(
gpu_page_size_bytes: int,
blocks_per_chunk: int,
num_gpu_blocks: int,
num_cpu_blocks: int,
num_cpu_chunks: int,
seed: int,
device: str,
) -> None:
Expand Down Expand Up @@ -470,37 +470,37 @@ def test_transfer_multi_group(
worker = CPUOffloadingWorker(
kv_caches=canonical_kv_caches,
blocks_per_chunk=blocks_per_chunk,
num_cpu_blocks=num_cpu_blocks,
num_cpu_chunks=num_cpu_chunks,
)

# group 0: aligned, group 1: empty, group 2: unaligned on CPU->GPU
group_sizes_in_cpu_blocks = [num_mappings_per_group, 0, num_mappings_per_group]
group_sizes_in_cpu_chunks = [num_mappings_per_group, 0, num_mappings_per_group]

total_cpu_blocks = sum(group_sizes_in_cpu_blocks)
total_gpu_blocks_needed = total_cpu_blocks * blocks_per_chunk
total_cpu_chunks = sum(group_sizes_in_cpu_chunks)
total_gpu_blocks_needed = total_cpu_chunks * blocks_per_chunk
gpu_blocks_all = random.sample(range(num_gpu_blocks), total_gpu_blocks_needed)
cpu_blocks_all = random.sample(range(num_cpu_blocks), total_cpu_blocks)
cpu_chunks_all = random.sample(range(num_cpu_chunks), total_cpu_chunks)

# split gpu/cpu blocks per group
# split gpu blocks / cpu chunks per group
gpu_blocks_per_group: list[list[int]] = []
cpu_blocks_per_group: list[list[int]] = []
cpu_chunks_per_group: list[list[int]] = []
gpu_offset = 0
cpu_offset = 0
for size in group_sizes_in_cpu_blocks:
for size in group_sizes_in_cpu_chunks:
gpu_count = size * blocks_per_chunk
gpu_blocks_per_group.append(gpu_blocks_all[gpu_offset : gpu_offset + gpu_count])
cpu_blocks_per_group.append(cpu_blocks_all[cpu_offset : cpu_offset + size])
cpu_chunks_per_group.append(cpu_chunks_all[cpu_offset : cpu_offset + size])
gpu_offset += gpu_count
cpu_offset += size

# expand cpu blocks to gpu-page granularity
cpu_blocks_expanded_per_group = [
# expand cpu chunks to gpu-page granularity
cpu_chunks_expanded_per_group = [
[
cpu_block * blocks_per_chunk + j
for cpu_block in cpu_blocks
cpu_chunk * blocks_per_chunk + j
for cpu_chunk in cpu_chunks
for j in range(blocks_per_chunk)
]
for cpu_blocks in cpu_blocks_per_group
for cpu_chunks in cpu_chunks_per_group
]

# skip sub-blocks from group 2 to test unaligned transfers.
Expand All @@ -509,7 +509,7 @@ def test_transfer_multi_group(
gpu_blocks_per_group[2] = gpu_blocks_per_group[2][
sub_blocks_to_skip:-sub_blocks_to_skip
]
cpu_blocks_expanded_per_group[2] = cpu_blocks_expanded_per_group[2][
cpu_chunks_expanded_per_group[2] = cpu_chunks_expanded_per_group[2][
sub_blocks_to_skip:-sub_blocks_to_skip
]

Expand All @@ -520,10 +520,10 @@ def test_transfer_multi_group(
gpu_blocks.extend(gpu_blks)
group_sizes.append(len(gpu_blks))

# build flat cpu_blocks list
cpu_blocks = []
for cpu_blks in cpu_blocks_per_group:
cpu_blocks.extend(cpu_blks)
# build flat cpu_chunks list
cpu_chunks = []
for cpu_chnks in cpu_chunks_per_group:
cpu_chunks.extend(cpu_chnks)

# block_indices: only relevant for unaligned transfers
block_indices: list[int] = [0, 0, sub_blocks_to_skip]
Expand All @@ -533,26 +533,26 @@ def test_transfer_multi_group(
src_spec = GPULoadStoreSpec(
gpu_blocks, group_sizes=group_sizes, block_indices=block_indices
)
dst_spec = CPULoadStoreSpec(cpu_blocks)
dst_spec = CPULoadStoreSpec(cpu_chunks)
# per-group mapping: cpu sub-block -> gpu sub-block
dst_to_src_per_group = [
dict(zip(expanded, gpu_blks))
for expanded, gpu_blks in zip(
cpu_blocks_expanded_per_group, gpu_blocks_per_group
cpu_chunks_expanded_per_group, gpu_blocks_per_group
)
]
num_dst_sub_blocks = num_cpu_blocks * blocks_per_chunk
num_dst_sub_blocks = num_cpu_chunks * blocks_per_chunk
else:
handler = worker._load_handler
src_spec = CPULoadStoreSpec(cpu_blocks)
src_spec = CPULoadStoreSpec(cpu_chunks)
dst_spec = GPULoadStoreSpec(
gpu_blocks, group_sizes=group_sizes, block_indices=block_indices
)
# per-group mapping: gpu sub-block -> cpu sub-block
dst_to_src_per_group = [
dict(zip(gpu_blks, expanded))
for gpu_blks, expanded in zip(
gpu_blocks_per_group, cpu_blocks_expanded_per_group
gpu_blocks_per_group, cpu_chunks_expanded_per_group
)
]
num_dst_sub_blocks = num_gpu_blocks
Expand Down Expand Up @@ -644,7 +644,7 @@ def test_load_waits_for_pending_compute_stream_writes(default_vllm_config) -> No
],
),
blocks_per_chunk=1,
num_cpu_blocks=num_blocks,
num_cpu_chunks=num_blocks,
)
worker._load_handler.src_tensors[0].fill_(sentinel)
expected = torch.full((page_size_bytes,), sentinel, dtype=torch.int8)
Expand Down
Loading
Loading