关于本文:01 篇是读 Aleksa Gordic 的 Inside vLLM: Anatomy of a High-Throughput LLM Inference System 后整理的概念笔记,为了好读,省掉了原文里大量的函数名、数据结构和数据流细节。可原文本身也不算好懂:很多函数一笔带过,不解释它干什么;数据在哪几个对象之间传来传去,也缺一张完整的图。
这一篇换个写法:挑一个请求,从
llm.generate()开始,一路跟到它拿到输出,途中遇到的每个函数都用大白话说清”输入什么、做了什么、输出什么”,能画图的地方都画了图(全部是我自己画的)。源码对照的是 vLLM main 分支 commitf5a78f2ad7(2026-09-24),和原文用的 2025 年 8 月版本相比改了不少名字,文中会顺带指出。建议先看 01 篇知道大概有哪些零件,再看这篇;如果想先看一个 1000 多行的”迷你版”,可以读 02 篇的 nano-vllm 精读。
0. 全文地图
整篇文章就是把图 1 从 ① 走到 ⑧。先记住三件事,后面就不容易迷路:
- 两个进程。你写代码的那个 Python 进程是”前端”,只负责文字和 token id 之间的转换;真正调度、跑模型的
EngineCore默认在另一个进程里(环境变量VLLM_ENABLE_V1_MULTIPROCESSING默认是 1)。两边用 ZMQ 传消息,传的全是数字。 - 一个循环。引擎核心就是在反复调用
step()。每转一圈,每个正在生成的请求最多往前走一小步(通常是 1 个 token)。 - 一个数。每个请求身上有个
num_computed_tokens(已经算好 KV 的 token 数)。调度器的全部工作,本质上就是让这个数追上请求的总长度。
下面是这次要跟踪的代码:
1
2
3
4
5
from vllm import LLM, SamplingParams
llm = LLM("Qwen/Qwen2.5-1.5B-Instruct")
outputs = llm.generate(["介绍一下南京"], SamplingParams(max_tokens=64))
print(outputs[0].outputs[0].text)
1. 启动:LLM(...) 那几十秒
请求还没来,引擎得先把自己搭起来。LLM.__init__ 会创建 LLMEngine(vllm/v1/engine/llm_engine.py),它手里有四样东西:
| 成员 | 是什么 | 干什么 |
|---|---|---|
renderer |
渲染器 | 套对话模板、调用 tokenizer 分词 |
input_processor |
输入处理器 InputProcessor |
检查参数、补默认值,打包成 EngineCoreRequest |
output_processor |
输出处理器 OutputProcessor |
把 token id 变回文字,拼最终结果 |
engine_core |
引擎客户端 | 默认是 SyncMPClient:在后台起一个 EngineCoreProc 进程,用 ZMQ 和它说话 |
原文里的
Processor现在叫InputProcessor。如果把VLLM_ENABLE_V1_MULTIPROCESSING设成 0,客户端会换成InprocClient,引擎核心就在同一个进程里直接调用,调试时方便打断点。
后台进程里构造 EngineCore(vllm/v1/engine/core.py),顺序如图 4:
逐个解释一下:
Worker.init_device()(vllm/v1/worker/gpu_worker.py):绑定到某张 GPU,初始化分布式通信(单卡也要走一遍),拍一张显存快照,按gpu_memory_utilization(默认 0.92)算出允许用多少,最后创建GPUModelRunner。Worker.load_model():把权重读进显存。determine_available_memory():用假数据跑一次profile_run,量出激活峰值,再扣掉权重和 CUDA Graph 的预估,剩下的就是 KV 缓存的地盘。get_kv_cache_configs+initialize_from_config:剩余显存 ÷ 每块大小 = 一共能切多少块(num_gpu_blocks),然后真的分配这块大张量。compile_or_warm_up_model():预热 kernel,录 CUDA Graph(capture_model)。录图的批大小默认是[1, 2, 4, 8, 16, …, 248, 256, 272, …, 512]这样一串,以后实际批大小会被补齐到最近的一档,直接回放。
这几步做完才创建 Scheduler,因为调度器必须知道”一共有多少块可以分”。
另外有个默认开启的开关要提前说:异步调度(async scheduling)。它开着时引擎用的不是 step(),而是 step_with_batch_queue(),第 3 节再讲。
2. 请求进门:从一句话到一个 Request
2.1 前端:generate() 做了什么
LLM.generate() 本身很薄,真正的逻辑在 vllm/entrypoints/offline_utils.py:
_add_completion_requests→renderer.render_cmpl:套模板、分词。”介绍一下南京”变成一串 token id,比如[100, 23, 57, 9, …]。_add_request:给请求编个号(从 0 开始的计数器),把输出模式设成FINAL_ONLY(离线模式只要最终结果,不要中间的流式增量),然后调LLMEngine.add_request。LLMEngine.add_request做三件事:input_processor.process_inputs(...):如果你没设max_tokens,就补成”上下文上限 − 提示词长度”;从模型的 generation config 里读出 EOS token;最后打包成EngineCoreRequest。output_processor.add_request(...):前端也登记一下这个请求,准备好它的反分词器。engine_core.add_request(...):通过 ZMQ 发给后台进程。
_run_engine:一个朴素的循环:
1
2
3
4
while self.llm_engine.has_unfinished_requests():
step_outputs = self.llm_engine.step()
... # 收集已完成的输出
# 最后按请求编号排序返回
注意这里的 LLMEngine.step() 是前端的 step:它只是去后台拿一批结果(engine_core.get_output()),交给输出处理器,真正的调度和计算在后台进程自己的循环里。
2.2 三副面孔
同一个请求一路上换了几次数据结构,这是读源码时最容易晕的地方:
EngineCoreRequest(vllm/v1/engine/__init__.py):跨进程发送用的”包裹”,是一个msgspec.Struct,序列化很快,字段都是只读的原始数据。Request(vllm/v1/request.py):后台进程收到包裹后,用Request.from_engine_core_request(req, block_hasher)拆包得到。它是调度器天天打交道的”病历本”,有状态、有计数器,会不断被修改。
Request 里最重要的几个字段:
1
2
3
4
5
6
7
8
9
10
11
12
13
self.status = RequestStatus.WAITING
self._all_token_ids = self.prompt_token_ids.copy() # 提示词 + 之后生成的
self.num_computed_tokens = 0 # 已经算好 KV 的个数
self.spec_token_ids = [] # 投机解码的草稿
self.block_hashes = ... # 每个"满块"的哈希,建对象时就算好
@property
def num_tokens(self): # 当前总长度
return len(self._all_token_ids)
@property
def num_tokens_with_spec(self): # 算上草稿的长度
return len(self._all_token_ids) + len(self.spec_token_ids)
每采样出一个新 token,就调用 append_output_token_ids():token 追加到 _all_token_ids,如果刚好填满一个块,顺手把这个块的哈希也算出来(update_block_hashes(),前缀缓存要用,见 4.3 节)。
2.3 请求的状态
RequestStatus 一共 12 个值,分三类看就简单了:
- 等待类:普通的
WAITING,以及三种”卡住的等待”(等语法编译、等远端 KV、等流式输入)。卡住的请求调度器会先跳过,等条件满足再转成WAITING。 - 运行类:
RUNNING,以及被抢占后的PREEMPTED。 - 结束类:六种
FINISHED_*。源码里判断”结束了没”就一句status > RequestStatus.PREEMPTED,因为结束类的枚举值都排在后面。
新请求建好后状态是 WAITING(要做约束解码的是 WAITING_FOR_STRUCTURED_OUTPUT_GRAMMAR),被放进调度器的 waiting 队列。默认 FCFS 策略下,这个队列就是一个 deque:add 放队尾,pop 从队头取,prepend(被抢占的请求用)插回队头。
3. 引擎的心跳:EngineCore.step()
后台进程的主循环 run_busy_loop 很短:从输入队列里取新请求 → 执行一次 step → 把结果放进输出队列。先看最简单的 step(),完整代码只有十几行(core.py):
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
def step(self):
if not self.scheduler.has_requests():
return {}, False
# 1) 调度:决定这一步算谁、算几个
scheduler_output = self.scheduler.schedule(...)
# 2) 让 GPU 开始跑前向。non_block=True:发出去立刻返回一个 future
future = self.model_executor.execute_model(scheduler_output, non_block=True)
# 3) 趁 GPU 在忙,CPU 算约束解码的掩码
grammar_output = self.scheduler.get_grammar_bitmask(scheduler_output)
# 4) 等前向结束;它返回 None,表示"logits 算好了,请采样"
model_output = future.result()
if model_output is None:
model_output = self.model_executor.sample_tokens(grammar_output)
# 5) 处理这期间收到的取消请求,再把结果记到各个请求上
self._process_aborts_queue()
engine_core_outputs = self.scheduler.update_from_output(
scheduler_output, model_output)
return engine_core_outputs, ...
原文把 step 讲成”调度、前向、收尾”三步,新版在中间多拆了一刀:前向(execute_model)和采样(sample_tokens)是两次调用。拆开的目的是让 CPU 在 GPU 跑前向的同时去算语法掩码,两件事重叠起来。
默认情况下用的其实是图 5(b) 的 step_with_batch_queue()。它维护一个最多装 2 批的队列:
- 队列没满,就调度一批新的、发给 GPU,然后不等结果直接返回;
- 队列满了(或者没东西可调度了),才阻塞等最早那一批的结果,调用
update_from_output。
效果是 GPU 在算第 N 批时,CPU 已经在给第 N+1 批做调度,GPU 几乎不空转。代价是调度第 N+1 批时还不知道第 N 批采到了什么 token,所以请求上有个 num_output_placeholders 字段,先给”还没回来的 token”占个位。初读可以先按图 5(a) 理解,逻辑是一样的。
下面三节分别拆开 step 里的三大块:调度(第 4 节)、GPU 执行(第 5 节)、收尾(第 6 节)。
4. 调度:Scheduler.schedule()
4.1 核心思想:追赶游戏
schedule()(vllm/v1/core/sched/scheduler.py)开头有一段很重要的注释:
1
2
3
4
5
6
7
# NOTE(woosuk) on the scheduling algorithm:
# There's no "decoding phase" nor "prefill phase" in the scheduler.
# Each request just has the num_computed_tokens and
# num_tokens_with_spec. ...
# At each step, the scheduler tries to assign tokens to the requests
# so that each request's num_computed_tokens can catch up its
# num_tokens_with_spec.
翻译成人话:调度器不区分预填充和解码,它只看每个请求的两个数,想办法让 num_computed_tokens 追上 num_tokens_with_spec。
图 6 是我们那个请求的前几步。提示词有 10 个 token,假设这一步的预算只剩 6:
- 第 1 步只能追 6 个,这就是分块预填充。提示词还没读完,不采样。
- 第 2 步追完剩下 4 个,追平了,对最后一个位置采样出第一个新 token。
- 之后每采样出一个 token,总长度 +1,于是又差 1 个,下一步就追这 1 个,这就是解码。
每个请求这一步要算几个 token,公式就是一行:
1
2
3
num_new_tokens = request.num_tokens_with_spec + request.num_output_placeholders \
- request.num_computed_tokens
num_new_tokens = min(num_new_tokens, token_budget) # 还要被预算、长度上限等截断
token 预算来自 max_num_batched_tokens,一批最多几个请求来自 max_num_seqs。这两个默认值和显卡有关:LLM 离线用法下,消费级卡和 A100 是 8192 / 256,H100 这类更大的卡是 16384 / 1024。
4.2 两段循环
schedule() 的主体就是图 7 的两段循环:
第一段:先照顾 running 里的请求。 它们已经占着 KV 块,让它们尽快跑完才能腾出地方。对每个请求:算 num_new_tokens → 调 kv_cache_manager.allocate_slots() 要块 → 要到了就扣预算,下一个;要不到就抢占。
抢占的逻辑在 _preempt_request(),去掉细节后就是:
1
2
3
4
5
6
self._free_request_blocks(request) # 块全部还回去
request.status = RequestStatus.PREEMPTED
request.num_computed_tokens = 0 # 已算的全部作废,以后重来
request.spec_token_ids = []
request.num_preemptions += 1
self.waiting.prepend_request(request) # 插回等待队列的最前面
踢谁?默认 FCFS 策略下踢 running 列表的最后一个,也就是最晚进来的那个;如果用 PRIORITY 策略,踢优先级最低(数值最大)、同优先级里最晚到的。踢完再重试 allocate_slots,直到要到块,或者把自己也踢掉了。被踢的请求重新调度时,前缀缓存往往还能把大部分块捡回来(4.3 节),所以”从头重算”没有听起来那么亏。
第二段:预算还有剩,再从 waiting 接新请求。 只有这一步没人被抢占时才会进入这一段,道理很简单:既然已经穷到要踢人了,就别再接新人。对队首的请求:
- 还卡在”等语法编译 / 等远端 KV”之类的状态?先挪到一边(
skipped_waiting),看下一个。 kv_cache_manager.get_computed_blocks(request):查前缀缓存,看开头有多少 token 的 KV 已经有现成的。命中 k 个,就相当于开局num_computed_tokens = k。- 如果配置了 KV 连接器,再问它外部有没有更多现成的(
get_num_new_matched_tokens,7.3 节)。 num_new_tokens = 总长度 − 已命中,再被预算截断。allocate_slots(...):要块。要不到就break,注意这里不抢占别人。- 要到了:从
waiting移到running,状态改RUNNING。
两段都走完,打包成 SchedulerOutput(4.5 节),最后调用 _update_after_schedule():先乐观地把每个被调度请求的 num_computed_tokens 加上本步的 token 数。为什么能先加?因为 GPU 一定会把这些 token 的 KV 算出来;唯一的例外是投机解码被拒绝的草稿,第 6 节会减回去。
4.3 分块:allocate_slots
KVCacheManager.allocate_slots()(vllm/v1/core/kv_cache_manager.py)是调度器最常调的函数。它的 docstring 里有一张 ASCII 布局图,把一个请求的 token 分成五段:
1
| < comp > | < new_comp > | < ext_comp > | < new > | < lookahead > |
comp:之前的步已经算好的;new_comp:这一步刚在前缀缓存里命中的;ext_comp:KV 连接器从外部搬来的;new:这一步要算的(包括投机解码的草稿);lookahead:给 EAGLE 这类起草模型预留的位置。
新手先只看 comp、new_comp、new 三段就够了。函数做的事按顺序是(对应图 8 的 ① ~ ⑦):
- 算这个请求一共需要多少个”槽位”(token 位置);
- 算出需要几块,扣掉手里已有的,得出还要新领几块;
- 空闲块不够?返回
None。调度器就是靠这个None决定抢占还是放弃; - 前缀命中的块调用
block_pool.touch(),把引用计数加 1,防止被别人领走; - 调
block_pool.get_new_blocks(n)领新块; cache_blocks():把已经填满的块登记进前缀缓存;- 返回这次新领的块,调度器把块号写进
SchedulerOutput带给 Worker。
每个 KV 块多大?page_size_bytes = 2 × block_size × num_kv_heads × head_size × dtype 字节数,这是一层的大小,乘以层数才是一个块的总占用。block_size 默认 16。
4.4 块池:BlockPool
BlockPool(vllm/v1/core/block_pool.py)管着所有物理块,里面三样东西(图 9):
free_block_queue:空闲块的双向链表(FreeKVCacheBlockQueue,在kv_cache_utils.py)。链表带一个假头和一个假尾,所以从中间摘掉某个块也是 O(1),前缀命中时就需要这个操作。cached_block_hash_to_block:哈希 → 块 的表,前缀缓存靠它查找。- 每个
KVCacheBlock的引用计数ref_cnt:几个请求正在用它。块 0 是null_block,初始化时就被拿走当占位符,永远不会分配出去。
几个关键函数:
1
2
3
4
5
6
def get_new_blocks(self, num_blocks):
ret = self.free_block_queue.popleft_n(num_blocks) # 从队头拿
for block in ret:
self._maybe_evict_cached_block(block) # 块若还带旧哈希,从哈希表删掉
block.ref_cnt += 1
return ret
1
2
3
4
5
6
7
8
9
10
def free_blocks(self, ordered_blocks):
for block in ordered_blocks:
block.ref_cnt -= 1
if block.ref_cnt == 0:
if block.block_hash is None:
blocks_to_evict_first.append(block) # 没哈希:以后没人能命中它
else:
blocks_to_evict_last.append(block) # 有哈希:尽量留着
self.free_block_queue.prepend_n(blocks_to_evict_first) # 插队头,先被复用
self.free_block_queue.append_n(blocks_to_evict_last) # 放队尾,晚被复用
这里藏着前缀缓存的一个关键设计:块被释放后并没有清空,它回到空闲队列,但哈希还留在表里。只要它还没被别人领走(领走时 _maybe_evict_cached_block 才删哈希),后来的请求就能通过哈希找到它,touch 一下再捡回来用。有哈希的块放队尾,等于”最近用过的最后被驱逐”,近似 LRU。
请求结束时,块是按倒序释放的(free_blocks(reversed(blocks)))。这样链尾的块先进入被驱逐的位置,开头的块留得更久。开头的块是最有可能被别人共享的(比如系统提示词),这样安排最划算。
4.5 前缀缓存:链式哈希
哈希函数是 kv_cache_utils.py 里的 hash_block_tokens:
1
hash_block_tokens(parent_block_hash or NONE_HASH, tuple(block_token_ids), extra_keys)
它把前一块的哈希也算进去,所以”第 k 块的哈希命中”就意味着”从开头到第 k 块全部一样”。extra_keys 里有 LoRA 名、多模态图片的哈希、cache_salt,避免”token 一样但含义不同”的误命中。默认哈希算法是 sha256,第一块用的种子 NONE_HASH 是一个固定值。
查找发生在 get_computed_blocks():从第一块开始逐块查表,碰到第一个查不到的就停,前面连续命中的就是可复用的前缀。图 10(b) 用两个请求演示了一遍:A 算完后把 3 个满块登记进表;B 和 A 开头一样,前两块直接命中,只需要算后面的。
有一个容易忽略的细节:最多只能命中 num_tokens − 1 个 token。因为最后一个 token 必须真的过一遍模型,才有 logits 可以采样下一个 token。
4.6 调度结果:SchedulerOutput
schedule() 的产物 SchedulerOutput(vllm/v1/core/sched/output.py)里,请求分成两种:
scheduled_new_reqs:本步第一次被调度的请求,每个是一个NewRequestData,带着全部信息:提示词、采样参数、整张块表……scheduled_cached_reqs:之前已经调度过的请求,打包成一个CachedRequestData,只带变化量:新领的块号new_block_ids、最新的num_computed_tokens、哪些是抢占后回来的resumed_req_ids。
为什么要分?一个请求可能活上几百步,每步都把整段提示词发给 Worker 太浪费。Worker 那边自己维护了一份”影子账本”,调度器只需要告诉它”变了什么”。其他字段还有 num_scheduled_tokens(每个请求本步算几个)、scheduled_spec_decode_tokens(草稿)、finished_req_ids(上一步结束的,让 Worker 清理)等。
5. GPU 执行:从 SchedulerOutput 到新 token
5.1 调用链
execute_model 从引擎核心到 GPU,要经过好几层转发:
1
2
3
4
5
EngineCore.step()
→ Executor.execute_model() # 单卡是 UniProcExecutor,多卡是 MultiprocExecutor
→ collective_rpc("execute_model") # "在每个 Worker 上调用这个方法"
→ Worker.execute_model() # gpu_worker.py
→ GPUModelRunner.execute_model() # gpu_model_runner.py,真正干活
单卡时前面几层都只是转手。新版默认的 ModelRunner 是重写过的 V2 版本,但结构和思路与 V1(gpu_model_runner.py)一致,下面按 V1 讲,原文也是按它讲的。
GPUModelRunner 的工作分成两个大函数(图 14),依次解释。
5.2 _update_states:同步影子账本
Worker 这边有两份数据结构:
self.requests:{req_id: CachedRequestState},每个请求的 token、块表、采样参数;self.input_batch:一个InputBatch,把所有请求的信息按”槽位”摊成大数组,比如token_ids_cpu是一个[最大请求数 × 最大长度]的二维 numpy 数组,每一行是一个请求的全部 token。
每步开头的 _update_states() 按 SchedulerOutput 更新它们:删掉结束的请求 → 为新请求建 CachedRequestState 并放进批 → 给老请求追加新块、更新已算数(被抢占回来的请求,块表整张替换)→ 最后 condense(),把后面的请求挪到前面的空槽里,保证数组前 N 行都是有效的。
5.3 _prepare_inputs:把一批请求拼成一条长序列
这是 01 篇说”第一次读不用追求全懂”的那个函数。其实源码注释里带了一个很好的数字例子,我把它画成了图 12:
三个请求本步分别算 2、5、3 个 token,一共 10 个,首尾相接拼成一条长序列。关键几行(注释是源码里的原话):
1
2
3
4
5
6
7
8
9
10
11
12
13
14
# Get request indices.
# E.g., [2, 5, 3] -> [0, 0, 1, 1, 1, 1, 1, 2, 2, 2]
req_indices = np.repeat(self.arange_np[:num_reqs], num_scheduled_tokens)
# cu_num_tokens: [2, 5, 3] -> [2, 7, 10]
# query_pos (段内序号): [0, 1, 0, 1, 2, 3, 4, 0, 1, 2]
# 每个 token 在自己请求里的真实位置 = 已算数 + 段内序号
positions_np = self.input_batch.num_computed_tokens_cpu[req_indices] + query_pos
# E.g., [0, 1, 0, 1, 2, 3, 4, 0, 1, 2]
# -> [0, 1, M, M + 1, M + 2, M + 3, M + 4, 2 * M, 2 * M + 1, 2 * M + 2]
# where M is the max_model_len.
token_indices = positions_np + req_indices * self.input_batch.token_ids_cpu.shape[1]
token_indices 这一步最巧:token_ids_cpu 是二维表,摊平后第 r 行第 p 列的下标就是 r × M + p。算出所有下标后一次 index_select,就把 10 个 token id 取出来了,不用写 Python 循环。
顺手还算出几样给注意力算子用的”说明书”:
query_start_loc = [0, 2, 7, 10]:每段从哪开始,算子靠它分清段与段的边界;seq_lens:每个请求这步要看多长的历史(已算 + 本步);slot_mapping:每个新 token 的 KV 写到显存哪个格子(5.4 节);logits_indices = query_start_loc[1:] − 1:每段最后一个位置,只在这些位置算 logits、采样。
这些运算全在 CPU 上用 numpy 做,结果放进预先分配好的 CpuGpuBuffer(CPU、GPU 各有一份同样大小的缓冲区),再一次性拷到 GPU。
5.4 slot_mapping:KV 写到哪
把 KV 缓存摊平看,就是一长排”格子”,每个格子装一个 token 的 K 和 V,第 b 块占第 b × block_size 到 b × block_size + block_size − 1 格。所以一个 token 该写到哪一格,只要查块表:
1
slot = block_table[req, pos // block_size] * block_size + pos % block_size
真实代码里这个公式由一个 Triton kernel(block_table.compute_slot_mapping)在 GPU 上批量算。注意力层里用到它两次:
- 写:
reshape_and_cache_flash(key, value, kv_cache, slot_mapping)把本步算出的 K、V 散写进对应格子; - 读:
flash_attn_varlen_func(..., block_table=..., seqused_k=seq_lens)按块表把历史 KV 找回来做注意力。”varlen”就是”变长”的意思,一批里每段长度可以不同,靠的就是query_start_loc。
FlashAttention 后端的 KV 缓存张量形状是 [块数, KV 头数, block_size, 2 × head_size],K 和 V 拼在最后一维(vllm/v1/attention/backends/flash_attn.py)。
5.5 前向,然后返回 None
准备好输入后,execute_model 接着:
- 决定 CUDA Graph 模式和填充:比如这步有 10 个 token,就补齐到录好的 16,回放录制好的图(默认模式
FULL_AND_PIECEWISE:纯解码批用整图,混合批用分段图); _build_attention_metadata:把上面那些”说明书”打包成注意力后端需要的FlashAttentionMetadata;set_forward_context+_model_forward:把元数据放进一个全局上下文(每层注意力从这里取),然后真正跑一遍 Transformer,得到每个位置的hidden_states;hidden_states[logits_indices]→compute_logits:只在每段最后一个位置算 logits;- 把 logits 暂存到
self.execute_model_state,返回None。
返回 None 就是在告诉引擎核心:”前向完了,请调 sample_tokens。”
5.6 sample_tokens 和 Sampler
sample_tokens(grammar_output) 依次:盖上语法掩码(apply_grammar_bitmask,7.1 节)→ _sample 采样 → 如果开了投机解码,起草下一轮的猜测 → _bookkeeping_sync 把新 token 写回 input_batch 并拷回 CPU → 返回 ModelRunnerOutput(主要是 req_ids 和每个请求的 sampled_token_ids)。
Sampler.forward(vllm/v1/sample/sampler.py)的顺序如图 15:
- 如果用户要 logprobs,先用原始 logits 算(在惩罚和温度之前,这点和老版本不同);
- 转成 float32,依次应用各种 logits 处理器:禁用词、
logit_bias、min_tokens(还没到最小长度就把 EOS 屏蔽掉)、重复/频率/存在惩罚; - 贪心请求直接
argmax,如果整批都是贪心,到这就返回了; - 随机请求:除以 temperature → min_p → top_k / top_p 截断 → 按概率抽一个。
最后”按概率抽一个”有个小技巧:没有用 torch.multinomial,而是给每个词生成一个指数分布的随机数 q,取 probs / q 最大的那个。数学上等价于按概率抽样,但全程在 GPU 上完成,不需要 CPU 同步。同一批里贪心和随机的请求,最后用 torch.where 各取各的结果。
6. 收尾:记账和”该停了吗”
6.1 引擎侧:update_from_output
ModelRunnerOutput 回到引擎核心,Scheduler.update_from_output() 对每个请求:
- 投机解码回滚:如果本步验证了草稿,被拒绝了几个,就把
num_computed_tokens减回去几个(4.2 节里那个”乐观地先加上”,在这里修正)。 - 追加新 token:
request.append_output_token_ids(token)。 check_stop(vllm/v1/core/sched/utils.py):只靠 token id 就能判断的停止条件:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
if last_token_id == sampling_params.eos_token_id:
request.status = RequestStatus.FINISHED_STOPPED
return True
if last_token_id in (sampling_params.stop_token_ids or ()):
request.status = RequestStatus.FINISHED_STOPPED
request.stop_reason = last_token_id
return True
if (request.num_tokens >= max_model_len
or request.num_output_tokens >= request.max_tokens):
request.status = RequestStatus.FINISHED_LENGTH_CAPPED
return True
if request.num_output_tokens < sampling_params.min_tokens:
return False
...
- 结束了就释放:
_free_request()把请求 id 记进finished_req_ids(下一步告诉 Worker 删掉它),调用kv_cache_manager.free()还块。 - 打包:每个请求一个
EngineCoreOutput(request_id, new_token_ids, finish_reason, …),整批装进EngineCoreOutputs,通过 ZMQ 发回前端。
6.2 前端侧:process_outputs
前端的 OutputProcessor.process_outputs()(vllm/v1/engine/output_processor.py)对每个输出:
- 增量反分词:
detokenizer.update()只把新增的 token 解成文字接在后面,不用每次重解整段; check_stop_strings:检查停止字符串,比如stop=["\n\n"]。停止字符串可能横跨好几个 token,必须先有文字才能判断,所以放在前端。命中了就把多余的文字截掉,并通知引擎 abort 这个请求(LLMEngine.step里的abort_requests);make_request_output:离线的FINAL_ONLY模式下,没结束就返回None;结束了才拼出RequestOutput,里面的outputs[0]是一个CompletionOutput,有text、token_ids、finish_reason。
_run_engine 收齐所有结束的请求,按编号排好序,generate() 返回。到这里,我们的请求走完了图 1 的 ① ~ ⑧。
7. 在骨架上加功能
有了上面的骨架,高级功能就是在某几个环节上各加一点逻辑。
7.1 约束解码(structured outputs)
用法:
1
2
3
from vllm.sampling_params import StructuredOutputsParams
params = SamplingParams(structured_outputs=StructuredOutputsParams(choice=["是", "否"]))
# 也可以是 json=schema、regex=...、grammar=...
走过的环节(图 17):
- 进门时:
StructuredOutputManager.grammar_init()把”编译语法”这件事交给一个线程池,拿到一个Future挂在请求上,请求状态设成WAITING_FOR_STRUCTURED_OUTPUT_GRAMMAR; - 调度时:遇到这种状态的请求,调度器检查 Future 完成没有(
_try_promote_blocked_waiting_request),完成了就转成WAITING,否则先跳过; - step 里:就是第 3 节那行
get_grammar_bitmask(),趁 GPU 跑前向,CPU 问每个受约束请求的状态机”现在允许哪些 token”,填成一张位掩码:每个 token 一个比特,32 个 token 挤进一个 int32; - 采样前:
apply_grammar_bitmask在 GPU 上把比特为 0 的 token 的 logits 改成负无穷; - 收尾时:
update_from_output里调用accept_tokens(),把采到的 token 喂给状态机,状态机前进一步。
7.2 投机解码(speculative decoding)
以最简单的 n-gram 起草为例(vllm/v1/spec_decode/ngram_proposer.py):在已有的 token 里找和结尾”最长匹配”的片段,把它后面的 k 个 token 抄过来当猜测(图 18a)。不需要任何额外模型。
数据怎么流:
sample_tokens采完样后,runner 调propose_draft_token_ids()起草,存在自己身上;- 引擎核心通过另一次 RPC
take_draft_token_ids()把草稿取回来,update_draft_token_ids()放进request.spec_token_ids; - 下一次
schedule():num_tokens_with_spec变大了 k,于是调度器自然给这个请求排 1 + k 个 token,草稿写进scheduled_spec_decode_tokens; - GPU 一次前向算出这 1 + k 个位置的概率,拒绝采样器(
vllm/v1/sample/rejection_sampler.py)从左往右逐个判定(图 18b):大模型概率 p 除以草稿概率 q ≥ 随机数 u 就接受(n-gram 没有草稿概率,q 当作 1);一旦拒绝,就在这个位置按max(p − q, 0)重采一个,后面的草稿全部作废;如果全部接受,还白送一个 bonus token; update_from_output:num_computed_tokens −= 被拒绝的个数。
这套接受规则保证了输出分布和不开投机解码完全一样。除了 n-gram,spec_decode/ 目录下还有 EAGLE / EAGLE-3、Medusa、MTP、draft model、suffix decoding 等方法。
7.3 KV 连接器与 PD 分离
KV 连接器的基类在 vllm/distributed/kv_transfer/kv_connector/v1/base.py,它的方法分成两组(图 19):
- 调度器侧决定”要不要搬、搬哪些”:接新请求时
get_num_new_matched_tokens问外部有没有现成的 KV;allocate_slots之后update_state_after_alloc记下搬到哪些块;schedule()末尾build_connector_meta把搬运清单塞进SchedulerOutput;请求结束时request_finished可以说”块先别还”。 - Worker 侧真正动手:前向前
start_load_kv,每层注意力前wait_for_layer_load、之后save_kv_layer,前向结束wait_for_save,最后get_finished汇报搬完了哪些。每层的两个钩子是用装饰器maybe_transfer_kv_layer挂在注意力层上的。
学习用的 example_connector.py(原文里叫 SharedStorageConnector)逻辑很直白:把每层 KV 存成 /tmp 下的 safetensors 文件,文件夹名是提示词的哈希;下次同样的提示词进来,发现文件夹存在,就读回来。PD 分离就是:预填充实例只存不读,max_tokens=1;解码实例只读,读完直接开始解码。
8. 多卡:MultiprocExecutor
开 tensor_parallel_size=4 后,执行器换成 MultiprocExecutor(vllm/v1/executor/multiproc_executor.py)。调度器、引擎核心完全不变,变化都在”怎么把一次 execute_model 发给 4 张卡”:
- 启动握手:主进程先建一个广播队列
rpc_broadcast_mq→ 为每张卡起一个WorkerProc进程,带两根管道:ready 管道(子进程报告”我准备好了”)和 death 管道(子进程靠它发现父进程死了)→ 每个 Worker 加载模型、建自己的应答队列,通过 ready 管道把队列句柄发回来 → 主进程连上所有应答队列 → Worker 进入worker_busy_loop死循环,等任务。 - 每步:
collective_rpc往广播队列里放一个元组(方法名, 参数, …, output_rank),4 个 Worker 同时读到、同时执行。卡之间每层用 NCCL all-reduce 拼结果。 - 只收一份:4 张卡采样结果一样,所以只有”输出 rank”(最后一个流水线阶段的第一个 TP Worker,只有 TP 时就是 Worker 0)把
ModelRunnerOutput放回应答队列。
广播队列 MessageQueue(vllm/distributed/device_communicators/shm_broadcast.py)是一个共享内存环形缓冲区:一个写者、多个读者,每个格子带几个标志字节(”写好了”、”谁读过了”),写一次所有进程都能读到,不用复制 4 份。消息超过 24 MiB(比如特别大的语法掩码)时,格子里只写一个”溢出”标记,真实数据改走 ZMQ。
9. 在线服务:一个 curl 请求的来回
vllm serve 起的服务,和离线用法相比只是最外层换了:
- 路由:
POST /v1/completions进到vllm/entrypoints/openai/completion/api_router.py的create_completion,再交给OpenAIServingCompletion.create_completion。(原文里的openai/api_server.py现在只是一个兼容转发,真正的启动代码在vllm/entrypoints/launchers/api_server/。) AsyncLLM.generate()(vllm/v1/engine/async_llm.py):LLMEngine的异步版。分词、为这个请求建一个输出队列、发给引擎,然后在队列上await,拿到一个就yield一个,这就是流式返回。- 挑引擎:数据并行(DP)时,
DPLBAsyncMPClient.get_core_engine_for_request()给每个引擎打分,挑最低的:
1
2
3
4
score = max(self.client_count * inflight, waiting + running)
if waiting:
# KV 占用率 50% 以下不罚;到 100% 时,等待数最多再罚 3 倍
score += waiting * 6.0 * max(0.0, kv_cache_usage - 0.5)
原文里的公式是 waiting × 4 + running,新版多了 KV 占用率这一项:KV 快满的引擎,排队的请求消化得慢,新请求应该尽量避开。
- 引擎进程
EngineCoreProc:一个输入线程收包解码、放进队列;主循环run_busy_loop反复取请求、step();一个输出线程把结果编码发回。ZMQ 收发会释放 GIL,所以网络 IO 和 GPU 计算可以重叠。 output_handler:AsyncLLM里的一个后台协程,不断从引擎拿EngineCoreOutputs,交给OutputProcessor反分词,再把结果放进各请求的队列,第 2 步那个await就醒了。
DP 大于 1 时还有一个 DP 协调器进程(vllm/v1/engine/coordinator.py):收集各引擎的等待数、运行数、KV 占用率,最多每 100 ms 广播给前端做负载均衡用;另外管理”请求波次”:所有引擎都空闲时会一起休眠,有新请求时由它唤醒大家。在引擎这边(DPEngineCoreProc),如果自己没活干但别的 rank 还在忙,会跑一个”空批”(dummy batch),因为 MoE 这类模型每层都需要所有 rank 一起通信,少一个就会卡住。
压测命令还是那几个:vllm bench latency、vllm bench throughput、vllm bench serve,新版还多了 startup(测启动时间)和 sweep(批量扫参数、画图)。
10. 小结:一张速查表
| 环节 | 关键函数 | 文件 | 一句话 |
|---|---|---|---|
| 进门 | InputProcessor.process_inputs |
v1/engine/input_processor.py |
文字 → EngineCoreRequest |
| 建档 | Request.from_engine_core_request |
v1/request.py |
包裹 → 可修改的”病历本” |
| 心跳 | EngineCore.step / step_with_batch_queue |
v1/engine/core.py |
调度 → 前向 → 采样 → 记账 |
| 调度 | Scheduler.schedule |
v1/core/sched/scheduler.py |
让 num_computed_tokens 追上总长度 |
| 分块 | KVCacheManager.allocate_slots |
v1/core/kv_cache_manager.py |
要块,不够返回 None |
| 块池 | get_new_blocks / touch / free_blocks |
v1/core/block_pool.py |
空闲链表 + 哈希表 + 引用计数 |
| 前缀 | get_computed_blocks / hash_block_tokens |
同上 + kv_cache_utils.py |
链式哈希,逐块查到第一个未命中 |
| 同步 | GPUModelRunner._update_states |
v1/worker/gpu_model_runner.py |
按差量更新 Worker 的影子账本 |
| 拼输入 | _prepare_inputs |
同上 | 一批请求 → 一条长序列 + 说明书 |
| 采样 | Sampler.forward |
v1/sample/sampler.py |
惩罚 → 温度 → top-k/p → 抽样 |
| 记账 | Scheduler.update_from_output + check_stop |
scheduler.py、sched/utils.py |
追加 token,判断停不停 |
| 出门 | OutputProcessor.process_outputs |
v1/engine/output_processor.py |
反分词 + 停止字符串 |
如果只想记一句话:请求在前端变成数字,在引擎里被调度器一步步”追平”,每一步都要先分到 KV 块、再拼进一条长序列跑一遍模型,采出的 token 回到请求上,直到某条停止规则说”够了”。
读源码时我自己的顺序建议:先看 core.py 的 step()(二十行),再看 schedule() 开头那段注释和两段循环,然后跳到 _prepare_inputs 对着图 12 看那几行 numpy,最后回头看 block_pool.py。高级功能和多卡、在线服务都可以等骨架熟了再看。
参考
- Aleksa Gordic, Inside vLLM: Anatomy of a High-Throughput LLM Inference System, vLLM Blog, 2025-09-05
- vLLM 源码仓库(本文对照 commit
f5a78f2ad7) - Leviathan et al., Fast Inference from Transformers via Speculative Decoding, ICML 2023
- 本系列:01 一个请求在 vLLM 里的一生、02 nano-vllm 源码逐段精读