Skip to content

vllm.v1.kv_offload.tiering.spec

TieringOffloadingSpec: Spec for multi-tier KV cache offloading.

This spec creates a TieringOffloadingManager with a CPU primary tier and configurable secondary tiers (e.g., Storage, Network).

Configuration via kv_connector_extra_config
  • cpu_bytes_to_use: (required) Bytes to allocate for CPU primary tier
  • block_size: (optional) Tokens per offloaded chunk (default: GPU block size)
  • eviction_policy: (optional) Primary tier eviction policy: built-in "lru"/ "arc", or the name of a policy registered via CachePolicyFactory, or an out-of-tree CachePolicy class name paired with cache_policy_module_path (default: "lru")
  • cache_policy_module_path: (optional) Python import path to load eviction_policy from when it names an out-of-tree CachePolicy not registered via CachePolicyFactory
  • secondary_tiers: (optional) List of secondary tier configurations Each secondary tier config is a dict with:
    • type: (required) Type of secondary tier (e.g., "example", "fs", "p2p", "obj"), or the class name of an out-of-tree SecondaryTierManager paired with module_path.
    • module_path: (optional) Python import path to load 'type' from when it names an out-of-tree SecondaryTierManager not registered via SecondaryTierFactory.register_tier()
    • Additional tier-specific parameters are passed directly to the tier constructor. See each tier's documentation for supported parameters.

