大语言模型推理服务的多级优先级调度 — 从请求路由到内存预留的生产级实现

当 LLM 推理集群同时为实时对话、在线 API 任务和离线批量推理服务时,如何确保高优先级请求不被饿死,同时最大化 GPU 利用率?本文从生产视角拆解 LLM 推理服务的多级优先级调度架构。

一、问题定义:为什么 LLM 推理需要优先级调度?

传统 Web 服务的请求调度相对简单——无状态、短期执行、资源可预测。但 LLM 推理服务有三个本质不同:

  1. 资源占用不对等:一条生成长度 2048 token 的请求可能独占 GPU KV Cache 数分钟,而一条分类任务只需 50ms
  2. 抢占成本极高:中断正在生成的 KV Cache 意味着已计算的中间状态作废,显存无法即时回收
  3. 批处理延迟敏感:Continuous Batching 下新请求必须等当前 batch 完成才能加入,排队延迟直接叠加

生产级 LLM 推理服务的典型场景:

优先级延迟要求典型占比场景
P0 实时<500ms5%多轮对话/工具调用
P1 在线<5s60%API 聊天/摘要
P2 批量<60s30%数据处理/分类
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 策略本文策略提升
320.80.90.8—
642.11.51.2P0 -43%
1288.53.21.8P0 -79%
25625+8.73.5P0 -86%
512OOM18.36.2P0 -

GPU 利用率对比

并发请求数无优先级DRR 策略本文策略说明
3272%70%71%利用率持平
12885%82%84%仅 -1% 损耗
25678%88%91%+3%(减少OOM)
512OOM75%86%+11%

关键洞察:

  • 在低负载场景下(<64 并发),优先级调度几乎无开销
  • 在高压场景下(128+ 并发),高优先级延迟降低 79-86%,GPU 利用率仅降低 1-3%
  • Swap vs Recompute 的决策优化可以减少约 15% 的抢占开销

七、未来方向

  1. 预测性调度:基于请求历史模式预测 KV Cache 占用,提前预留资源
  2. 分布式优先级感知:跨多节点全局协调优先级队列(类似 Google Borg 的跨机调度)
  3. SLA-aware Auto Scaling:优先级维度的自动扩缩容,P0 队列增长时优先扩容
  4. Speculative Decoding 与调度的联合优化:为高优先级请求启用草稿模型,进一步降低 TTFT

总结

LLM 推理服务的优先级调度是一个跨越排队论、操作系统内存管理和 GPU 编程的多层系统工程。生产实践中的关键 lessons learned:

  1. 调度必须与内存管理协同设计——纯队列调度无法解决 KV Cache 抢占问题
  2. 混合策略优于单一算法——P0 严格隔离 + P1-P2 WFQ + P3 Best-Effort 在实践中表现最佳
  3. 可观测性是优先级系统的生命线——没有精确的 queue depth、wait time、preemption count 指标,无法证明调度策略有效
  4. Preemption 需要与引擎深度集成——粗粒度驱逐会增加延迟抖动,细粒度协作是未来的方向

下一代推理引擎的调度系统将不再是简单的"谁先执行",而是一个多目标优化问题:在满足 SLA 的同时,最大化吞吐、最小化能耗、保证公平性。


本文基于 vLLM 0.6.x 和自制调度插件的生产实践总结,代码示例均为简化版本,仅用于说明算法逻辑。

点赞(0) 打赏

评论列表 共有 0 条评论

暂无评论
立即
投稿

微信公众账号

微信扫一扫加关注

发表
评论
返回
顶部