vLLM 学习笔记 03:跟着一个请求走完 vLLM 源码

给新手的图解版:每个函数干什么、数据怎么流

Posted by Liu Mengxuan on September 29, 2026

关于本文:01 篇是读 Aleksa Gordic 的 Inside vLLM: Anatomy of a High-Throughput LLM Inference System 后整理的概念笔记,为了好读,省掉了原文里大量的函数名、数据结构和数据流细节。可原文本身也不算好懂:很多函数一笔带过,不解释它干什么;数据在哪几个对象之间传来传去,也缺一张完整的图。

这一篇换个写法:挑一个请求,从 llm.generate() 开始,一路跟到它拿到输出,途中遇到的每个函数都用大白话说清”输入什么、做了什么、输出什么”,能画图的地方都画了图(全部是我自己画的)。源码对照的是 vLLM main 分支 commit f5a78f2ad7(2026-09-24),和原文用的 2025 年 8 月版本相比改了不少名字,文中会顺带指出。

建议先看 01 篇知道大概有哪些零件,再看这篇;如果想先看一个 1000 多行的”迷你版”,可以读 02 篇的 nano-vllm 精读。


0. 全文地图

图 1

整篇文章就是把图 1 从 ① 走到 ⑧。先记住三件事,后面就不容易迷路:

  1. 两个进程。你写代码的那个 Python 进程是”前端”,只负责文字和 token id 之间的转换;真正调度、跑模型的 EngineCore 默认在另一个进程里(环境变量 VLLM_ENABLE_V1_MULTIPROCESSING 默认是 1)。两边用 ZMQ 传消息,传的全是数字。
  2. 一个循环。引擎核心就是在反复调用 step()。每转一圈,每个正在生成的请求最多往前走一小步(通常是 1 个 token)。
  3. 一个数。每个请求身上有个 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:

图 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:

  1. _add_completion_requests → renderer.render_cmpl:套模板、分词。”介绍一下南京”变成一串 token id,比如 [100, 23, 57, 9, …]。
  2. _add_request:给请求编个号(从 0 开始的计数器),把输出模式设成 FINAL_ONLY(离线模式只要最终结果,不要中间的流式增量),然后调 LLMEngine.add_request。
  3. LLMEngine.add_request 做三件事:
    • input_processor.process_inputs(...):如果你没设 max_tokens,就补成”上下文上限 − 提示词长度”;从模型的 generation config 里读出 EOS token;最后打包成 EngineCoreRequest。
    • output_processor.add_request(...):前端也登记一下这个请求,准备好它的反分词器。
    • engine_core.add_request(...):通过 ZMQ 发给后台进程。
  4. _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 三副面孔

图 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 请求的状态

图 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

默认情况下用的其实是图 5(b) 的 step_with_batch_queue()。它维护一个最多装 2 批的队列:

  1. 队列没满,就调度一批新的、发给 GPU,然后不等结果直接返回;
  2. 队列满了(或者没东西可调度了),才阻塞等最早那一批的结果,调用 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

图 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 两段循环

图 7

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 接新请求。 只有这一步没人被抢占时才会进入这一段,道理很简单:既然已经穷到要踢人了,就别再接新人。对队首的请求:

  1. 还卡在”等语法编译 / 等远端 KV”之类的状态?先挪到一边(skipped_waiting),看下一个。
  2. kv_cache_manager.get_computed_blocks(request):查前缀缓存,看开头有多少 token 的 KV 已经有现成的。命中 k 个,就相当于开局 num_computed_tokens = k。
  3. 如果配置了 KV 连接器,再问它外部有没有更多现成的(get_num_new_matched_tokens,7.3 节)。
  4. num_new_tokens = 总长度 − 已命中,再被预算截断。
  5. allocate_slots(...):要块。要不到就 break,注意这里不抢占别人。
  6. 要到了:从 waiting 移到 running,状态改 RUNNING。

两段都走完,打包成 SchedulerOutput(4.5 节),最后调用 _update_after_schedule():先乐观地把每个被调度请求的 num_computed_tokens 加上本步的 token 数。为什么能先加?因为 GPU 一定会把这些 token 的 KV 算出来;唯一的例外是投机解码被拒绝的草稿,第 6 节会减回去。

4.3 分块:allocate_slots

图 8

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 的 ① ~ ⑦):

  1. 算这个请求一共需要多少个”槽位”(token 位置);
  2. 算出需要几块,扣掉手里已有的,得出还要新领几块;
  3. 空闲块不够?返回 None。调度器就是靠这个 None 决定抢占还是放弃;
  4. 前缀命中的块调用 block_pool.touch(),把引用计数加 1,防止被别人领走;
  5. 调 block_pool.get_new_blocks(n) 领新块;
  6. cache_blocks():把已经填满的块登记进前缀缓存;
  7. 返回这次新领的块,调度器把块号写进 SchedulerOutput 带给 Worker。

每个 KV 块多大?page_size_bytes = 2 × block_size × num_kv_heads × head_size × dtype 字节数,这是一层的大小,乘以层数才是一个块的总占用。block_size 默认 16。

4.4 块池:BlockPool

图 9

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 前缀缓存:链式哈希

图 10

哈希函数是 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

图 11

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 讲,原文也是按它讲的。

图 14

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:

图 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 写到哪

图 13

把 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 接着:

  1. 决定 CUDA Graph 模式和填充:比如这步有 10 个 token,就补齐到录好的 16,回放录制好的图(默认模式 FULL_AND_PIECEWISE:纯解码批用整图,混合批用分段图);
  2. _build_attention_metadata:把上面那些”说明书”打包成注意力后端需要的 FlashAttentionMetadata;
  3. set_forward_context + _model_forward:把元数据放进一个全局上下文(每层注意力从这里取),然后真正跑一遍 Transformer,得到每个位置的 hidden_states;
  4. hidden_states[logits_indices] → compute_logits:只在每段最后一个位置算 logits;
  5. 把 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)。