Example configuration: { "cpu_bytes_to_use": 10737418240, # 10 GB "block_size": 16, "eviction_policy": "lru", "secondary_tiers": [ { "type": "example", "custom_param": 67 } ] }

Example out-of-tree tier configuration: { "cpu_bytes_to_use": 10737418240, "secondary_tiers": [ { "type": "MyCustomTier", "module_path": "my_package.my_module", "custom_param": "value" } ] }

Classes:

TieringOffloadingSpec

Bases: CPUOffloadingSpec

Spec for multi-tier KV cache offloading.

Creates a TieringOffloadingManager with: - Primary tier: CPU (LRU or ARC eviction policy) - Secondary tiers: Configurable via extra_config

The CPU primary tier has direct GPU access and serves as the gateway for all GPU↔offload operations. Secondary tiers cannot directly access GPU memory and must transfer data through the primary tier.

Methods:

Source code in vllm/v1/kv_offload/tiering/spec.py
 83
 84
 85
 86
 87
 88
 89
 90
 91
 92
 93
 94
 95
 96
 97
 98
 99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
class TieringOffloadingSpec(CPUOffloadingSpec):
    """Spec for multi-tier KV cache offloading.

    Creates a TieringOffloadingManager with:
    - Primary tier: CPU (LRU or ARC eviction policy)
    - Secondary tiers: Configurable via extra_config

    The CPU primary tier has direct GPU access and serves as the gateway for
    all GPU↔offload operations. Secondary tiers cannot directly access GPU
    memory and must transfer data through the primary tier.
    """

    BLOCK_SIZE_ALIGNMENT = SharedOffloadRegion.BLOCK_SIZE_ALIGNMENT

    @classmethod
    @override
    def build_metric_definitions(
        cls, extra_config: dict[str, Any]
    ) -> dict[str, OffloadingMetricMetadata]:
        metrics = super().build_metric_definitions(extra_config)
        metrics[TieringOffloadingMetrics.LOOKUP_SYNC_DELAY] = (
            OffloadingHistogramMetadata(
                documentation=(
                    "Histogram of blocking time spent in a per-chunk tier lookup "
                    "that resolved as a hit or miss, labeled by tier, in seconds."
                ),
                labelnames=("tier",),
                buckets=(
                    0.00001,
                    0.00005,
                    0.0001,
                    0.0005,
                    0.001,
                    0.005,
                    0.01,
                    0.05,
                    0.1,
                    0.5,
                    1,
                ),
            )
        )
        metrics[TieringOffloadingMetrics.LOOKUP_ASYNC_DELAY] = (
            OffloadingHistogramMetadata(
                documentation=(
                    "Histogram of wall-clock time from a per-chunk tier lookup "
                    "first returning retry until that same tier lookup resolves "
                    "as a hit or miss, labeled by tier, in seconds."
                ),
                labelnames=("tier",),
                buckets=(
                    0.0001,
                    0.0005,
                    0.001,
                    0.005,
                    0.01,
                    0.05,
                    0.1,
                    0.5,
                    1,
                    5,
                    10,
                ),
            )
        )
        metrics[TieringOffloadingMetrics.READ_BYTES] = OffloadingCounterMetadata(
            documentation=(
                "Total bytes read from secondary tiers into the primary tier, "
                "labeled by tier."
            ),
            labelnames=("tier",),
        )
        metrics[TieringOffloadingMetrics.READ_TIME] = OffloadingCounterMetadata(
            documentation=(
                "Total time spent reading from secondary tiers into the primary "
                "tier, in seconds, labeled by tier."
            ),
            labelnames=("tier",),
        )
        metrics[TieringOffloadingMetrics.WRITE_BYTES] = OffloadingCounterMetadata(
            documentation=(
                "Total bytes written from the primary tier to secondary tiers, "
                "labeled by tier."
            ),
            labelnames=("tier",),
        )
        metrics[TieringOffloadingMetrics.WRITE_TIME] = OffloadingCounterMetadata(
            documentation=(
                "Total time spent writing from the primary tier to secondary "
                "tiers, in seconds, labeled by tier."
            ),
            labelnames=("tier",),
        )
        metrics[TieringOffloadingMetrics.PROMOTION_JOB_FAILURES] = (
            OffloadingCounterMetadata(
                documentation=(
                    "Number of failed secondary-tier promotion jobs, labeled by tier."
                ),
                labelnames=("tier",),
            )
        )
        metrics[TieringOffloadingMetrics.CASCADE_JOB_FAILURES] = (
            OffloadingCounterMetadata(
                documentation=(
                    "Number of failed secondary-tier cascade jobs, labeled by tier."
                ),
                labelnames=("tier",),
            )
        )
        metrics[TieringOffloadingMetrics.CHUNK_QUERIES] = OffloadingCounterMetadata(
            documentation=(
                "Number of chunk lookup queries sent to a tier, labeled by tier."
            ),
            labelnames=("tier",),
        )
        metrics[TieringOffloadingMetrics.CHUNK_HITS] = OffloadingCounterMetadata(
            documentation="Number of chunk lookup hits in a tier, labeled by tier.",
            labelnames=("tier",),
        )
        metrics[TieringOffloadingMetrics.PROMOTION_ALLOCATION_FAILURES] = (
            OffloadingCounterMetadata(
                documentation=(
                    "Number of promotion attempts that failed because the "
                    "primary tier could not allocate space."
                ),
            )
        )
        metrics[TieringOffloadingMetrics.PRIMARY_WRITE_USAGE_PERC] = (
            OffloadingGaugeMetadata(
                documentation=(
                    "Current fraction of primary-tier space used by writes from "
                    "secondary tiers, labeled by tier."
                ),
                labelnames=("tier",),
            )
        )
        metrics[TieringOffloadingMetrics.PRIMARY_READ_USAGE_PERC] = (
            OffloadingGaugeMetadata(
                documentation=(
                    "Current fraction of primary-tier space used by reads to "
                    "secondary tiers, labeled by tier."
                ),
                labelnames=("tier",),
            )
        )
        metrics[TieringOffloadingMetrics.ACTIVE_PROMOTION_JOBS] = (
            OffloadingGaugeMetadata(
                documentation=(
                    "Number of active secondary-tier promotion jobs, labeled by tier."
                ),
                labelnames=("tier",),
            )
        )
        metrics[TieringOffloadingMetrics.ACTIVE_CASCADE_JOBS] = OffloadingGaugeMetadata(
            documentation=(
                "Number of active secondary-tier cascade jobs, labeled by tier."
            ),
            labelnames=("tier",),
        )
        secondary_tier_configs = extra_config.get("secondary_tiers", [])
        if not isinstance(secondary_tier_configs, list):
            raise ValueError("secondary_tiers must be a list of tier configurations")

        for tier_config in secondary_tier_configs:
            assert isinstance(tier_config, dict)
            tier_cls = SecondaryTierFactory.get_tier_class(tier_config)
            metrics.update(tier_cls.build_metric_definitions(tier_config))

        metrics[TieringOffloadingMetrics.BACKPRESSURE_STORE_LATENCY_EMA] = (
            OffloadingGaugeMetadata(
                documentation=(
                    "Exponential moving average of store latency "
                    "for back-pressure detection, in s/MiB."
                ),
                labelnames=("tier",),
            )
        )
        metrics[TieringOffloadingMetrics.BACKPRESSURE_STORES_DROPPED] = (
            OffloadingCounterMetadata(
                documentation=(
                    "Number of store operations dropped due to "
                    "back-pressure on a secondary tier."
                ),
                labelnames=("tier",),
            )
        )
        metrics[TieringOffloadingMetrics.BACKPRESSURE_BLOCKS_DROPPED] = (
            OffloadingCounterMetadata(
                documentation=(
                    "Number of blocks dropped due to back-pressure on a secondary tier."
                ),
                labelnames=("tier",),
            )
        )

        return metrics

    def __init__(self, config: OffloadingConfig):
        super().__init__(config)
        # Redeclare for mypy: parent sets this but `--follow-imports skip` hides it
        self._manager: OffloadingManager | None = None

        # Parse secondary tier configurations
        self.secondary_tier_configs = self.extra_config.get("secondary_tiers", [])
        if not isinstance(self.secondary_tier_configs, list):
            raise ValueError("secondary_tiers must be a list of tier configurations")

        # Backpressure config is merged field-by-field in priority order
        # (highest first):
        #   1. Per-tier ``backpressure`` dict in the tier config
        #   2. Top-level ``backpressure`` in kv_connector_extra_config
        # Merging per field (rather than per whole dict) means a partial
        # tier override still inherits missing fields from the top-level
        # default, so the resolved dict reaching the factory is complete.
        # Within each tier's resolved dict, tier-type-aware water marks
        # are filled in last so a bare ``"backpressure": {}`` picks up
        # sensible thresholds for the storage medium.
        bp_defaults = self.extra_config.get("backpressure")

        for tier_cfg in self.secondary_tier_configs:
            tier_override = tier_cfg.get("backpressure")
            # Overlay from lowest to highest precedence so higher-precedence
            # fields win while lower-precedence ones fill in the gaps.
            merged: dict[str, Any] = {}
            for source in (bp_defaults, tier_override):
                if source:
                    merged.update(source)
            # Only set a resolved dict when at least one source contributed
            # (or an explicit ``backpressure`` key was present, e.g. ``{}``),
            # so tiers without any backpressure config stay unconfigured.
            if merged or "backpressure" in tier_cfg:
                tier_cfg["backpressure"] = merged

        # Scheduler-side mmap (rank=None); kept for cleanup
        self._scheduler_mmap: SharedOffloadRegion | None = None

        # Set by create_worker when canonical_layout is enabled: True when
        # every layer's canonical bytes are parallelism-agnostic (portable),
        # False when some layers use the opaque fallback (exact-topology only)
        self.all_layers_portable: bool | None = None

        # engine_id is unique per DP replica (suffixed with _dp{rank} in both
        # the Ray and multiprocessing paths), so it names a per-replica offload
        # region.
        self._engine_id = config.engine_id

    @override
    def get_manager(self) -> OffloadingManager:
        """Get the TieringOffloadingManager.

        Creates a TieringOffloadingManager with:
        - Primary tier: CPU (LRU or ARC)
        - Secondary tiers: As configured in extra_config

        Returns:
            TieringOffloadingManager instance

        """
        if not self._manager:
            if int(self.extra_config.get("store_threshold", 0)) >= 2:
                raise ValueError(
                    "store_threshold is not supported for TieringOffloadingSpec"
                )

            scheduler_mmap: SharedOffloadRegion | None = None
            primary_tier: CPUPrimaryTierOffloadingManager | None = None
            secondary_tiers = []
            try:
                # Create scheduler-side SharedOffloadRegion (rank=None) so the
                # primary tier can eagerly create a memoryview over _base.
                scheduler_mmap = SharedOffloadRegion(
                    engine_id=self._engine_id,
                    num_chunks=self.num_chunks,
                    rank=None,
                    kv_bytes_per_chunk=self.kv_bytes_per_chunk,
                    cpu_page_size=self.cpu_page_size_per_worker,
                )
                self._scheduler_mmap = scheduler_mmap

                # Create primary tier (CPU-based)
                primary_tier = CPUPrimaryTierOffloadingManager(
                    num_chunks=self.num_chunks,
                    cache_policy=self.eviction_policy,
                    cache_policy_module_path=self.cache_policy_module_path,
                    enable_events=self.kv_events_config.enable_kv_cache_events,
                    mmap_region=scheduler_mmap,
                )

                # Create secondary tiers
                primary_kv_view = primary_tier.get_kv_memoryview()
                for i, tier_config in enumerate(self.secondary_tier_configs):
                    tier = SecondaryTierFactory.create_secondary_tier(
                        tier_config, primary_kv_view, self
                    )
                    secondary_tiers.append(tier)
                    logger.info(
                        "Created secondary tier #%d (%s)",
                        i,
                        tier.tier_type,
                    )

                # Create TieringOffloadingManager. GPU↔CPU transfers use the inherited
                # get_worker(). Secondary tier transfers are handled by the
                # secondary tier managers and need no additional workers here.
                tiering_manager = TieringOffloadingManager(
                    primary_tier=primary_tier,
                    secondary_tiers=secondary_tiers,
                )
                self._manager = tiering_manager
            except Exception:
                for tier in reversed(secondary_tiers):
                    try:
                        tier.shutdown()
                    except Exception:
                        logger.exception(
                            "Failed to shut down secondary tier during "
                            "initialization cleanup"
                        )
                if primary_tier is not None:
                    try:
                        primary_tier.shutdown()
                    except Exception:
                        logger.exception(
                            "Failed to shut down primary tier during "
                            "initialization cleanup"
                        )
                elif scheduler_mmap is not None:
                    try:
                        scheduler_mmap.cleanup()
                    except Exception:
                        logger.exception(
                            "Failed to clean up scheduler mmap during "
                            "initialization cleanup"
                        )
                self._scheduler_mmap = None
                raise

            logger.info(
                "Created TieringOffloadingManager with primary tier "
                "(%s, %s chunks) and %s secondary tier(s)",
                self.eviction_policy,
                self.num_chunks,
                len(secondary_tiers),
            )

        return self._manager

    @override
    def _uses_shared_region(self) -> bool:
        # Tiering always allocates on the shared region (every platform), so the
        # replicated-layout gate must not be narrowed by the CPU spec's
        # CUDA-alike check.
        return True

    @override
    def create_worker(self, kv_caches: CanonicalKVCaches) -> CPUOffloadingWorker:
        world_size = self.config.parallel.world_size
        if self.replicated_layout:
            rank = 0
        else:
            # Fold the global physical device index into the replica-local
            # [0, world_size) slot range.
            rank = torch.accelerator.current_device_index() % world_size
        worker_mmap = SharedOffloadRegion(
            engine_id=self._engine_id,
            num_chunks=self.num_chunks,
            rank=rank,
            kv_bytes_per_chunk=self.kv_bytes_per_chunk,
            cpu_page_size=self.cpu_page_size_per_worker,
        )
        try:
            if self.config.canonical_layout:
                self._validate_canonical_refs(kv_caches)
            return CPUOffloadingWorker(
                kv_caches=kv_caches,
                blocks_per_chunk=self.blocks_per_chunk,
                num_cpu_chunks=self.num_chunks,
                mmap_region=worker_mmap,
                canonical_layout=self.config.canonical_layout,
            )
        except Exception:
            worker_mmap.cleanup()
            raise

    def _validate_canonical_refs(self, kv_caches: CanonicalKVCaches) -> None:
        """Require a mapping on every ref and record layer portability.

        Fails loudly rather than persist direct-layout bytes under a
        canonical format identity."""
        all_refs = [
            ref for group_refs in kv_caches.group_data_refs for ref in group_refs
        ]
        if any(ref.mapping is None for ref in all_refs):
            raise RuntimeError(
                "canonical_layout was requested but the KV cache layout "
                "could not be certified for canonical offload (offload "
                "workers must be exactly the TP group, and packed / "
                "cross-layer KV layouts are not supported). Remove "
                "canonical_layout from kv_connector_extra_config."
            )
        self.all_layers_portable = all(
            ref.mapping is not None and ref.mapping.parallelism_agnostic
            for ref in all_refs
        )
        if self.config.parallel.is_parallelism_agnostic and not (
            self.all_layers_portable
        ):
            # The scheduler-side storage namespace was already collapsed on
            # the static portability claim; opaque per-topology bytes must
            # not land in it.
            raise RuntimeError(
                "canonical_layout could not certify every layer as "
                "parallelism-agnostic, but the storage namespace is shared "
                "across topologies. Remove canonical_layout from "
                "kv_connector_extra_config or disable parallel-agnostic "
                "secondary tiers."
            )
        logger.info(
            "Canonical KV layout enabled (all_layers_portable=%s)",
            self.all_layers_portable,
        )

_validate_canonical_refs(kv_caches)

Require a mapping on every ref and record layer portability.

Fails loudly rather than persist direct-layout bytes under a canonical format identity.

Source code in vllm/v1/kv_offload/tiering/spec.py
def _validate_canonical_refs(self, kv_caches: CanonicalKVCaches) -> None:
    """Require a mapping on every ref and record layer portability.

    Fails loudly rather than persist direct-layout bytes under a
    canonical format identity."""
    all_refs = [
        ref for group_refs in kv_caches.group_data_refs for ref in group_refs
    ]
    if any(ref.mapping is None for ref in all_refs):
        raise RuntimeError(
            "canonical_layout was requested but the KV cache layout "
            "could not be certified for canonical offload (offload "
            "workers must be exactly the TP group, and packed / "
            "cross-layer KV layouts are not supported). Remove "
            "canonical_layout from kv_connector_extra_config."
        )
    self.all_layers_portable = all(
        ref.mapping is not None and ref.mapping.parallelism_agnostic
        for ref in all_refs
    )
    if self.config.parallel.is_parallelism_agnostic and not (
        self.all_layers_portable
    ):
        # The scheduler-side storage namespace was already collapsed on
        # the static portability claim; opaque per-topology bytes must
        # not land in it.
        raise RuntimeError(
            "canonical_layout could not certify every layer as "
            "parallelism-agnostic, but the storage namespace is shared "
            "across topologies. Remove canonical_layout from "
            "kv_connector_extra_config or disable parallel-agnostic "
            "secondary tiers."
        )
    logger.info(
        "Canonical KV layout enabled (all_layers_portable=%s)",
        self.all_layers_portable,
    )

get_manager()

Get the TieringOffloadingManager.

Creates a TieringOffloadingManager with: - Primary tier: CPU (LRU or ARC) - Secondary tiers: As configured in extra_config

Returns:

Source code in vllm/v1/kv_offload/tiering/spec.py
@override
def get_manager(self) -> OffloadingManager:
    """Get the TieringOffloadingManager.

    Creates a TieringOffloadingManager with:
    - Primary tier: CPU (LRU or ARC)
    - Secondary tiers: As configured in extra_config

    Returns:
        TieringOffloadingManager instance

    """
    if not self._manager:
        if int(self.extra_config.get("store_threshold", 0)) >= 2:
            raise ValueError(
                "store_threshold is not supported for TieringOffloadingSpec"
            )

        scheduler_mmap: SharedOffloadRegion | None = None
        primary_tier: CPUPrimaryTierOffloadingManager | None = None
        secondary_tiers = []
        try:
            # Create scheduler-side SharedOffloadRegion (rank=None) so the
            # primary tier can eagerly create a memoryview over _base.
            scheduler_mmap = SharedOffloadRegion(
                engine_id=self._engine_id,
                num_chunks=self.num_chunks,
                rank=None,
                kv_bytes_per_chunk=self.kv_bytes_per_chunk,
                cpu_page_size=self.cpu_page_size_per_worker,
            )
            self._scheduler_mmap = scheduler_mmap

            # Create primary tier (CPU-based)
            primary_tier = CPUPrimaryTierOffloadingManager(
                num_chunks=self.num_chunks,
                cache_policy=self.eviction_policy,
                cache_policy_module_path=self.cache_policy_module_path,
                enable_events=self.kv_events_config.enable_kv_cache_events,
                mmap_region=scheduler_mmap,
            )

            # Create secondary tiers
            primary_kv_view = primary_tier.get_kv_memoryview()
            for i, tier_config in enumerate(self.secondary_tier_configs):
                tier = SecondaryTierFactory.create_secondary_tier(
                    tier_config, primary_kv_view, self
                )
                secondary_tiers.append(tier)
                logger.info(
                    "Created secondary tier #%d (%s)",
                    i,
                    tier.tier_type,
                )

            # Create TieringOffloadingManager. GPU↔CPU transfers use the inherited
            # get_worker(). Secondary tier transfers are handled by the
            # secondary tier managers and need no additional workers here.
            tiering_manager = TieringOffloadingManager(
                primary_tier=primary_tier,
                secondary_tiers=secondary_tiers,
            )
            self._manager = tiering_manager
        except Exception:
            for tier in reversed(secondary_tiers):
                try:
                    tier.shutdown()
                except Exception:
                    logger.exception(
                        "Failed to shut down secondary tier during "
                        "initialization cleanup"
                    )
            if primary_tier is not None:
                try:
                    primary_tier.shutdown()
                except Exception:
                    logger.exception(
                        "Failed to shut down primary tier during "
                        "initialization cleanup"
                    )
            elif scheduler_mmap is not None:
                try:
                    scheduler_mmap.cleanup()
                except Exception:
                    logger.exception(
                        "Failed to clean up scheduler mmap during "
                        "initialization cleanup"
                    )
            self._scheduler_mmap = None
            raise

        logger.info(
            "Created TieringOffloadingManager with primary tier "
            "(%s, %s chunks) and %s secondary tier(s)",
            self.eviction_policy,
            self.num_chunks,
            len(secondary_tiers),
        )

    return self._manager