vLLM调度器源码实战:手把手教你用Python模拟三种队列(waiting/running/swapped)的流转
·
vLLM调度器源码实战:Python模拟三种队列状态流转
在大型语言模型推理系统中,高效的调度算法是保证吞吐量的关键。vLLM作为开源推理引擎,其核心调度器通过waiting/running/swapped三个队列的动态流转,实现了对计算资源的优化利用。本文将带您用Python构建简化版调度器,深入理解这三种队列的状态迁移机制。
1. 调度系统基础架构
现代LLM推理系统面临的核心挑战是如何在有限GPU内存下最大化吞吐量。vLLM采用的解决方案是将请求分为三种状态:
- Waiting队列:新到达的请求,等待初始处理
- Running队列:正在执行推理的请求
- Swapped队列:因资源不足被暂时换出的请求
class SequenceGroup:
def __init__(self, request_id: int, prompt: str):
self.request_id = request_id
self.prompt = prompt
self.status = "WAITING" # WAITING/RUNNING/SWAPPED
self.blocks = [] # 占用的内存块
self.output_tokens = [] # 已生成的token
调度器的核心是一个循环处理流程,每次推理前都会执行以下操作:
- 检查swapped队列是否有待恢复的请求
- 从waiting队列选取合适的新请求
- 评估running队列请求的资源需求
- 根据资源情况决定队列间的请求迁移
2. 队列调度核心逻辑实现
2.1 预算管理机制
调度器通过Budget类跟踪当前资源使用情况,防止过载:
class SchedulingBudget:
def __init__(self, token_budget: int, seq_budget: int):
self.token_budget = token_budget
self.seq_budget = seq_budget
self.used_tokens = 0
self.used_seqs = 0
def can_schedule(self, new_tokens: int, new_seqs: int) -> bool:
return (self.used_tokens + new_tokens <= self.token_budget and
self.used_seqs + new_seqs <= self.seq_budget)
2.2 三阶段调度方法
Waiting队列处理
新请求首先进入waiting队列,调度时需检查:
- Prompt长度是否超限
- 是否有足够内存块
- 是否超出当前预算
def schedule_prefills(self, budget: SchedulingBudget):
scheduled = []
while self.waiting:
seq_group = self.waiting[0]
# 检查prompt长度
if len(seq_group.prompt) > self.max_prompt_len:
seq_group.status = "FINISHED_IGNORED"
self.waiting.popleft()
continue
# 检查内存块
if not self.can_allocate(seq_group):
break
# 检查预算
new_tokens = len(seq_group.prompt)
new_seqs = 1 # 初始只有1个序列
if not budget.can_schedule(new_tokens, new_seqs):
break
# 转移至running队列
self.waiting.popleft()
seq_group.status = "RUNNING"
scheduled.append(seq_group)
budget.used_tokens += new_tokens
budget.used_seqs += new_seqs
return scheduled
Running队列维护
每次推理后需要重新评估running队列中的请求:
def schedule_running(self, budget: SchedulingBudget):
preempted = []
swapped = []
for seq_group in list(self.running):
# 每个序列生成1个token
if not budget.can_schedule(1, 1):
# 资源不足时的抢占处理
if self.swap_strategy == "RECOMPUTE":
seq_group.status = "WAITING"
preempted.append(seq_group)
else:
seq_group.status = "SWAPPED"
swapped.append(seq_group)
self.running.remove(seq_group)
else:
budget.used_tokens += 1
seq_group.output_tokens.append("[new_token]")
return preempted, swapped
Swapped队列恢复
当资源释放后,应优先处理swapped队列中的请求:
def schedule_swapped(self, budget: SchedulingBudget):
restored = []
while self.swapped:
seq_group = self.swapped[0]
# 检查恢复条件
new_tokens = 1 # 每个序列生成1个token
new_seqs = len(seq_group.output_tokens) + 1
if not budget.can_schedule(new_tokens, new_seqs):
break
# 恢复执行
self.swapped.popleft()
seq_group.status = "RUNNING"
restored.append(seq_group)
budget.used_tokens += new_tokens
budget.used_seqs += new_seqs
return restored
3. 完整调度流程模拟
将各组件整合成完整调度循环:
class Scheduler:
def __init__(self):
self.waiting = deque()
self.running = []
self.swapped = deque()
self.max_prompt_len = 2048
self.swap_strategy = "SWAP" # or "RECOMPUTE"
def run_scheduling_cycle(self):
# 初始化预算
budget = SchedulingBudget(token_budget=4096, seq_budget=32)
# 1. 优先处理swapped队列
restored = self.schedule_swapped(budget)
self.running.extend(restored)
# 2. 处理waiting队列
if not self.swapped: # 只有swapped为空时才处理新请求
new_seqs = self.schedule_prefills(budget)
self.running.extend(new_seqs)
# 3. 维护running队列
preempted, swapped = self.schedule_running(budget)
self.waiting.extendleft(preempted)
self.swapped.extend(swapped)
# 执行推理
self.run_inference()
4. 实战演示与状态观察
让我们通过具体案例观察队列状态变化:
# 初始化调度器
scheduler = Scheduler()
# 添加初始请求
scheduler.waiting.append(SequenceGroup(1, "Hello"))
scheduler.waiting.append(SequenceGroup(2, "Hi"))
scheduler.waiting.append(SequenceGroup(3, "Greetings"))
# 第一次调度
scheduler.run_scheduling_cycle()
print(f"Running: {[sg.request_id for sg in scheduler.running]}")
# 输出: Running: [1, 2]
# 模拟资源紧张
for _ in range(5):
scheduler.run_scheduling_cycle()
# 观察队列状态
print(f"Waiting: {[sg.request_id for sg in scheduler.waiting]}")
print(f"Running: {[sg.request_id for sg in scheduler.running]}")
print(f"Swapped: {[sg.request_id for sg in scheduler.swapped]}")
典型输出可能显示:
Waiting: [3]
Running: [1, 2]
Swapped: []
当继续模拟资源不足时,会看到请求从running转移到swapped队列。
5. 高级调度策略优化
实际生产环境中还需要考虑以下优化点:
-
优先级调度:为不同请求设置优先级
class SequenceGroup: def __init__(self, ..., priority: int = 0): self.priority = priority -
动态批处理:根据请求特征智能合并
def dynamic_batching(self): # 合并相似长度的请求 groups = sorted(self.waiting, key=lambda x: len(x.prompt)) batches = [groups[i:i+4] for i in range(0, len(groups), 4)] -
内存压缩:对swapped请求进行内存优化
def compress_swapped(self): for seq_group in self.swapped: if len(seq_group.blocks) > 10: seq_group.blocks = self.compress(seq_group.blocks)
通过本文的Python实现,我们完整模拟了vLLM调度器的核心机制。在实际项目中,这种队列流转策略能够有效提升GPU利用率,将语言模型推理的吞吐量提升3-5倍。读者可以在此基础上扩展更复杂的调度策略,或集成到自己的推理框架中。
更多推荐


所有评论(0)