图 15

Sampler.forward(vllm/v1/sample/sampler.py)的顺序如图 15:

  1. 如果用户要 logprobs,先用原始 logits 算(在惩罚和温度之前,这点和老版本不同);
  2. 转成 float32,依次应用各种 logits 处理器:禁用词、logit_bias、min_tokens(还没到最小长度就把 EOS 屏蔽掉)、重复/频率/存在惩罚;
  3. 贪心请求直接 argmax,如果整批都是贪心,到这就返回了;
  4. 随机请求:除以 temperature → min_p → top_k / top_p 截断 → 按概率抽一个。

最后”按概率抽一个”有个小技巧:没有用 torch.multinomial,而是给每个词生成一个指数分布的随机数 q,取 probs / q 最大的那个。数学上等价于按概率抽样,但全程在 GPU 上完成,不需要 CPU 同步。同一批里贪心和随机的请求,最后用 torch.where 各取各的结果。


6. 收尾:记账和”该停了吗”

图 16

6.1 引擎侧:update_from_output

ModelRunnerOutput 回到引擎核心,Scheduler.update_from_output() 对每个请求:

  1. 投机解码回滚:如果本步验证了草稿,被拒绝了几个,就把 num_computed_tokens 减回去几个(4.2 节里那个”乐观地先加上”,在这里修正)。
  2. 追加新 token:request.append_output_token_ids(token)。
  3. 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
...
  1. 结束了就释放:_free_request() 把请求 id 记进 finished_req_ids(下一步告诉 Worker 删掉它),调用 kv_cache_manager.free() 还块。
  2. 打包:每个请求一个 EngineCoreOutput(request_id, new_token_ids, finish_reason, …),整批装进 EngineCoreOutputs,通过 ZMQ 发回前端。

6.2 前端侧:process_outputs

前端的 OutputProcessor.process_outputs()(vllm/v1/engine/output_processor.py)对每个输出:

  1. 增量反分词:detokenizer.update() 只把新增的 token 解成文字接在后面,不用每次重解整段;
  2. check_stop_strings:检查停止字符串,比如 stop=["\n\n"]。停止字符串可能横跨好几个 token,必须先有文字才能判断,所以放在前端。命中了就把多余的文字截掉,并通知引擎 abort 这个请求(LLMEngine.step 里的 abort_requests);
  3. make_request_output:离线的 FINAL_ONLY 模式下,没结束就返回 None;结束了才拼出 RequestOutput,里面的 outputs[0] 是一个 CompletionOutput,有 text、token_ids、finish_reason。

_run_engine 收齐所有结束的请求,按编号排好序,generate() 返回。到这里,我们的请求走完了图 1 的 ① ~ ⑧。


7. 在骨架上加功能

有了上面的骨架,高级功能就是在某几个环节上各加一点逻辑。

7.1 约束解码(structured outputs)

图 17

用法:

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)

图 18

以最简单的 n-gram 起草为例(vllm/v1/spec_decode/ngram_proposer.py):在已有的 token 里找和结尾”最长匹配”的片段,把它后面的 k 个 token 抄过来当猜测(图 18a)。不需要任何额外模型。

数据怎么流:

  1. sample_tokens 采完样后,runner 调 propose_draft_token_ids() 起草,存在自己身上;
  2. 引擎核心通过另一次 RPC take_draft_token_ids() 把草稿取回来,update_draft_token_ids() 放进 request.spec_token_ids;
  3. 下一次 schedule():num_tokens_with_spec 变大了 k,于是调度器自然给这个请求排 1 + k 个 token,草稿写进 scheduled_spec_decode_tokens;
  4. GPU 一次前向算出这 1 + k 个位置的概率,拒绝采样器(vllm/v1/sample/rejection_sampler.py)从左往右逐个判定(图 18b):大模型概率 p 除以草稿概率 q ≥ 随机数 u 就接受(n-gram 没有草稿概率,q 当作 1);一旦拒绝,就在这个位置按 max(p − q, 0) 重采一个,后面的草稿全部作废;如果全部接受,还白送一个 bonus token;
  5. update_from_output:num_computed_tokens −= 被拒绝的个数。

这套接受规则保证了输出分布和不开投机解码完全一样。除了 n-gram,spec_decode/ 目录下还有 EAGLE / EAGLE-3、Medusa、MTP、draft model、suffix decoding 等方法。

7.3 KV 连接器与 PD 分离

图 19

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

图 20

开 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 请求的来回

图 21

vllm serve 起的服务,和离线用法相比只是最外层换了:

  1. 路由:POST /v1/completions 进到 vllm/entrypoints/openai/completion/api_router.py 的 create_completion,再交给 OpenAIServingCompletion.create_completion。(原文里的 openai/api_server.py 现在只是一个兼容转发,真正的启动代码在 vllm/entrypoints/launchers/api_server/。)
  2. AsyncLLM.generate()(vllm/v1/engine/async_llm.py):LLMEngine 的异步版。分词、为这个请求建一个输出队列、发给引擎,然后在队列上 await,拿到一个就 yield 一个,这就是流式返回。
  3. 挑引擎:数据并行(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 快满的引擎,排队的请求消化得慢,新请求应该尽量避开。

  1. 引擎进程 EngineCoreProc:一个输入线程收包解码、放进队列;主循环 run_busy_loop 反复取请求、step();一个输出线程把结果编码发回。ZMQ 收发会释放 GIL,所以网络 IO 和 GPU 计算可以重叠。
  2. 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。高级功能和多卡、在线服务都可以等骨架熟了再看。


参考