[BugFix] Revert "[KV Offload] Use background thread for mmap / cpu_tensors pinning" (#46958)

Signed-off-by: <>
Co-authored-by: Varun Sundar Rabindranath <varun-sundar-rabindranath@h100-01.nemg-001.lab.rdu2.dc.redhat.com>
This commit is contained in:
Varun Sundar Rabindranath
2026-06-29 07:54:34 -07:00
committed by GitHub
co-authored by Varun Sundar Rabindranath
parent 6149187a4c
commit 36bbecd643
+38 -89
View File
@@ -1,7 +1,6 @@
# SPDX-License-Identifier: Apache-2.0
# SPDX-FileCopyrightText: Copyright contributors to the vLLM project
import functools
import threading
import time
from collections import deque
from dataclasses import dataclass
@@ -121,6 +120,36 @@ def compute_sub_block_ptrs(
output[:] = flat[skip_count : skip_count + num_sub_blocks]
def pin_mmap_region(region: SharedOffloadRegion) -> None:
"""Register the entire mmap as CUDA pinned memory via cudaHostRegister."""
if not current_platform.is_cuda_alike():
logger.info(
"Skipping mmap host registration on %s; cudaHostRegister is only "
"available on CUDA/ROCm.",
current_platform.device_name,
)
return
rank = region.rank
base_ptr = region._base.data_ptr()
result = torch.cuda.cudart().cudaHostRegister(base_ptr, region.total_size_bytes, 0)
if result.value != 0:
logger.warning(
"cudaHostRegister failed for rank=%d (code=%d) — "
"transfers will still work but may be slower (unpinned DMA)",
rank,
result,
)
else:
logger.debug(
"cudaHostRegister rank=%d %.2f GB",
rank,
region.total_size_bytes / 1e9,
)
region.is_pinned = True
def _new_descriptor_buffers(
num_copy_ops: int,
) -> tuple[torch.Tensor, torch.Tensor, torch.Tensor]:
@@ -150,8 +179,6 @@ class SingleDirectionOffloadingHandler:
kv_cache_groups_data_refs: list[list[CanonicalKVCacheRef]],
gpu_to_cpu: bool,
mmap_region: SharedOffloadRegion | None = None,
pin_thread: threading.Thread | None = None,
manually_pinned_tensors: list[torch.Tensor] | None = None,
):
"""
Initialize a SingleDirectionOffloadingHandler.
@@ -199,8 +226,6 @@ class SingleDirectionOffloadingHandler:
# mmap_region to clean up on shutdown (gpu_to_cpu handler owns it)
self._mmap_region = mmap_region
self._pin_thread = pin_thread
self._manually_pinned_tensors = manually_pinned_tensors
# job_id -> event
self._transfer_events: dict[int, torch.Event] = {}
# queue of transfers (job_id, stream, event)
@@ -433,23 +458,8 @@ class SingleDirectionOffloadingHandler:
self._stream_pool.clear()
self._event_pool.clear()
self._buffer_pool.clear()
if self._pin_thread is not None:
self._pin_thread.join()
self._pin_thread = None
if self._manually_pinned_tensors is not None:
for tensor in self._manually_pinned_tensors:
result = torch.cuda.cudart().cudaHostUnregister(tensor.data_ptr())
if result.value != 0:
logger.warning(
"cudaHostUnregister failed for CPU tensor (code=%d)",
result.value,
)
self.src_tensors.clear()
self.dst_tensors.clear()
if self._mmap_region is not None:
self._mmap_region.cleanup()
self._mmap_region = None
@@ -471,14 +481,12 @@ class CPUOffloadingWorker(OffloadingWorker):
mmap_region: SharedOffloadRegion | None = None,
):
pin_memory = PIN_MEMORY
self.pin_thread: threading.Thread | None = None
self._manually_pinned_tensors: list[torch.Tensor] = []
logger.info("Allocating %d CPU tensors...", len(kv_caches.tensors))
self._mmap_region = mmap_region
if mmap_region is not None and pin_memory:
pin_mmap_region(mmap_region)
gpu_tensors: list[torch.Tensor] = []
self.cpu_tensors: list[torch.Tensor] = []
cpu_tensors: list[torch.Tensor] = []
for kv_cache_tensor in kv_caches.tensors:
gpu_page_size_bytes = kv_cache_tensor.page_size_bytes
gpu_tensor = kv_cache_tensor.tensor.view(torch.int8).view(
@@ -494,13 +502,10 @@ class CPUOffloadingWorker(OffloadingWorker):
(num_cpu_blocks, cpu_page_size_bytes),
dtype=torch.int8,
device="cpu",
# CUDA/ROCm memory is registered asynchronously below.
# Pinning here would block worker initialization; other
# hardware need PyTorch allocation-time pinning.
pin_memory=PIN_MEMORY and not current_platform.is_cuda_alike(),
pin_memory=pin_memory,
)
logger.debug(
"torch.zeros tensor %d×%d (%.2f GB): %.3f s",
"torch.zeros pinned tensor %d×%d (%.2f GB): %.3f s",
num_cpu_blocks,
cpu_page_size_bytes,
num_cpu_blocks * cpu_page_size_bytes / 1e9,
@@ -508,81 +513,25 @@ class CPUOffloadingWorker(OffloadingWorker):
)
gpu_tensors.append(gpu_tensor)
self.cpu_tensors.append(cpu_tensor)
if pin_memory:
if not current_platform.is_cuda_alike():
logger.info(
"Skipping host registration on %s; cudaHostRegister is only "
"available on CUDA/ROCm.",
current_platform.device_name,
)
else:
self.pin_thread = threading.Thread(
target=self._pin_cpu_tensors,
name="CPUTensorPinThread",
)
self.pin_thread.start()
logger.info("Starting to pin memory in background...")
cpu_tensors.append(cpu_tensor)
self._store_handler = SingleDirectionOffloadingHandler(
gpu_tensors=gpu_tensors,
cpu_tensors=self.cpu_tensors,
cpu_tensors=cpu_tensors,
block_size_factor=block_size_factor,
kv_cache_groups_data_refs=kv_caches.group_data_refs,
gpu_to_cpu=True,
mmap_region=mmap_region,
pin_thread=self.pin_thread,
manually_pinned_tensors=self._manually_pinned_tensors,
)
self._load_handler = SingleDirectionOffloadingHandler(
gpu_tensors=gpu_tensors,
cpu_tensors=self.cpu_tensors,
cpu_tensors=cpu_tensors,
block_size_factor=block_size_factor,
kv_cache_groups_data_refs=kv_caches.group_data_refs,
gpu_to_cpu=False,
)
def _pin_cpu_tensors(self) -> None:
"""Register the CPU offload memory as CUDA pinned memory."""
t0 = time.monotonic()
tensors_to_pin = (
[self._mmap_region._base]
if self._mmap_region is not None
else self.cpu_tensors
)
num_pinned = 0
for tensor in tensors_to_pin:
total_size_bytes = tensor.numel() * tensor.element_size()
result = torch.cuda.cudart().cudaHostRegister(
tensor.data_ptr(), total_size_bytes, 0
)
if result.value != 0:
logger.warning(
"cudaHostRegister failed for host tensor (code=%d) "
"- transfers will still work but may be slower (unpinned DMA)",
result.value,
)
continue
if self._mmap_region is not None:
self._mmap_region.is_pinned = True
else:
self._manually_pinned_tensors.append(tensor)
num_pinned += 1
logger.debug(
"cudaHostRegister pin %.2f GB",
total_size_bytes / 1e9,
)
logger.info(
"Completed CPU memory pinning: %d tensors pinned in %.3f s",
num_pinned,
time.monotonic() - t0,
)
def submit_store(
self, job_id: int, src_spec: GPULoadStoreSpec, dst_spec: LoadStoreSpec
) -> bool: