Scheduler 初始化与三队列
Scheduler 是 vLLM 的核心调度组件,负责在每个推理步骤(step)中决定哪些请求进入 GPU 执行。
它维护三个 Deque[SequenceGroup] 队列,构成整个调度系统的骨架。
class Scheduler:
def __init__(
self,
scheduler_config: SchedulerConfig,
cache_config: CacheConfig,
lora_config: Optional[LoRAConfig],
) -> None:
...
# Sequence groups in the WAITING state.
# Contain new prefill or preempted requests.
self.waiting: Deque[SequenceGroup] = deque()
# Sequence groups in the RUNNING state.
# Contain decode requests.
self.running: Deque[SequenceGroup] = deque()
# Sequence groups in the SWAPPED state.
# Contain decode requests that are swapped out.
self.swapped: Deque[SequenceGroup] = deque()
三个队列的角色与语义
这三个队列分别对应请求生命周期中的三种状态,每个队列的元素都是 SequenceGroup(一个请求对应一个 SequenceGroup):
存储所有尚未开始执行的新请求,以及被抢占后需要从头重新计算(Recompute 模式)的请求。
新请求通过 add_seq_group() 进入此队列:
def add_seq_group(self, seq_group: SequenceGroup) -> None:
# Add sequence groups to the waiting queue.
self.waiting.append(seq_group)
使用 deque 而非 list 的原因:deque 的两端插入删除均为 O(1),
在调度时可以高效地从左端弹出(popleft)最高优先级请求,同时从右端追加新请求。
存储当前批次(batch)中正在 GPU 上执行的请求。这些请求已完成至少一次 Prefill, 正处于 Decode 阶段(或 Chunked Prefill 模式下的分块 Prefill 阶段)。 每个调度周期结束后,已选中执行的请求被移入此队列。
注释特别说明"Contain decode requests",即该队列中的请求均处于逐 token 生成阶段, 每次迭代每个请求只消耗 1 个新 token 的预算(1 个 decode token)。
存储因 GPU 显存不足而被换出(Swap Out)到 CPU 内存的请求。 注释"Contain decode requests that are swapped out"表明, 只有正在 Decode 的请求才会被换出(Prefill 阶段请求通过 Recompute 抢占,不经过 Swap)。
当 GPU 显存恢复充裕时,swapped 队列中的请求会被换回(Swap In)至 running 队列继续执行。
KV Cache 块通过 block_manager 在 CPU ↔ GPU 之间搬运。
为什么是三个队列,而不是一个状态列表?
这是一个关键的设计决策。将不同状态的请求分离到不同队列有几个重要优势:
-
调度逻辑清晰:每个队列对应独立的调度子程序
(
_schedule_running()、_schedule_swapped()、_schedule_prefills()), 职责分离,便于维护。 - 优先级明确:三队列天然实现了优先级顺序:running > swapped > waiting。 先保障正在运行的请求完成 decode,再考虑换入已换出的请求,最后才处理新请求。
-
性能高效:
deque的 O(1) 两端操作远优于在单一列表中按状态筛选(O(n))。
初始化时的其他关键组件
if self.scheduler_config.chunked_prefill_enabled:
self.prompt_limit = self.scheduler_config.max_model_len
else:
self.prompt_limit = min(
self.scheduler_config.max_model_len,
self.scheduler_config.max_num_batched_tokens)
BlockSpaceManagerImpl = BlockSpaceManager.get_block_space_manager_class(
version="v2" if self.scheduler_config.use_v2_block_manager else "v1")
self.block_manager = BlockSpaceManagerImpl(
block_size=self.cache_config.block_size,
num_gpu_blocks=self.cache_config.num_gpu_blocks,
num_cpu_blocks=self.cache_config.num_cpu_blocks,
sliding_window=self.cache_config.sliding_window,
enable_caching=self.cache_config.enable_prefix_caching,
...)
prompt_limit 的双路计算逻辑值得关注:在 Chunked Prefill 模式下,prompt 可以被分块处理,
因此上限放宽到 max_model_len;而在普通模式下,单次批处理的 token 数受
max_num_batched_tokens 约束,因此取两者的最小值。
这一设计防止超长 prompt 在一次 step 中耗尽所有 token 预算,导致其他请求饥饿。
SchedulingBudget 预算机制
SchedulingBudget 是调度器在每个 step 开始时创建的临时对象,用于追踪当前批次已消耗的资源,
并以此判断新请求是否还能被纳入本次批次。它是防止单批次过载的核心守门人。
@dataclass
class SchedulingBudget:
"""The available slots for scheduling."""
token_budget: int # 本批次最多处理多少 token
max_num_seqs: int # 本批次最多同时运行多少个序列
_requeset_ids_num_batched_tokens: Set[str] = field(default_factory=set)
_requeset_ids_num_curr_seqs: Set[str] = field(default_factory=set)
_num_batched_tokens: int = 0
_num_curr_seqs: int = 0
def can_schedule(self, *, num_new_tokens: int, num_new_seqs: int):
assert num_new_tokens != 0
assert num_new_seqs != 0
return (self.num_batched_tokens + num_new_tokens <= self.token_budget
and self.num_curr_seqs + num_new_seqs <= self.max_num_seqs)
def remaining_token_budget(self):
return self.token_budget - self.num_batched_tokens
两个维度的预算约束
SchedulingBudget 同时约束两个资源维度,任何一个超出都会阻止新请求入队:
对应配置参数 max_num_batched_tokens,控制单次 forward pass 最多处理多少个 token。
这直接决定 GPU 的计算负载(矩阵乘法的批大小)。
- Decode 请求:每个 sequence 贡献 1 个 token(只生成下一个 token)
- Prefill 请求:每个请求贡献 prompt 长度个 token(一次性处理整个 prompt)
- Chunked Prefill:每次最多贡献 chunk_size 个 token(分块处理)
对应配置参数 max_num_seqs,控制同时在 GPU 上活跃的序列数量。
这直接决定 KV Cache 的占用量(每个序列需要维护自己的 KV Cache)。
注意:num_curr_seqs 计的是序列(Sequence)数,而不是请求数。
一个 SequenceGroup(对应 Beam Search 等场景)可以包含多个 Sequence,
所以序列数 ≥ 请求数。
Request-ID 去重机制
SchedulingBudget 的一个微妙设计是基于 request_id 的去重:
def add_num_batched_tokens(self, req_id: str, num_batched_tokens: int):
if req_id in self._requeset_ids_num_batched_tokens:
return # 同一请求不重复计数
self._requeset_ids_num_batched_tokens.add(req_id)
self._num_batched_tokens += num_batched_tokens
def subtract_num_batched_tokens(self, req_id: str, num_batched_tokens: int):
if req_id in self._requeset_ids_num_batched_tokens:
self._requeset_ids_num_batched_tokens.remove(req_id)
self._num_batched_tokens -= num_batched_tokens
def add_num_seqs(self, req_id: str, num_curr_seqs: int):
if req_id in self._requeset_ids_num_curr_seqs:
return # 同一请求不重复计数
...
代码注释解释了为什么需要这个去重机制:
在 Normal Scheduling 路径中,调度器对 RUNNING 队列中的请求提前更新了
num_seqs 预算(在正式调度前就计入),这导致同一个请求可能被计数两次。
通过维护已记录 request_id 的 Set,确保每个请求只贡献一次到预算中。
注释还提到:如果未来 Chunked Prefill 成为默认模式,这个问题不会出现,届时可以移除此机制。 这是一个向前兼容的临时设计,保证新旧调度路径都能正确工作。
预算的生命周期
SchedulingBudget 在每次 _schedule() 调用时创建,调度完成后随即销毁。
它是一个局部状态对象,不在 Scheduler 实例中持久化。
这确保了每个 step 的调度决策相互独立,不会受历史预算消耗的影响。
# 在 _schedule_default() 或 _schedule_chunked_prefill() 开始时
budget = SchedulingBudget(
token_budget=self.scheduler_config.max_num_batched_tokens,
max_num_seqs=self.scheduler_config.max_num_seqs,
)
Policy 调度策略(FCFS)
policy.py 定义了请求的排序策略。当前 vLLM 只实现了一种策略:
FCFS(First Come First Served,先来先服务),但架构设计预留了扩展接口。
class Policy:
def get_priority(self, now: float, seq_group: SequenceGroup) -> float:
raise NotImplementedError
def sort_by_priority(
self,
now: float,
seq_groups: Deque[SequenceGroup],
) -> Deque[SequenceGroup]:
return deque(
sorted(
seq_groups,
key=lambda seq_group: self.get_priority(now, seq_group),
reverse=True, # 优先级值越大越先调度
))
class FCFS(Policy):
def get_priority(self, now: float, seq_group: SequenceGroup) -> float:
return now - seq_group.metrics.arrival_time
# 等待时间越长 → 优先级越高(值越大)
class PolicyFactory:
_POLICY_REGISTRY = {'fcfs': FCFS}
@classmethod
def get_policy(cls, policy_name: str, **kwargs) -> Policy:
return cls._POLICY_REGISTRY[policy_name](**kwargs)
FCFS 优先级计算的数学含义
FCFS 的优先级公式非常简洁:
等待时间越长,优先级越高。这保证了请求按到达顺序被服务, 不会因为后来的短请求插队而造成长请求永久等待(饥饿)。
sort_by_priority 使用 reverse=True,即优先级值最大(等待最久)的请求排在队首。
这与直觉一致:先到先得。
为什么只有 FCFS?
代码注释坦承当前 LoRA 调度策略"extremely simple and NOT fair",存在饥饿风险, 未来需要改进。当前只有 FCFS 的原因是:
- 实现简单,行为可预测:FCFS 无需额外配置,用户容易理解服务顺序。
- 对 LLM 推理足够公平:LLM 请求通常没有严格 SLA 要求,FCFS 是合理的默认策略。
-
架构预留扩展点:
PolicyFactory._POLICY_REGISTRY是一个注册表, 未来可以添加 Shortest Job First(SJF)、优先级队列等策略, 只需继承Policy并注册即可,不需要修改调度器主逻辑。
Policy 在调度中的实际使用
Policy 在调度时对 waiting 队列进行排序,确保每次都处理等待最久的请求。 运行中(running/swapped)的请求不需要重新排序,因为它们已经在执行中, 调度器的首要任务是保证它们能继续运行(维持低延迟)。
policy = PolicyFactory.get_policy(policy_name="fcfs")
# waiting 队列按到达时间排序,优先处理等待最久的请求
waiting_queue = policy.sort_by_priority(now=time.time(),
seq_groups=self.waiting)
_schedule() 方法总览
_schedule() 是调度器的顶层入口,每个推理 step 调用一次,
负责选择本次 forward pass 要处理的请求集合,并返回调度决策。
它本身不包含调度逻辑,而是作为策略分发器,根据配置选择具体的调度路径。
def _schedule(self) -> SchedulerOutputs:
"""Schedule sequence groups.
The current policy is designed to optimize the throughput. First,
it batches as many prefill requests as possible. And it schedules
decodes. If there's a pressure on GPU memory, preemption or
swapping is performed.
"""
if self.scheduler_config.chunked_prefill_enabled:
return self._schedule_chunked_prefill()
else:
return self._schedule_default()
两条调度路径的对比
调用 _schedule_default(),执行三阶段调度。
注意:代码的 docstring 明确写道 "First, it batches as many prefill requests as possible.
And it schedules decodes."——Prefill 优先于 Decode。
- Phase 1 — Waiting 优先(Prefill):如果 swapped 队列为空, 先从 waiting 队列取新请求做 Prefill。多个 prefill 可以 batch 在一起,循环取出 waiting 请求直到 budget 用完。 如果 swapped 队列非空则跳过此阶段。
- Phase 2 — Running 次之(Decode):如果 Phase 1 没有调度任何 prefill, 才为 running 队列中的请求分配 decode 所需的 KV Cache 块。 如果显存不足,从队尾(最近才升入 running 的请求)开始抢占。
- Phase 3 — Swapped 最后(Swap In):如果 Phase 2 没有发生抢占, 尝试将 swapped 请求换回 GPU。
互斥机制:做了 prefill 就跳过 decode 和 swap_in;有 swapped 就跳过 prefill。 这保证了同一批次中 Prefill 和 Decode 不会混合。
多个 prefill 能一起 batch 吗? Q&A
能。_schedule_prefills() 内部是一个 while 循环,
不断从 waiting 队列取请求,直到 budget 用完或块不够:
while waiting_queue:
seq_group = waiting_queue[0]
if block_manager.can_allocate(seq_group) == LATER:
break
if not budget.can_schedule(num_new_tokens, num_new_seqs):
break
# 够!调度这个 prefill
seq_groups.append(seq_group)
budget.add_num_batched_tokens(num_new_tokens)
例如 max_num_batched_tokens=4096,waiting 队列有三个请求(1000、800、1500 tokens),
它们会全部被塞进同一个 batch(共 3300 tokens < 4096)。
但这一步仍然没有 decode——prefill 和 decode 不混合。
"Prefill 和 Decode 不混合"具体怎么保证的? Q&A
通过两个 if 判断实现互斥:
# 有 swapped → 跳过 prefill
if not self.swapped:
prefills = self._schedule_prefills(...)
# 做了 prefill → 跳过 decode
if len(prefills.seq_groups) == 0:
running_scheduled = self._schedule_running(...)
四种场景分析:
- 有 swapped:不做 prefill → 做 decode → 尝试 swap_in → batch 里只有 decode
- 没 swapped,有 waiting,够 prefill:做 prefill → 跳过 decode → batch 里只有 prefill
- 没 swapped,有 waiting,不够 prefill:prefill 失败 → 做 decode → batch 里只有 decode
- 没 swapped,没 waiting:无 prefill → 做 decode → batch 里只有 decode
每个 step 要么全是 prefill,要么全是 decode,不会两种同时出现。 代价:prefill 那一步,正在 decode 的请求被饿死,用户体感输出卡了一下。
调用 _schedule_chunked_prefill(),允许 Prefill 和 Decode 在同一批次中共存。
长 prompt 被切分为多个小 chunk,每个 step 只处理一个 chunk,
剩余的 token 预算留给 decode 请求。
优势:即使有超长 prompt 请求,decode 请求的延迟也不会被大幅拉升, GPU 利用率更高。代价是实现复杂度更高,需要精确追踪每个请求的 chunk 进度。
Chunked Prefill 的调度优先级和 Default 有什么区别? Q&A
优先级完全反转——Default 是 prefill 优先,Chunked 是 decode 优先:
# Default 路径:
① _schedule_prefills (waiting) ← prefill 优先
② _schedule_running (decode)
③ _schedule_swapped (swap_in)
# Chunked Prefill 路径:
① _schedule_running (decode) ← decode 优先
② _schedule_swapped (swap_in)
③ _schedule_prefills (waiting) ← prefill 用剩余 budget
具体例子(max_num_batched_tokens=2048,3 个 decode 请求 + 1 个 5000 token 新请求):
- Default:seq_D 的 prefill(比如 2000 tokens)独占整步 batch, seq_A/B/C 这步不 decode,用户看到输出卡了一下。
- Chunked:先调度 seq_A/B/C 的 decode(3 tokens), budget 剩 2045 → seq_D 取 2045 个 token 做 chunk1。 下一步再取 chunk2(2045 tokens),第三步取剩余 910 tokens。 decode 请求每步都在跑,不会被饿死。
公开入口:schedule() 方法
外部调用者(LLMEngine.step())调用的是公开的 schedule() 方法,
它在 _schedule() 基础上添加了 NVTX 性能标注和日志记录:
@nvtx_tag(message="schedule", domain="vllm")
def schedule(self) -> Tuple[List[SequenceGroupMetadata], SchedulerOutputs]:
# Schedule sequence groups.
# NOTE: It is expected that the caller holds at least 1 GPU block
# before calling this method. Else, the system will deadlock.
scheduler_outputs = self._schedule()
now = time.time()
# Create input data structures.
seq_group_metadata_list: List[SequenceGroupMetadata] = []
for i, scheduled_seq_group in enumerate(
scheduler_outputs.scheduled_seq_groups):
seq_group = scheduled_seq_group.seq_group
token_chunk_size = scheduled_seq_group.token_chunk_size
seq_group.maybe_set_first_scheduled_time(now)
# seq_id -> SequenceData
seq_data: Dict[int, SequenceData] = {}
slot_mapping: Dict[int, List[int]] = {}
...
seq_group_metadata_list.append(
SequenceGroupMetadata(
request_id=seq_group.request_id,
is_prompt=seq_group.is_prefill(),
seq_data=seq_data,
sampling_params=seq_group.sampling_params,
block_tables=block_tables,
token_chunk_size=token_chunk_size,
...
))
return seq_group_metadata_list, scheduler_outputs
这里有个重要的注释:"It is expected that the caller holds at least 1 GPU block"。
这是一个前置条件契约——如果 GPU 完全没有可用块,
调度器会陷入死锁(所有请求都无法被调度,但没有请求完成来释放块)。
实际上 LLMEngine 在初始化时会验证 GPU 块数量的合法性。
SchedulerOutputs 数据结构
SchedulerOutputs 是调度决策的最终产物,它将调度器的决策"打包"后传递给执行层(Worker)。
理解其字段是理解 vLLM 执行流水线的关键。
@dataclass
class SchedulerOutputs:
"""The scheduling decision made from a scheduler."""
# Scheduled sequence groups.
scheduled_seq_groups: Iterable[ScheduledSequenceGroup]
# Number of prefill groups scheduled.
num_prefill_groups: int
# Total number of batched tokens.
num_batched_tokens: int
# Blocks to swap in. Dict of CPU -> GPU block number.
blocks_to_swap_in: Dict[int, int]
# Blocks to swap out. Dict of GPU -> CPU block number.
blocks_to_swap_out: Dict[int, int]
# Blocks to copy. Source to a list of dest blocks.
blocks_to_copy: Dict[int, List[int]]
# Sequence groups that are going to be ignored.
ignored_seq_groups: List[SequenceGroup]
# The number of slots for lookahead decoding.
num_lookahead_slots: int
# The number of requests in the running queue
running_queue_size: int
字段逐一解析
类型为 Iterable[ScheduledSequenceGroup],其中 ScheduledSequenceGroup 是:
@dataclass
class ScheduledSequenceGroup:
seq_group: SequenceGroup # 请求本体
token_chunk_size: int # 本次处理多少个 token
# Decode: token_chunk_size = 1
# Prefill: token_chunk_size = prompt 长度(或 chunk 大小)
token_chunk_size 是 Chunked Prefill 的核心字段:
它告诉 Worker 这个请求本次处理多少个 token,允许调度器精细控制每次迭代的计算量。
这两个字典分别指定本次 step 需要执行的 CPU↔GPU 内存搬运操作:
blocks_to_swap_in:{cpu_block_id: gpu_block_id},将 CPU 块复制到 GPU 指定位置blocks_to_swap_out:{gpu_block_id: cpu_block_id},将 GPU 块备份到 CPU
__post_init__ 中有一个重要断言:
assert not (self.blocks_to_swap_in and self.blocks_to_swap_out)。
Swap In 和 Swap Out 不可以在同一个 step 同时发生。
这避免了数据混乱,也简化了 GPU 端的内存搬运实现。
类型为 Dict[int, List[int]],值是目标块列表而非单个块。
当多个 Sequence(如 Beam Search 中)共享同一个 KV Cache 块,
而其中一个需要写入新 token 时,就触发 Copy-on-Write:
将源块复制到一个新块,再写入新数据,保证其他 Sequence 不受影响。
当一个请求的 prompt 超过 prompt_limit,或者即使独占整个 GPU 也无法分配足够的 KV Cache 块,
该请求就会被放入 ignored_seq_groups,直接返回给上层并标记为错误。
这防止了系统因为一个超大请求而永久阻塞。
SchedulerOutputs 与执行层的接口
schedule() 方法返回 (seq_group_metadata_list, scheduler_outputs):
-
seq_group_metadata_list:每个调度请求的完整元数据(token IDs、block tables、采样参数等), 直接传入 Worker 的execute_model(),作为模型 forward pass 的输入。 -
scheduler_outputs:调度决策的摘要,用于执行前的 KV Cache 搬运(swap in/out/copy), 以及执行后的结果处理(更新队列状态、释放完成请求的资源)。
中间调度数据结构
除了 SchedulerOutputs,调度器内部还定义了三个辅助数据类,
分别对应三个队列各自的调度子结果:
@dataclass
class SchedulerRunningOutputs:
"""The requests scheduled from the running queue."""
decode_seq_groups: List[SequenceGroup] # 正常 decode 的请求
prefill_seq_groups: List[SequenceGroup] # Chunked Prefill 中分块执行的请求
preempted: List[SequenceGroup] # 被 Recompute 抢占的请求
swapped_out: List[SequenceGroup] # 被 Swap Out 的请求
blocks_to_swap_out: Dict[int, int]
blocks_to_copy: Dict[int, List[int]]
num_lookahead_slots: int
@dataclass
class SchedulerSwappedInOutputs:
"""The requests scheduled from the swapped queue."""
decode_seq_groups: List[SequenceGroup]
prefill_seq_groups: List[SequenceGroup]
blocks_to_swap_in: Dict[int, int] # 需要换入的块
blocks_to_copy: Dict[int, List[int]]
num_lookahead_slots: int
infeasible_seq_groups: List[SequenceGroup] # 无法换入的请求(CPU 块已损坏等)
@dataclass
class SchedulerPrefillOutputs:
"""The requests scheduled from the waiting queue."""
seq_groups: List[SequenceGroup] # 被选中做 Prefill 的新请求
ignored_seq_groups: List[SequenceGroup] # prompt 过长被忽略的请求
num_lookahead_slots: int
这三个数据类的设计使得各子调度函数(_schedule_running() 等)
可以返回清晰定义的结果,最后在顶层汇总合并为一个 SchedulerOutputs。
这是分治 + 聚合的设计模式,各子函数职责单一,便于独立测试。
队列流转图
下面两张图分别展示:请求在三队列之间的状态流转,以及每次调度 step 内部的执行顺序。
请求生命周期:三队列状态机
新请求到达 waiting --> running: _schedule_prefills()
Prefill 完成,分配 GPU KV Cache running --> running: _schedule_running()
Decode(正常情况) running --> waiting: _schedule_running()
内存不足 → Recompute 抢占
(丢弃 KV Cache,重新排队) running --> swapped: _schedule_running()
内存不足 → Swap Out 抢占
(KV Cache 移至 CPU) swapped --> running: _schedule_swapped()
GPU 内存充裕 → Swap In
(KV Cache 移回 GPU) running --> [*]: 序列生成完毕
free_seq_group() waiting --> [*]: prompt 过长 → ignored
直接返回错误
单次 Step 的调度执行顺序(Default 路径)
设计重点:为什么有 swapped 请求时跳过 Prefill?
在 Default 调度路径中,当 swapped 队列非空时,不会从 waiting 队列取新请求做 Prefill。 这个约束来自两个层面:
- 显存压力:swapped 队列非空说明 GPU 显存正处于紧张状态(否则那些请求不会被换出)。 此时再引入新的 Prefill 请求,会进一步加剧显存压力,可能导致刚换回的请求又被换出, 造成"换进换出震荡"。
- 公平性:被换出的请求(swapped)已经在系统中等待更久, 理应比新到达的 waiting 请求有更高优先级恢复执行。
Chunked Prefill 路径打破了这一限制,允许 Prefill chunk 和 Decode(包括换入请求)共存, 通过精细的 token 预算管理来控制显存占用,是对 Default 路径的重要优化。
抢占策略选择:Swap vs Recompute
将请求的 KV Cache 块从 GPU 显存完整地复制到 CPU 内存,保留所有已计算的中间结果。 优点是换回时无需重新计算(低延迟恢复),缺点是 PCIe 带宽成为瓶颈,且需要足够的 CPU 内存。
直接丢弃请求的 KV Cache 块(释放 GPU 显存),将请求重新放回 waiting 队列, 下次调度时从 Prefill 阶段重新开始。优点是不占用 CPU 内存,也不消耗 PCIe 带宽; 缺点是已完成的 Decode token 对应的 prompt KV Cache 计算白费,产生重复计算开销。
PreemptionMode 枚举定义了这两种模式,具体选择由
scheduler_config.preemption_mode 和 CPU 内存可用情况共同决定。
class PreemptionMode(enum.Enum):
"""Preemption modes.
1. Swapping: Swap out the blocks of the preempted sequences to CPU memory
and swap them back in when the sequences are resumed.
2. Recomputation: Discard the blocks of the preempted sequences and
recompute them when the sequences are resumed, treating the sequences as
new prompts.
"""
SWAP = enum.auto()
RECOMPUTE = enum.auto()