Remove full-cache scans from CP owner-lane allocation
The CP shared-KV allocator was still doing total-cache-sized CPU work in the scheduler hot path. That cannot be hidden by GPU overlap, so owner-lane allocation now maintains per-owner free/release buckets and consumes request-sized prefixes instead of rebuilding masks over the full free-page tensor on each request.\n\nThe benchmark was extended to isolate L1 stats, selection, and allocation costs, and the CPU layout tests now install a complete sgl_kernel stub before importing SGLang helpers so remote unit collection does not abort in native extension loading.\n\nConstraint: Allocator CPU work blocks scheduler progress and cannot overlap with GPU forward execution.\nConstraint: CPU unit tests must not load native sgl_kernel on remote images where the loader can SIGABRT.\nRejected: Keep contiguous-run search over full free_pages | still scales with cache capacity and measured multi-ms overhead.\nRejected: Treat remote collection abort as an environment-only issue | it prevented allocator regression coverage and was fixable with a test-local stub.\nConfidence: high\nScope-risk: moderate\nDirective: CP owner-lane allocation is bucket-based; do not reintroduce full free_pages scans on the hot path without benchmark evidence.\nTested: Local py_compile for touched files\nTested: Local benchmark unit test, 6 passed\nTested: Remote benchmark unit test, 6 passed\nTested: Remote test_alloc_pages_with_owners.py, 10 passed\nTested: Remote test_cp_shared_kv_layout.py, 27 passed\nTested: Remote production allocator microbench shows select/alloc p50 reduced from ms-scale to sub-ms scale\nNot-tested: Full ETE traffic run after allocator bucket change
This commit is contained in:
@@ -14,10 +14,11 @@ Examples:
|
||||
--bench host --host-sizes-gb 220 --request-pages 1,8,64,512 \
|
||||
--patterns contiguous_fifo,fragmented_prefix_later_run,random_fragmented
|
||||
|
||||
# Production L1 allocator path on CUDA, stubbing sgl_kernel import if needed.
|
||||
# L1 allocator path on CUDA, stubbing sgl_kernel import if needed.
|
||||
PYTHONPATH=python:. python benchmark/hicache/bench_cp_hicache_allocator_overhead.py \
|
||||
--bench l1 --device cuda --stub-sgl-kernel --physical-pages 8192,32768 \
|
||||
--request-pages 8,64,512 --l1-impl current,fifo
|
||||
--request-pages 8,64,512 --l1-impl current,fifo \
|
||||
--l1-ops stats,free_room_stats,select_only,alloc_pages
|
||||
"""
|
||||
|
||||
import argparse
|
||||
@@ -244,6 +245,296 @@ class StandaloneHostAllocator:
|
||||
return select_index
|
||||
|
||||
|
||||
def _compute_owner_lane_free_room_deficits(
|
||||
*,
|
||||
required: list[int],
|
||||
available: list[int],
|
||||
capacities: list[int],
|
||||
target_ratio: float,
|
||||
trigger_ratio: float,
|
||||
) -> list[int]:
|
||||
deficits: list[int] = []
|
||||
for req, avail, capacity in zip(required, available, capacities):
|
||||
target_room = (
|
||||
int(math.ceil(float(capacity) * float(target_ratio)))
|
||||
if capacity > 0 and target_ratio > 0
|
||||
else 0
|
||||
)
|
||||
trigger_room = (
|
||||
int(math.ceil(float(capacity) * float(trigger_ratio)))
|
||||
if capacity > 0 and trigger_ratio > 0
|
||||
else 0
|
||||
)
|
||||
if int(avail) >= int(req) + trigger_room:
|
||||
deficits.append(0)
|
||||
else:
|
||||
deficits.append(max(0, int(req) + target_room - int(avail)))
|
||||
return deficits
|
||||
|
||||
|
||||
class StandaloneCPSharedPagedAllocator:
|
||||
"""Metadata-only copy of the CP shared-KV page owner allocator.
|
||||
|
||||
This intentionally mirrors the current Python/Torch control path used by
|
||||
``CPSharedPagedTokenToKVPoolAllocator`` so CPU-only environments can measure
|
||||
the allocator shape without importing the full SGLang runtime dependency
|
||||
stack. The benchmark still uses the production allocator when imports are
|
||||
available.
|
||||
"""
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
*,
|
||||
physical_pages: int,
|
||||
page_size: int,
|
||||
cp_size: int,
|
||||
device: torch.device,
|
||||
):
|
||||
self.physical_size = int(physical_pages) * int(page_size)
|
||||
self.page_size = int(page_size)
|
||||
self.cp_size = int(cp_size)
|
||||
self.device = device
|
||||
logical_pages = int(physical_pages) * int(cp_size)
|
||||
self._owner_free_pages = None
|
||||
self._owner_release_pages = None
|
||||
self._flat_free_pages_cache = None
|
||||
self._flat_release_pages_cache = None
|
||||
self.free_pages = torch.arange(
|
||||
1, logical_pages + 1, dtype=torch.int64, device=device
|
||||
)
|
||||
self.release_pages = torch.empty((0,), dtype=torch.int64, device=device)
|
||||
self.debug_mode = False
|
||||
|
||||
def _empty_pages(self) -> torch.Tensor:
|
||||
return torch.empty((0,), dtype=torch.int64, device=self.device)
|
||||
|
||||
def _split_owner_buckets(
|
||||
self, pages: Optional[torch.Tensor]
|
||||
) -> Optional[list[torch.Tensor]]:
|
||||
if pages is None:
|
||||
return None
|
||||
if pages.numel() == 0:
|
||||
return [torch.empty_like(pages) for _ in range(self.cp_size)]
|
||||
owner_ids = torch.remainder(pages - 1, self.cp_size)
|
||||
buckets: list[torch.Tensor] = []
|
||||
for owner in range(self.cp_size):
|
||||
owner_pages = pages[owner_ids == owner]
|
||||
if owner_pages.numel() > 1:
|
||||
owner_pages, _ = torch.sort(owner_pages)
|
||||
buckets.append(owner_pages)
|
||||
return buckets
|
||||
|
||||
def _materialize_owner_buckets(
|
||||
self, buckets: Optional[list[torch.Tensor]]
|
||||
) -> Optional[torch.Tensor]:
|
||||
if buckets is None:
|
||||
return None
|
||||
non_empty = [bucket for bucket in buckets if bucket.numel() > 0]
|
||||
if not non_empty:
|
||||
return self._empty_pages()
|
||||
return torch.cat(non_empty)
|
||||
|
||||
@property
|
||||
def free_pages(self):
|
||||
if self._owner_free_pages is not None:
|
||||
if self._flat_free_pages_cache is None:
|
||||
self._flat_free_pages_cache = self._materialize_owner_buckets(
|
||||
self._owner_free_pages
|
||||
)
|
||||
return self._flat_free_pages_cache
|
||||
return self._flat_free_pages_cache
|
||||
|
||||
@free_pages.setter
|
||||
def free_pages(self, pages):
|
||||
self._flat_free_pages_cache = pages
|
||||
self._owner_free_pages = self._split_owner_buckets(pages)
|
||||
|
||||
@property
|
||||
def release_pages(self):
|
||||
if self._owner_release_pages is not None:
|
||||
if self._flat_release_pages_cache is None:
|
||||
self._flat_release_pages_cache = self._materialize_owner_buckets(
|
||||
self._owner_release_pages
|
||||
)
|
||||
return self._flat_release_pages_cache
|
||||
return self._flat_release_pages_cache
|
||||
|
||||
@release_pages.setter
|
||||
def release_pages(self, pages):
|
||||
self._flat_release_pages_cache = pages
|
||||
self._owner_release_pages = self._split_owner_buckets(pages)
|
||||
|
||||
def _owner_bucket_counts(self, buckets: Optional[list[torch.Tensor]]) -> list[int]:
|
||||
if buckets is None:
|
||||
return [0 for _ in range(self.cp_size)]
|
||||
return [int(bucket.numel()) for bucket in buckets]
|
||||
|
||||
def _owner_available_counts(self) -> list[int]:
|
||||
free_counts = self._owner_bucket_counts(self._owner_free_pages)
|
||||
release_counts = self._owner_bucket_counts(self._owner_release_pages)
|
||||
return [
|
||||
free_count + release_count
|
||||
for free_count, release_count in zip(free_counts, release_counts)
|
||||
]
|
||||
|
||||
def _consume_owner_bucket_prefix(
|
||||
self,
|
||||
*,
|
||||
release: bool,
|
||||
counts_by_owner: list[int],
|
||||
) -> None:
|
||||
target_attr = "_owner_release_pages" if release else "_owner_free_pages"
|
||||
cache_attr = "_flat_release_pages_cache" if release else "_flat_free_pages_cache"
|
||||
buckets = getattr(self, target_attr)
|
||||
mutated = False
|
||||
for owner, count in enumerate(counts_by_owner):
|
||||
if count <= 0:
|
||||
continue
|
||||
buckets[owner] = buckets[owner][count:]
|
||||
mutated = True
|
||||
if mutated:
|
||||
setattr(self, target_attr, buckets)
|
||||
setattr(self, cache_attr, None)
|
||||
|
||||
def compute_owner_lane_stats(
|
||||
self,
|
||||
page_compute_owners: list[int],
|
||||
) -> tuple[list[int], list[int], list[int]]:
|
||||
required = [0 for _ in range(self.cp_size)]
|
||||
for owner in page_compute_owners:
|
||||
if owner < 0 or owner >= self.cp_size:
|
||||
raise ValueError(
|
||||
f"compute owner must be in [0, {self.cp_size}), got {owner}"
|
||||
)
|
||||
required[owner] += 1
|
||||
|
||||
available = self._owner_available_counts()
|
||||
deficits = [
|
||||
max(0, required_count - available_count)
|
||||
for required_count, available_count in zip(required, available)
|
||||
]
|
||||
return required, available, deficits
|
||||
|
||||
def compute_owner_lane_capacity_pages(self) -> list[int]:
|
||||
capacity_pages = int(self.physical_size // self.page_size)
|
||||
return [capacity_pages for _ in range(self.cp_size)]
|
||||
|
||||
def compute_owner_lane_free_room_stats(
|
||||
self,
|
||||
page_compute_owners: list[int],
|
||||
*,
|
||||
target_ratio: float,
|
||||
trigger_ratio: float,
|
||||
) -> tuple[list[int], list[int], list[int]]:
|
||||
required, available, _exact_deficits = self.compute_owner_lane_stats(
|
||||
page_compute_owners
|
||||
)
|
||||
deficits = _compute_owner_lane_free_room_deficits(
|
||||
required=required,
|
||||
available=available,
|
||||
capacities=self.compute_owner_lane_capacity_pages(),
|
||||
target_ratio=target_ratio,
|
||||
trigger_ratio=trigger_ratio,
|
||||
)
|
||||
return required, available, deficits
|
||||
|
||||
def _select_owner_free_pages_prefer_contiguous(
|
||||
self,
|
||||
owner_pages: torch.Tensor,
|
||||
required_count: int,
|
||||
) -> torch.Tensor:
|
||||
return owner_pages[: min(required_count, int(owner_pages.numel()))]
|
||||
|
||||
def _select_compute_owner_pages(
|
||||
self,
|
||||
page_compute_owners: list[int],
|
||||
) -> Optional[tuple[torch.Tensor, list[int], list[int]]]:
|
||||
if not page_compute_owners:
|
||||
return (
|
||||
torch.empty((0,), dtype=torch.int64, device=self.device),
|
||||
[0 for _ in range(self.cp_size)],
|
||||
[0 for _ in range(self.cp_size)],
|
||||
)
|
||||
|
||||
required_by_owner = [0 for _ in range(self.cp_size)]
|
||||
positions_by_owner: list[list[int]] = [[] for _ in range(self.cp_size)]
|
||||
for position, owner in enumerate(page_compute_owners):
|
||||
if owner < 0 or owner >= self.cp_size:
|
||||
raise ValueError(
|
||||
f"compute owner must be in [0, {self.cp_size}), got {owner}"
|
||||
)
|
||||
required_by_owner[owner] += 1
|
||||
positions_by_owner[owner].append(position)
|
||||
|
||||
lane_pages = [None for _ in range(self.cp_size)]
|
||||
selected_free_counts = [0 for _ in range(self.cp_size)]
|
||||
selected_release_counts = [0 for _ in range(self.cp_size)]
|
||||
for owner, required_count in enumerate(required_by_owner):
|
||||
if required_count == 0:
|
||||
continue
|
||||
|
||||
selected_owner_free_mask = self._select_owner_free_pages_prefer_contiguous(
|
||||
self._owner_free_pages[owner], required_count
|
||||
)
|
||||
selected_owner_pages = selected_owner_free_mask
|
||||
|
||||
free_count = int(selected_owner_pages.numel())
|
||||
remaining_count = required_count - free_count
|
||||
if remaining_count > 0:
|
||||
release_bucket = self._owner_release_pages[owner]
|
||||
if remaining_count > release_bucket.numel():
|
||||
return None
|
||||
selected_owner_release_pages = release_bucket[:remaining_count]
|
||||
selected_owner_pages = torch.cat(
|
||||
(selected_owner_pages, selected_owner_release_pages)
|
||||
)
|
||||
selected_release_counts[owner] = remaining_count
|
||||
|
||||
selected_free_counts[owner] = free_count
|
||||
lane_pages[owner] = selected_owner_pages
|
||||
|
||||
selected_pages = torch.empty(
|
||||
(len(page_compute_owners),), dtype=torch.int64, device=self.device
|
||||
)
|
||||
for owner, positions in enumerate(positions_by_owner):
|
||||
if not positions:
|
||||
continue
|
||||
position_tensor = torch.tensor(
|
||||
positions, dtype=torch.int64, device=self.device
|
||||
)
|
||||
selected_pages[position_tensor] = lane_pages[owner]
|
||||
|
||||
return (
|
||||
selected_pages,
|
||||
selected_free_counts,
|
||||
selected_release_counts,
|
||||
)
|
||||
|
||||
def alloc_pages_with_owners(
|
||||
self,
|
||||
page_compute_owners: list[int],
|
||||
) -> Optional[torch.Tensor]:
|
||||
if not page_compute_owners:
|
||||
return torch.empty((0,), dtype=torch.int64, device=self.device)
|
||||
selected = self._select_compute_owner_pages(page_compute_owners)
|
||||
if selected is None:
|
||||
return None
|
||||
selected_pages, selected_free_counts, selected_release_counts = selected
|
||||
page_size = self.page_size
|
||||
base = selected_pages.to(torch.int64).unsqueeze(1) * page_size
|
||||
offsets = torch.arange(
|
||||
page_size, dtype=torch.int64, device=self.device
|
||||
).unsqueeze(0)
|
||||
out_indices = (base + offsets).reshape(-1)
|
||||
self._consume_owner_bucket_prefix(
|
||||
release=False, counts_by_owner=selected_free_counts
|
||||
)
|
||||
self._consume_owner_bucket_prefix(
|
||||
release=True, counts_by_owner=selected_release_counts
|
||||
)
|
||||
return out_indices
|
||||
|
||||
|
||||
def _is_page_contiguous_selection(selected: Optional[torch.Tensor], page_size: int) -> bool:
|
||||
if selected is None or selected.numel() == 0:
|
||||
return False
|
||||
@@ -477,7 +768,22 @@ def _install_sgl_kernel_stubs() -> None:
|
||||
sys.modules[submodule] = sub
|
||||
|
||||
|
||||
def _make_l1_allocator(*, physical_pages: int, page_size: int, cp_size: int, device: torch.device):
|
||||
def _make_l1_allocator(
|
||||
*,
|
||||
physical_pages: int,
|
||||
page_size: int,
|
||||
cp_size: int,
|
||||
device: torch.device,
|
||||
production: bool,
|
||||
):
|
||||
if not production:
|
||||
return StandaloneCPSharedPagedAllocator(
|
||||
physical_pages=physical_pages,
|
||||
page_size=page_size,
|
||||
cp_size=cp_size,
|
||||
device=device,
|
||||
)
|
||||
|
||||
from sglang.srt.mem_cache.allocator import CPSharedPagedTokenToKVPoolAllocator
|
||||
|
||||
return CPSharedPagedTokenToKVPoolAllocator(
|
||||
@@ -494,8 +800,8 @@ def _make_l1_allocator(*, physical_pages: int, page_size: int, cp_size: int, dev
|
||||
|
||||
|
||||
def _patch_l1_fifo_selector(allocator) -> None:
|
||||
def fifo_selector(owner_mask: torch.Tensor, required_count: int) -> torch.Tensor:
|
||||
return owner_mask & (torch.cumsum(owner_mask.to(torch.int64), dim=0) <= required_count)
|
||||
def fifo_selector(owner_pages: torch.Tensor, required_count: int) -> torch.Tensor:
|
||||
return owner_pages[: min(required_count, int(owner_pages.numel()))]
|
||||
|
||||
allocator._select_owner_free_pages_prefer_contiguous = fifo_selector
|
||||
|
||||
@@ -520,6 +826,7 @@ def _is_l1_selection_physically_contiguous(
|
||||
def _bench_l1_case(
|
||||
*,
|
||||
impl: str,
|
||||
op: str,
|
||||
physical_pages: int,
|
||||
request_pages: int,
|
||||
page_size: int,
|
||||
@@ -530,6 +837,9 @@ def _bench_l1_case(
|
||||
repeat: int,
|
||||
warmup: int,
|
||||
seed: int,
|
||||
free_room_ratio: float,
|
||||
free_room_trigger_ratio: float,
|
||||
production_allocator: bool,
|
||||
) -> BenchResult:
|
||||
request_owners = _make_page_compute_owners(request_pages, cp_size, owner_pattern)
|
||||
base_free_pages = _make_l1_free_pages(
|
||||
@@ -551,6 +861,7 @@ def _bench_l1_case(
|
||||
page_size=page_size,
|
||||
cp_size=cp_size,
|
||||
device=device,
|
||||
production=production_allocator,
|
||||
)
|
||||
allocator.free_pages = base_free_pages.clone()
|
||||
allocator.release_pages = torch.empty((0,), dtype=torch.int64, device=device)
|
||||
@@ -559,7 +870,28 @@ def _bench_l1_case(
|
||||
if use_cuda:
|
||||
torch.cuda.synchronize(device)
|
||||
start_ns = time.perf_counter_ns()
|
||||
selected = allocator.alloc_pages_with_owners(request_owners)
|
||||
selected = None
|
||||
if op == "stats":
|
||||
allocator.compute_owner_lane_stats(request_owners)
|
||||
elif op == "free_room_stats":
|
||||
allocator.compute_owner_lane_free_room_stats(
|
||||
request_owners,
|
||||
target_ratio=free_room_ratio,
|
||||
trigger_ratio=free_room_trigger_ratio,
|
||||
)
|
||||
elif op == "select_only":
|
||||
selected_result = allocator._select_compute_owner_pages(request_owners)
|
||||
if selected_result is not None:
|
||||
selected_pages = selected_result[0]
|
||||
base = selected_pages.to(torch.int64).unsqueeze(1) * page_size
|
||||
offsets = torch.arange(
|
||||
page_size, dtype=torch.int64, device=device
|
||||
).unsqueeze(0)
|
||||
selected = (base + offsets).reshape(-1)
|
||||
elif op == "alloc_pages":
|
||||
selected = allocator.alloc_pages_with_owners(request_owners)
|
||||
else:
|
||||
raise ValueError(f"unsupported l1 op: {op}")
|
||||
if use_cuda:
|
||||
torch.cuda.synchronize(device)
|
||||
elapsed_us = (time.perf_counter_ns() - start_ns) / 1000.0
|
||||
@@ -572,7 +904,7 @@ def _bench_l1_case(
|
||||
)
|
||||
return _summarize(
|
||||
bench="l1",
|
||||
impl=impl,
|
||||
impl=f"{impl}:{op}",
|
||||
pattern=f"{free_pattern}:{owner_pattern}",
|
||||
device=device.type,
|
||||
total_pages=physical_pages,
|
||||
@@ -648,6 +980,7 @@ def _run_l1(args) -> list[BenchResult]:
|
||||
free_patterns = [item.strip() for item in args.l1_free_patterns.split(",") if item.strip()]
|
||||
owner_patterns = [item.strip() for item in args.l1_owner_patterns.split(",") if item.strip()]
|
||||
impls = [item.strip() for item in args.l1_impl.split(",") if item.strip()]
|
||||
ops = [item.strip() for item in args.l1_ops.split(",") if item.strip()]
|
||||
|
||||
results: list[BenchResult] = []
|
||||
for physical_pages in physical_pages_list:
|
||||
@@ -657,21 +990,26 @@ def _run_l1(args) -> list[BenchResult]:
|
||||
for free_pattern in free_patterns:
|
||||
for owner_pattern in owner_patterns:
|
||||
for impl in impls:
|
||||
results.append(
|
||||
_bench_l1_case(
|
||||
impl=impl,
|
||||
physical_pages=physical_pages,
|
||||
request_pages=request_pages,
|
||||
page_size=args.page_size,
|
||||
cp_size=args.cp_size,
|
||||
device=device,
|
||||
free_pattern=free_pattern,
|
||||
owner_pattern=owner_pattern,
|
||||
repeat=args.repeat,
|
||||
warmup=args.warmup,
|
||||
seed=args.seed,
|
||||
for op in ops:
|
||||
results.append(
|
||||
_bench_l1_case(
|
||||
impl=impl,
|
||||
op=op,
|
||||
physical_pages=physical_pages,
|
||||
request_pages=request_pages,
|
||||
page_size=args.page_size,
|
||||
cp_size=args.cp_size,
|
||||
device=device,
|
||||
free_pattern=free_pattern,
|
||||
owner_pattern=owner_pattern,
|
||||
repeat=args.repeat,
|
||||
warmup=args.warmup,
|
||||
seed=args.seed,
|
||||
free_room_ratio=args.l1_free_room_ratio,
|
||||
free_room_trigger_ratio=args.l1_free_room_trigger_ratio,
|
||||
production_allocator=args.l1_allocator == "production",
|
||||
)
|
||||
)
|
||||
)
|
||||
return results
|
||||
|
||||
|
||||
@@ -705,6 +1043,15 @@ def _build_parser() -> argparse.ArgumentParser:
|
||||
parser.add_argument("--cuda-device", type=int, default=0)
|
||||
parser.add_argument("--stub-sgl-kernel", action="store_true")
|
||||
parser.add_argument("--l1-impl", default="fifo,current")
|
||||
parser.add_argument("--l1-ops", default="alloc_pages")
|
||||
parser.add_argument(
|
||||
"--l1-allocator",
|
||||
choices=("production", "standalone"),
|
||||
default="production",
|
||||
help="Use production allocator imports or the dependency-light metadata copy.",
|
||||
)
|
||||
parser.add_argument("--l1-free-room-ratio", type=float, default=0.15)
|
||||
parser.add_argument("--l1-free-room-trigger-ratio", type=float, default=0.05)
|
||||
parser.add_argument(
|
||||
"--l1-free-patterns",
|
||||
default="sequential,owner_fragmented_later_run,random",
|
||||
|
||||
Reference in New Issue
Block a user