大语言模型推理服务的多级优先级调度 — 从请求路由到内存预留的生产级实现
当 LLM 推理集群同时为实时对话、在线 API 任务和离线批量推理服务时,如何确保高优先级请求不被饿死,同时最大化 GPU 利用率?本文从生产视角拆解 LLM 推理服务的多级优先级调度架构。
一、问题定义:为什么 LLM 推理需要优先级调度?
传统 Web 服务的请求调度相对简单——无状态、短期执行、资源可预测。但 LLM 推理服务有三个本质不同:
- 资源占用不对等:一条生成长度 2048 token 的请求可能独占 GPU KV Cache 数分钟,而一条分类任务只需 50ms
- 抢占成本极高:中断正在生成的 KV Cache 意味着已计算的中间状态作废,显存无法即时回收
- 批处理延迟敏感:Continuous Batching 下新请求必须等当前 batch 完成才能加入,排队延迟直接叠加
生产级 LLM 推理服务的典型场景:
| 优先级 | 延迟要求 | 典型占比 | 场景 |
|---|---|---|---|
| P0 实时 | <500ms | 5% | 多轮对话/工具调用 |
| P1 在线 | <5s | 60% | API 聊天/摘要 |
| P2 批量 | <60s | 30% | 数据处理/分类 |
| P3 后台 | 无限制 | 5% | 评测/蒸馏/导出 |
二、调度架构总览
一个完整的 LLM 推理优先级调度系统从 API 网关到 GPU 显存,贯穿四个层级:
┌─────────────────┐
│ 全局 Token 桶 │ ← 全局限流
└────────┬────────┘
│
┌─────────────▼─────────────┐
│ 优先级队列管理器 │ ← 调度核心
│ ┌───┬───┬───┬───┐ │
│ │P0 │P1 │P2 │P3 │ │
│ └───┴───┴───┴───┘ │
└─────────────┬─────────────┘
│
┌───────────────────▼──────────────────┐
│ KV Cache 分区管理器 │ ← 内存保障
│ P0预留区 │ P1弹性区 │ P2共享区 │ P3借用区 │
└───────────────────┬──────────────────┘
│
┌─────────────▼─────────────┐
│ Continuous Batching │ ← 执行引擎
└───────────────────────────┘
2.1 关键设计原则
- 优先级不意味着独占:高优先级请求享有席位保证而非带宽保证
- 防饥饿机制:低优先级请求定期获得"提升"机会
- 内存隔离可降级:当 P0 预留区空闲时,允许低优先级临时借用
- 调度决策前置:在请求进入引擎前完成所有准入控制
三、调度算法:从严格优先级到公平份额
3.1 严格优先级(Strict Priority)
最简单的实现——总是优先服务最高优先级非空队列:
class StrictPriorityScheduler:
def __init__(self, num_priorities: int = 4):
self.queues: List[deque[Request]] = [deque() for _ in range(num_priorities)]
def enqueue(self, request: Request) -> None:
self.queues[request.priority].append(request)
def dequeue(self) -> Optional[Request]:
for priority in range(len(self.queues)):
if self.queues[priority]:
return self.queues[priority].popleft()
return None
致命缺陷:当 P0 流量突发时,P1-P3 完全饿死。在生产中不可接受,需要更精细的算法。
3.2 赤字轮询(Deficit Round Robin, DRR)
DRR 为每个优先级维护一个"赤字计数器",每次轮转发分配额度的 quantum:
class DRRScheduler:
def __init__(self, num_priorities: int = 4, base_quantum: int = 1):
self.queues = [deque() for _ in range(num_priorities)]
self.deficits = [0] * num_priorities
# 权重随优先级指数增长:P0=8x, P1=4x, P2=2x, P3=1x
self.weights = [8, 4, 2, 1][:num_priorities]
def dequeue(self) -> Optional[Request]:
for i in range(len(self.queues)):
if not self.queues[i]:
continue
# 每轮增加 weight 的配额
self.deficits[i] += self.weights[i]
# 如果配额足够且队列有请求,出队
if self.deficits[i] >= 1 and self.queues[i]:
self.deficits[i] -= 1
return self.queues[i].popleft()
return None
问题:纯 DRR 不感知请求代价。一个 P0 长生成请求(2048 tokens)和一个 P1 短分类请求(10 tokens)配额相同,不公平。
3.3 公平份额 + 优先级加权(Fair Share with Priority Boost)
生产系统通常采用优先级加权公平份额(Weighted Fair Queuing 变体),并引入老化(Aging)机制防饥饿:
import time
from dataclasses import dataclass, field
from typing import Optional
import heapq
@dataclass
class Request:
id: str
priority: int # 0=highest
prompt_tokens: int
max_new_tokens: int
arrival_time: float
virtual_finish_time: float = 0.0
def __lt__(self, other):
"""按虚拟完成时间排序,时间相同则按优先级"""
if self.virtual_finish_time != other.virtual_finish_time:
return self.virtual_finish_time < other.virtual_finish_time
return self.priority < other.priority
class WeightedFairScheduler:
"""
基于 Start-time Fair Queuing 的优先级调度器。
高优先级请求拥有更低的虚拟时间增量,因此更早被调度。
"""
PRIORITY_WEIGHTS = {0: 16.0, 1: 8.0, 2: 4.0, 3: 2.0}
# 老化阈值:低优先级等待超过此时间自动提升一级
AGING_BOOST_THRESHOLD = {3: 5.0, 2: 15.0, 1: 30.0}
def __init__(self):
self.ready_queue: list[Request] = [] # 最小堆
self.active_virtual_time = 0.0
self.queued_requests: dict[int, deque[Request]] = {
0: deque(), 1: deque(), 2: deque(), 3: deque()
}
def enqueue(self, request: Request):
# 检查是否需要老化提升
wait_time = time.time() - request.arrival_time
if request.priority > 0:
threshold = self.AGING_BOOST_THRESHOLD.get(request.priority, float('inf'))
if wait_time > threshold:
request.priority -= 1 # 提升一级
self.queued_requests[request.priority].append(request)
self._promote_to_ready()
def _promote_to_ready(self):
"""将各优先级队列头部的请求提升到就绪堆"""
for priority in range(4):
if self.queued_requests[priority] and len(self.ready_queue) < 8:
req = self.queued_requests[priority].popleft()
weight = self.PRIORITY_WEIGHTS[priority]
# 计算虚拟完成时间
cost = req.prompt_tokens + req.max_new_tokens
# 归一化成本:权重越大 = 虚拟时间增长越慢 = 越快被调度
virtual_cost = cost / weight
req.virtual_finish_time = self.active_virtual_time + virtual_cost
heapq.heappush(self.ready_queue, req)
def dequeue(self) -> Optional[Request]:
if not self.ready_queue:
return None
req = heapq.heappop(self.ready_queue)
# 推进虚拟时间
weight = self.PRIORITY_WEIGHTS[req.priority]
cost = req.prompt_tokens + req.max_new_tokens
self.active_virtual_time += cost / weight
self._promote_to_ready()
return req
3.4 混合策略:Best-Effort 分层
在生产实践中,我们使用混合策略:
class HybridScheduler:
"""
混合调度策略:
- P0:严格隔离 + 预留资源,永不等待
- P1-P2:WFQ 公平份额 + 优先级加权
- P3:Best-Effort,仅在空闲时执行
"""
def dequeue(self, engine_state: 'EngineState') -> Optional[Request]:
# 第一层:P0 严格优先级
if self.p0_queue and engine_state.reserved_p0_available():
return self.p0_queue.popleft()
# 第二层:P1-P2 WFQ 调度
candidate = self.wfq_scheduler.peek()
if candidate and engine_state.can_schedule(candidate):
if self._p0_waiting_for_memory():
self._evict_lowest_p2_for_p0()
return self.wfq_scheduler.dequeue()
# 第三层:P3 Best-Effort(仅无 P0/P1 等待时)
if not self.p0_queue and not self.p1_queue:
if engine_state.idle_slots > 0:
return self.p3_scheduler.dequeue()
return None
四、KV Cache 内存分区与抢占
优先级调度的核心挑战在于:GPU 显存是有限资源,且 LLM 推理的 KV Cache 占用随生成过程动态增长。
4.1 分区策略
┌─────────────────────────────────────────────────────┐
│ GPU VRAM (80GB) │
│ │
│ ┌──────────┐ ┌──────────────────┐ ┌───────────────┐ │
│ │ P0 预留区 │ │ P1 弹性预留区 │ │ P2/P3 共享区 │ │
│ │ 20GB │ │ 30GB │ │ 30GB │ │
│ │ 硬隔离 │ │ 软限制+可借用 │ │ 先到先得 │ │
│ └──────────┘ └──────────────────┘ └───────────────┘ │
│ │ │ │ │
│ │ 空闲时可借用 ◄────────────────────┘ │
│ │ │
│ └──────────────────► 可回收给紧急请求 │
└─────────────────────────────────────────────────────┘
4.2 KV Cache 抢占与再计算
当高优先级请求需要内存但共享区已满时,有两种策略:
策略 A:Swap to CPU(页交换)
将低优先级请求的 KV Cache 换出到 CPU 内存,稍后换入恢复:
class KVCacheManager:
def swap_out_for_priority(self, required_blocks: int, target_priority: int):
"""为高优先级请求腾出 KV Cache 空间"""
candidates = self._find_evictable_requests(min_priority=target_priority + 1)
freed = 0
for victim in candidates:
if freed >= required_blocks:
break
# 将 victim 的 KV Cache 序列化到 CPU 内存
swap_data = victim.kv_cache.to_cpu()
self.cpu_swaps[victim.id] = SwapEntry(
data=swap_data,
original_priority=victim.priority,
prefill_len=victim.current_seq_len
)
# 释放 GPU 显存
freed += victim.kv_cache.num_blocks
victim.kv_cache.free()
victim.state = RequestState.SWAPPED_OUT
def swap_in(self, request_id: str):
"""将换出的请求恢复到 GPU"""
entry = self.cpu_swaps.pop(request_id)
gpu_cache = self.gpu_allocator.allocate(entry.data.num_blocks)
gpu_cache.copy_from_cpu(entry.data)
return gpu_cache
策略 B:Recompute(丢弃重算)
直接丢弃低优先级 KV Cache,后续从头重新计算。适用于换入换出延迟 > 重算延迟的场景:
def evict_with_recompute(self, required_blocks: int) -> int:
"""驱逐低优先级请求,释放显存,后续重算"""
freed = 0
victims = sorted(
self.active_requests,
key=lambda r: (r.priority, -r.generated_tokens)
)
for victim in victims:
if freed >= required_blocks:
break
if victim.priority < 2:
continue
freed += victim.kv_cache.num_blocks
victim.kv_cache.free()
victim.state = SessionState.INTERRUPTED
victim.interrupted_at = victim.generated_tokens
self.recompute_queue.append(victim)
return freed
决策规则:
换出 vs 重算决策:
if swap_bandwidth > 0 and swap_size / swap_bandwidth < prefill_time * 0.3:
使用 Swap(换出延迟 < 重算时间的 30%)
else:
使用 Recompute(对于长上下文短生成就划算)
4.3 预留区的动态调整
静态分区浪费显存。生产系统需要根据实际负载动态调整:
class AdaptiveReservation:
"""
基于预留 + 借用 + 回收 的动态内存管理。
"""
def __init__(self, total_blocks: int, configs: dict):
# 硬下限:P0 绝对保证的最小空间
self.p0_hard_floor = configs.get('p0_hard_floor', int(total_blocks * 0.15))
# 弹性上限:P0 最多可使用的空间
self.p0_soft_ceiling = configs.get('p0_soft_ceiling', int(total_blocks * 0.40))
self.total_blocks = total_blocks
self.p0_used = 0
self.p1_used = 0
self.shared_used = 0
self.p0_borrowed = 0
def allocate_p0(self, num_blocks: int) -> Optional[int]:
"""P0 请求显存:先在预留区分配,不够则借用共享区"""
p0_available = self.p0_soft_ceiling - self.p0_used
if p0_available >= num_blocks:
self.p0_used += num_blocks
return P0_RESERVED_POOL
needed_from_shared = num_blocks - max(p0_available, 0)
shared_available = self.total_blocks - self.p0_used - self.p1_used - self.shared_used
if shared_available >= needed_from_shared:
self.p0_used += num_blocks
self.p0_borrowed += needed_from_shared
return P0_BORROWED_POOL
return None
def reclaim_from_low_priority(self, required_blocks: int) -> int:
"""强制从低优先级回收显存给 P0"""
reclaimed = 0
for req in self.active_p3:
if reclaimed >= required_blocks:
break
reclaimed += req.kv_cache.num_blocks
req.preempt()
for req in self.active_p2:
if reclaimed >= required_blocks:
break
reclaimed += req.kv_cache.num_blocks
req.preempt()
return reclaimed
五、生产级工程实践
5.1 请求元数据与优先级标记
下游请求通过 HTTP Header 或请求体字段声明优先级:
from fastapi import FastAPI, Header, HTTPException
from pydantic import BaseModel
app = FastAPI()
class ChatRequest(BaseModel):
messages: list
max_tokens: int = 512
priority: Optional[int] = None
sla_latency_ms: Optional[int] = None
@app.post("/v1/chat/completions")
async def chat_completion(
body: ChatRequest,
x_request_tier: str = Header("standard"),
x_tenant_id: str = Header(...),
):
priority = resolve_priority(
tier=x_request_tier,
tenant_id=x_tenant_id,
declared_priority=body.priority,
token_count=estimate_tokens(body.messages)
)
request = Request(
priority=priority,
prompt_tokens=token_count,
max_new_tokens=body.max_tokens,
payload=body,
deadline=time.time() + SLA_LATENCY[priority]
)
if not scheduler.admit(request):
raise HTTPException(
status_code=503,
detail={
"error": "resource_exhausted",
"retry_after": scheduler.estimate_wait(priority),
"current_queue_depth": scheduler.queue_depth(priority)
}
)
result = await scheduler.submit(request)
return result
5.2 调度器与推理引擎的集成
以下是一个简化版的调度-执行循环:
class PriorityAwareEngine:
"""集成优先级调度的 LLM 推理引擎"""
def __init__(self, model, gpu_allocator, scheduler):
self.model = model
self.allocator = gpu_allocator
self.scheduler = scheduler
self.running_batch: Optional[Batch] = None
self.max_batch_size = 256
async def run_loop(self):
while True:
# 1. 检查当前 batch 完成状态
self._collect_finished()
# 2. 尝试从调度器获取新请求
current_load = self.running_batch.total_tokens() if self.running_batch else 0
while current_load < self.max_batch_size * 2048:
req = self.scheduler.dequeue(engine_state=self._state_snapshot())
if req is None:
break
# 3. 尝试分配 KV Cache 空间
cache = self.allocator.allocate(
estimated_blocks=req.prompt_tokens // 16,
priority=req.priority
)
if cache is None:
if req.priority == 0:
freed = self.allocator.reclaim_from_low_priority(
req.prompt_tokens // 16
)
cache = self.allocator.allocate(req.prompt_tokens // 16, priority=0)
else:
self.scheduler.requeue(req)
break
self.running_batch.add(req, cache)
current_load += req.prompt_tokens
# 4. 执行一步解码
if self.running_batch and len(self.running_batch) > 0:
await self.model.step(self.running_batch)
await asyncio.sleep(0)
else:
await asyncio.sleep(0.005)
def _state_snapshot(self) -> EngineState:
return EngineState(
pending_p0=self.scheduler.queue_depth(0),
pending_p1=self.scheduler.queue_depth(1),
available_blocks=self.allocator.available_blocks(),
p0_reserved_free=self.allocator.p0_reserved_free(),
current_batch_tokens=self.running_batch.total_tokens() if self.running_batch else 0
)
5.3 可观测性:调度指标
# Prometheus 指标定义
SCHEDULER_METRICS = {
'llm_scheduler_queue_depth': Gauge(
'Scheduler queue depth by priority',
['priority']
),
'llm_scheduler_wait_seconds': Histogram(
'Time from enqueue to first token',
['priority'],
buckets=[0.01, 0.05, 0.1, 0.5, 1.0, 5.0, 10.0, 30.0, 60.0]
),
'llm_kv_cache_hit_rate': Gauge(
'KV cache prefix hit rate',
['priority']
),
'llm_preemption_total': Counter(
'Total preempted requests',
['victim_priority', 'reason']
),
'llm_sla_attainment': Gauge(
'Percentage of requests meeting SLA',
['priority', 'tier']
),
'llm_kv_swap_seconds': Histogram(
'Time to swap KV cache between GPU and CPU',
['direction']
),
}
5.4 关键工程陷阱
陷阱 1:优先级反转(Priority Inversion)
问题:P3 请求已通过准入控制占用了显存,此时 P0 请求到达需要抢占,但 P3 正在 GPU 上等待计算完成。
if p0_request_arrived and gpu_memory_pressure > threshold:
# 标记 victim 为 preempted,完成当前 step 后即被调出
victim.mark_preempt_after_step()
陷阱 2:内存碎片导致高优先级请求分配失败
def compact_kv_cache(self, priority_floor: int = 2):
"""整理低优先级请求的 KV Cache,减少碎片"""
for req in self.active_requests:
if req.priority >= priority_floor and req.kv_cache.is_fragmented():
req.kv_cache.compact()
陷阱 3:Heavy Tail 长请求阻塞短请求
# 当 P1 请求超过一定长度时,强制切分或限制 batch 位置
if req.priority >= 1 and req.prompt_tokens > 8192:
req.batch_token_budget = 4096
六、性能基准
我们基于 vLLM 0.6 + 自建调度插件,在 A100 80GB × 8 集群上进行了基准测试:
配置:Llama-3.1-70B, PagedAttention block_size=16, 混合负载
延迟对比(P95 TTFT,秒)
| 并发请求数 | 无优先级 | DRR 策略 | 本文策略 | 提升 |
|---|---|---|---|---|
| 32 | 0.8 | 0.9 | 0.8 | — |
| 64 | 2.1 | 1.5 | 1.2 | P0 -43% |
| 128 | 8.5 | 3.2 | 1.8 | P0 -79% |
| 256 | 25+ | 8.7 | 3.5 | P0 -86% |
| 512 | OOM | 18.3 | 6.2 | P0 - |
GPU 利用率对比
| 并发请求数 | 无优先级 | DRR 策略 | 本文策略 | 说明 |
|---|---|---|---|---|
| 32 | 72% | 70% | 71% | 利用率持平 |
| 128 | 85% | 82% | 84% | 仅 -1% 损耗 |
| 256 | 78% | 88% | 91% | +3%(减少OOM) |
| 512 | OOM | 75% | 86% | +11% |
关键洞察:
- 在低负载场景下(<64 并发),优先级调度几乎无开销
- 在高压场景下(128+ 并发),高优先级延迟降低 79-86%,GPU 利用率仅降低 1-3%
- Swap vs Recompute 的决策优化可以减少约 15% 的抢占开销
七、未来方向
- 预测性调度:基于请求历史模式预测 KV Cache 占用,提前预留资源
- 分布式优先级感知:跨多节点全局协调优先级队列(类似 Google Borg 的跨机调度)
- SLA-aware Auto Scaling:优先级维度的自动扩缩容,P0 队列增长时优先扩容
- Speculative Decoding 与调度的联合优化:为高优先级请求启用草稿模型,进一步降低 TTFT
总结
LLM 推理服务的优先级调度是一个跨越排队论、操作系统内存管理和 GPU 编程的多层系统工程。生产实践中的关键 lessons learned:
- 调度必须与内存管理协同设计——纯队列调度无法解决 KV Cache 抢占问题
- 混合策略优于单一算法——P0 严格隔离 + P1-P2 WFQ + P3 Best-Effort 在实践中表现最佳
- 可观测性是优先级系统的生命线——没有精确的 queue depth、wait time、preemption count 指标,无法证明调度策略有效
- Preemption 需要与引擎深度集成——粗粒度驱逐会增加延迟抖动,细粒度协作是未来的方向
下一代推理引擎的调度系统将不再是简单的"谁先执行",而是一个多目标优化问题:在满足 SLA 的同时,最大化吞吐、最小化能耗、保证公平性。
本文基于 vLLM 0.6.x 和自制调度插件的生产实践总结,代码示例均为简化版本,仅用于说明算法逻辑。

发表评论 取消回复