LLMEngine 类概览
LLMEngine 是 vLLM 整个推理系统的核心枢纽,它统一管理 tokenizer、调度器、执行器、输出处理器和指标记录器,向上对接 LLM(离线批量推理)和 AsyncLLMEngine(在线服务),向下驱动分布式模型执行。
- 请求管理:接收、编码、排队请求,维护
SequenceGroup生命周期 - 迭代调度:每次
step()调用 Scheduler 决定本轮执行哪些序列 - 模型执行:将调度决策打包成
ExecuteModelRequest交给 Executor - 输出处理:将原始 token id 解码、检查停止条件、生成
RequestOutput - 指标监控:每步收集延迟、吞吐、KV Cache 使用率等指标推送给 Prometheus
初始化时创建的核心组件
__init__ 方法(L86–L242)按顺序初始化了以下关键成员:
# 1. 各类 Config 对象直接赋值保存
self.model_config = model_config # 模型配置
self.cache_config = cache_config # KV Cache 配置
self.scheduler_config = scheduler_config
self.parallel_config = parallel_config
self.speculative_config = speculative_config
self.decoding_config = decoding_config or DecodingConfig()
# 2. Tokenizer(支持 LoRA 多 tokenizer 组)
self._init_tokenizer()
self.detokenizer = Detokenizer(self.tokenizer) # L153
# 3. 序列 ID 全局计数器
self.seq_counter = Counter() # L158
# 4. 分布式执行器(根据设备/并行配置动态选择实现)
self.model_executor = executor_class(...) # L162
# 5. 初始化 KV Cache(由 executor profiling 决定块数)
self._initialize_kv_caches() # L174
# 6. 调度器
self.scheduler = Scheduler(scheduler_config, cache_config, lora_config) # L219
# 7. 指标记录器
self.stat_logger = StatLogger(
local_interval=_LOCAL_LOGGING_INTERVAL_SEC, # 5秒
labels=dict(model_name=model_config.served_model_name),
max_model_len=self.model_config.max_model_len) # L223
# 8. 输出处理器(单步 or 多步,取决于 speculative decoding)
self.output_processor = SequenceGroupOutputProcessor.create_output_processor(
self.scheduler_config, self.detokenizer, self.scheduler,
self.seq_counter, self.get_tokenizer_for_seq,
stop_checker=StopChecker(...)) # L231
这 8 个组件构成了 LLMEngine 的完整依赖图。其中 model_executor 是最重的一个——它在 from_engine_args()(L266–L300)中根据运行环境动态选择:单 GPU 用 GPUExecutor,多机多卡用 RayGPUExecutor,也支持 CPU 和 Neuron。
_initialize_kv_caches()(L244–L264)先让 executor 做 profiling 试跑来确定可用 GPU/CPU block 数,再调用 executor.initialize_cache() 实际分配内存。这个两步法是 vLLM PagedAttention 内存管理的基础——block 数直到运行时才确定,随后写回 cache_config.num_gpu_blocks。
step() 方法详解
step()(L556–L630)是引擎的"心跳"——每次调用代表一个推理迭代。它只有 25 行核心逻辑,却串联起整个推理管线:调度 → 执行 → 输出处理 → 指标记录。
@nvtx_tag
def step(self) -> List[RequestOutput]:
# ── 阶段 1:调度 ──────────────────────────────────────────
seq_group_metadata_list, scheduler_outputs = self.scheduler.schedule()
# scheduler_outputs 包含:
# scheduled_seq_groups 本轮参与执行的序列组列表
# ignored_seq_groups 超长被忽略的组
# blocks_to_swap_in/out 需要在 GPU/CPU 间搬运的 KV block
# blocks_to_copy prefix caching 需要复制的 block
# num_lookahead_slots 投机解码的前瞻槽数
# running_queue_size 当前 running 队列大小
# ── 阶段 2:执行(如果有任务)────────────────────────────
if not scheduler_outputs.is_empty():
execute_model_req = ExecuteModelRequest(
seq_group_metadata_list=seq_group_metadata_list,
blocks_to_swap_in=scheduler_outputs.blocks_to_swap_in,
blocks_to_swap_out=scheduler_outputs.blocks_to_swap_out,
blocks_to_copy=scheduler_outputs.blocks_to_copy,
num_lookahead_slots=scheduler_outputs.num_lookahead_slots,
running_queue_size=scheduler_outputs.running_queue_size,
)
output = self.model_executor.execute_model(
execute_model_req=execute_model_req)
else:
output = [] # 队列为空,跳过执行
# ── 阶段 3:处理输出 ──────────────────────────────────────
request_outputs = self._process_model_outputs(
output,
scheduler_outputs.scheduled_seq_groups,
scheduler_outputs.ignored_seq_groups,
seq_group_metadata_list)
# ── 阶段 4:记录指标 ──────────────────────────────────────
self.do_log_stats(scheduler_outputs, output)
return request_outputs
四个阶段逐一解析
阶段 1 — 调度
self.scheduler.schedule() 返回两个对象。seq_group_metadata_list 是每个序列组的元数据列表,包含 token id、block table、采样参数等,将直接喂给模型。scheduler_outputs 则包含调度决策的全部上下文,包括需要在 GPU/CPU 之间迁移的 KV block 映射表。
阶段 2 — 执行
is_empty() 检查本轮是否有实际工作。如果调度结果为空(所有请求都在等待且无法分配内存),则直接跳过模型执行,返回空列表。这是 vLLM 在高负载下的重要保护机制。
ExecuteModelRequest 把调度决策打包成一个统一的数据结构传给 executor。executor 负责处理 KV block 的 swap 操作并调用实际的模型 forward pass。
阶段 3 — 输出处理
_process_model_outputs() 接收原始的 SamplerOutput 列表,经过 detokenization、stop check、beam search fork/free 等操作,最终生成客户端可读的 RequestOutput 列表。
阶段 4 — 指标记录
do_log_stats() 每步都会调用,向 Prometheus 推送 gauge/counter/histogram,并在每 5 秒(_LOCAL_LOGGING_INTERVAL_SEC,L38)向 stdout 打印吞吐量摘要。
step() 和多个内部方法都用 @nvtx_tag 标注。这是 vLLM 内部封装的 NVIDIA NVTX 性能分析注解——在 Nsight Systems 等工具中,每次 step() 调用都会在 timeline 上形成一个独立的命名区间,便于定位性能瓶颈。
add_request() — 请求如何进入引擎
add_request()(L357–L466)是请求进入引擎的唯一入口。它完成从原始字符串到 SequenceGroup 的完整构建,然后交给调度器排队。
@nvtx_tag
def add_request(
self,
request_id: str,
prompt: Optional[str],
sampling_params: SamplingParams,
prompt_token_ids: Optional[List[int]] = None,
arrival_time: Optional[float] = None,
lora_request: Optional[LoRARequest] = None,
multi_modal_data: Optional[MultiModalData] = None,
) -> None:
# 1. 校验 logprobs 上限
if arrival_time is None:
arrival_time = time.time()
# 2. 编码 prompt → token ids(如果未提供)
prompt_token_ids = self.encode_request(
request_id=request_id,
prompt=prompt,
prompt_token_ids=prompt_token_ids,
lora_request=lora_request)
# 3. 构建 Sequence 对象
block_size = self.cache_config.block_size
seq_id = next(self.seq_counter) # 全局递增的序列 ID
eos_token_id = self.tokenizer.get_lora_tokenizer(lora_request).eos_token_id
seq = Sequence(seq_id, prompt, prompt_token_ids, block_size,
eos_token_id, lora_request)
# 4. 克隆 SamplingParams(防止调用方修改影响引擎内部状态)
sampling_params = sampling_params.clone()
sampling_params.all_stop_token_ids.add(seq.eos_token_id)
sampling_params.update_from_generation_config(self.generation_config_fields)
# 5. 构建 SequenceGroup(一个请求对应一个 SequenceGroup)
seq_group = SequenceGroup(request_id, [seq], sampling_params,
arrival_time, lora_request, multi_modal_data)
# 6. 提交给调度器的 waiting 队列
self.scheduler.add_seq_group(seq_group)
关键设计细节
为什么要 clone SamplingParams?
L450 的 sampling_params.clone() 是防御性复制。调用方传入的 SamplingParams 对象可能在引擎执行期间被外部修改(尤其是在异步服务场景),克隆一份确保引擎内部状态隔离。注释也特别说明这不深拷贝 LogitsProcessor 对象——因为那些通常是无状态的函数。
EOS Token 的处理
L453–L454 把 eos_token_id 加入 all_stop_token_ids。这个集合在生成时被 stop checker 使用,支持 min_tokens 参数——在达到最小生成长度之前,即使采样出 EOS token 也不会停止。此外 L438–L443 还支持通过环境变量 MANUAL_EOS_TOKEN 手动覆盖 EOS token id,这是一个用于调试/测试的扩展点。
SequenceGroup vs Sequence
一个请求对应一个 SequenceGroup,初始时只包含一个 Sequence(L459)。当使用 beam search 时,调度器后续会 fork 出多个 Sequence——这就是 best_of 参数的实现基础。
_process_model_outputs() — 模型输出如何被处理
_process_model_outputs()(L503–L553)接收模型执行的原始结果 List[SamplerOutput],完成三件事:更新序列状态、释放已完成的序列组、生成客户端结果。
@nvtx_tag
def _process_model_outputs(
self,
output: List[SamplerOutput],
scheduled_seq_groups: List[ScheduledSequenceGroup],
ignored_seq_groups: List[SequenceGroup],
seq_group_metadata_list: List[SequenceGroupMetadata],
) -> List[RequestOutput]:
now = time.time()
# ── 步骤 1:重组输出维度 ──────────────────────────────────
# 模型输出格式是 [step][seq_group],重组为 [seq_group][step]
# 这对于多步执行(speculative decoding)很重要
output_by_sequence_group = create_output_by_sequence_group(
sampler_outputs=output,
num_seq_groups=len(scheduled_seq_groups))
# ── 步骤 2:逐组处理输出 ──────────────────────────────────
for scheduled_seq_group, outputs, seq_group_meta in zip(
scheduled_seq_groups, output_by_sequence_group,
seq_group_metadata_list):
seq_group = scheduled_seq_group.seq_group
# 更新本组已计算的 token 数(用于 chunked prefill 追踪)
seq_group.update_num_computed_tokens(
scheduled_seq_group.token_chunk_size)
# 处理 prompt logprobs(无论是否采样都要更新)
self.output_processor.process_prompt_logprob(seq_group, outputs)
# 只有 do_sample=True 时才处理生成的 token
if seq_group_meta.do_sample:
self.output_processor.process_outputs(seq_group, outputs)
# ── 步骤 3:释放已完成的序列组 ───────────────────────────
self.scheduler.free_finished_seq_groups()
# ── 步骤 4:构建 RequestOutput ───────────────────────────
request_outputs: List[RequestOutput] = []
for scheduled_seq_group in scheduled_seq_groups:
seq_group = scheduled_seq_group.seq_group
seq_group.maybe_set_first_token_time(now) # 记录 TTFT 时间戳
request_output = RequestOutput.from_seq_group(seq_group)
request_outputs.append(request_output)
# 被忽略的组(超过 max_model_len)也需要返回一个错误输出
for seq_group in ignored_seq_groups:
request_output = RequestOutput.from_seq_group(seq_group)
request_outputs.append(request_output)
return request_outputs
关键细节解析
输出维度重组
create_output_by_sequence_group()(来自 vllm/engine/output_processor/util.py)将模型输出从 [step × seq_group] 转置为 [seq_group × step]。对于普通单步解码,每个 SamplerOutput 只有一步,所以转置结果是 [seq_group][1]。但对于投机解码(speculative decoding),一次 forward pass 可能生成多个 token,此时步骤数 > 1,转置后每个序列组的 outputs 列表就包含多步结果。
do_sample 的条件判断
L534 检查 seq_group_meta.do_sample。在 chunked prefill 场景下,一个序列组可能还处于 prefill 阶段(还没到采样步骤),此时 do_sample=False,只需要更新 prompt logprobs,不需要 append 新 token。这个判断确保了 chunked prefill 与正常 decode 的统一处理路径。
TTFT 时间戳
maybe_set_first_token_time(now) 在每个序列组生成第一个 token 时记录时间戳。这个时间戳后续在 _get_stats() 中被用来计算 TTFT(Time To First Token),是推理服务最关键的延迟指标之一。
当请求的 prompt 长度超过 max_model_len 时,调度器会把它放入 ignored_seq_groups 而不是 scheduled_seq_groups。_process_model_outputs() 依然会为这些组调用 RequestOutput.from_seq_group() 并返回给调用方——调用方会看到一个标记为 finished(错误状态)的输出,从而可以向用户报告错误。
输出处理管线 — SequenceGroupOutputProcessor 的角色
SequenceGroupOutputProcessor(vllm/engine/output_processor/interfaces.py)是输出处理的抽象层,将"如何处理模型输出"这一复杂逻辑从 LLMEngine 中剥离出来,允许针对不同解码策略提供不同实现。
class SequenceGroupOutputProcessor(ABC):
"""处理序列组中新 token id 的接口:
管理 detokenization、stop checking 以及 scheduler 中的序列 fork/free。
与 LLMEngine 高度耦合,可视为 LLMEngine 的功能性扩展,
分离出来是为了简化 LLMEngine 并支持不同解码策略的独立实现。
"""
@staticmethod
def create_output_processor(
scheduler_config: SchedulerConfig,
detokenizer: Detokenizer,
scheduler: Scheduler,
seq_counter: Counter,
get_tokenizer_for_seq: Callable[[Sequence], PreTrainedTokenizer],
stop_checker: "StopChecker",
):
# 根据 num_lookahead_slots 决定实现类型
if scheduler_config.num_lookahead_slots == 0:
return SingleStepOutputProcessor(...) # 普通解码
else:
return MultiStepOutputProcessor(...) # 投机解码
@abstractmethod
def process_outputs(self, sequence_group: SequenceGroup,
outputs: List[SequenceGroupOutput]) -> None:
"""处理新 token id:detokenize、stop check、fork/free 序列。"""
@abstractmethod
def process_prompt_logprob(self, seq_group: SequenceGroup,
outputs: List[SequenceGroupOutput]) -> None:
"""将 outputs 中的 prompt logprobs 更新到 seq_group。"""
两种实现的职责划分
路径:vllm/engine/output_processor/single_step.py。用于标准自回归解码(含 beam search)。process_outputs() 负责:
- 从
SamplerOutput中提取每个序列的 token id - 调用
Detokenizer增量解码新 token - 调用
StopChecker判断是否满足停止条件(EOS、停止词、最大长度) - 对 beam search:fork 概率高的 beam,free 已淘汰的 beam
- 将完成的序列状态标记为 FINISHED,触发后续的 KV block 回收
路径:vllm/engine/output_processor/multi_step.py。用于投机解码(speculative decoding)。一次 forward pass 可能输出多个 token(草稿模型提出 + 验证模型接受),process_outputs() 需要循环处理每一步的结果,并在遇到被拒绝的 token 时截断序列。不支持 beam search。
工厂方法的决策依据
工厂方法 create_output_processor() 通过 scheduler_config.num_lookahead_slots 来判断:如果是 0,则使用单步处理器;否则使用多步处理器。num_lookahead_slots 正是投机解码中草稿模型前瞻的 token 槽数量。这个设计将"解码策略的复杂性"完全封装在处理器内部,LLMEngine.step() 无需感知当前使用何种解码策略。
self.output_processor = (
SequenceGroupOutputProcessor.create_output_processor(
self.scheduler_config,
self.detokenizer,
self.scheduler,
self.seq_counter,
self.get_tokenizer_for_seq,
stop_checker=StopChecker(
self.scheduler_config.max_model_len,
self.get_tokenizer_for_seq,
),
))
统计与指标 — StatLogger 的监控机制
vLLM 的监控系统分为两层:数据采集层(_get_stats())和数据上报层(StatLogger.log())。每次 step() 都会触发完整的采集+上报流程。
_get_stats():采集当前迭代的全量数据
def _get_stats(self, scheduler_outputs, model_output=None) -> Stats:
# ── 类型 1:系统状态(System Stats)──────────────────────
num_running_sys = len(self.scheduler.running) # 正在 GPU 上运行的请求数
num_swapped_sys = len(self.scheduler.swapped) # 已换出到 CPU 的请求数
num_waiting_sys = len(self.scheduler.waiting) # 等待队列中的请求数
# KV Cache 使用率
num_free_gpu = self.scheduler.block_manager.get_num_free_gpu_blocks()
gpu_cache_usage_sys = 1.0 - (num_free_gpu / num_total_gpu)
# ── 类型 2:迭代级指标(Iteration Stats)────────────────
# 遍历 scheduled_seq_groups,区分 prefill 组和 decode 组:
for idx, scheduled_seq_group in enumerate(scheduled_seq_groups):
group_was_prefill = idx < scheduler_outputs.num_prefill_groups
if group_was_prefill:
num_prompt_tokens_iter += scheduled_seq_group.token_chunk_size
if not seq_group.is_prefill(): # prefill 刚结束 → 记录 TTFT
latency = seq_group.get_last_latency(now)
time_to_first_tokens_iter.append(latency)
else:
latency = seq_group.get_last_latency(now) # decode → 记录 TPOT
time_per_output_tokens_iter.append(latency)
# ── 类型 3:请求级指标(Request Stats,仅在请求完成时记录)
if seq_group.is_finished():
time_e2e_requests.append(now - seq_group.metrics.arrival_time)
num_prompt_tokens_requests.append(len(seq_group.prompt_token_ids))
# ... best_of, n, finished_reason 等
return Stats(now=now, num_running_sys=..., gpu_cache_usage_sys=..., ...)
StatLogger.log():双通道上报
def log(self, stats: Stats) -> None:
# 通道 1:每步上报到 Prometheus(实时)
self._log_prometheus(stats)
# 累积 token 计数用于吞吐量计算
self.num_prompt_tokens.append(stats.num_prompt_tokens_iter)
self.num_generation_tokens.append(stats.num_generation_tokens_iter)
# 通道 2:每 5 秒上报到 stdout(人工可读摘要)
if self._local_interval_elapsed(stats.now):
prompt_throughput = self._get_throughput(self.num_prompt_tokens, stats.now)
generation_throughput = self._get_throughput(self.num_generation_tokens, stats.now)
logger.info(
"Avg prompt throughput: %.1f tokens/s, "
"Avg generation throughput: %.1f tokens/s, "
"Running: %d reqs, Swapped: %d reqs, "
"Pending: %d reqs, GPU KV cache usage: %.1f%%, "
"CPU KV cache usage: %.1f%%",
prompt_throughput, generation_throughput,
stats.num_running_sys, stats.num_swapped_sys,
stats.num_waiting_sys,
stats.gpu_cache_usage_sys * 100,
stats.cpu_cache_usage_sys * 100,
)
# 重置累积计数
self.num_prompt_tokens = []
self.num_generation_tokens = []
self.last_local_log = stats.now
Prometheus 指标体系
Metrics 类(engine/metrics.py L25–L138)定义了完整的 Prometheus 指标集,覆盖三个维度:
| 指标名 | 类型 | 含义 |
|---|---|---|
vllm:num_requests_running |
Gauge | 当前在 GPU 上执行的请求数 |
vllm:num_requests_waiting |
Gauge | 等待队列中的请求数 |
vllm:num_requests_swapped |
Gauge | 已换出到 CPU 的请求数 |
vllm:gpu_cache_usage_perc |
Gauge | GPU KV Cache 使用率(0~1) |
vllm:prompt_tokens_total |
Counter | 累计处理的 prefill token 数 |
vllm:generation_tokens_total |
Counter | 累计生成的 decode token 数 |
vllm:time_to_first_token_seconds |
Histogram | 首 token 延迟(TTFT)分布 |
vllm:time_per_output_token_seconds |
Histogram | 每 token 延迟(TPOT)分布 |
vllm:e2e_request_latency_seconds |
Histogram | 端到端请求延迟分布 |
vllm:request_success_total |
Counter | 按完成原因统计的成功请求数 |
TTFT(Time to First Token)在 prefill 恰好完成的那一步采集:group_was_prefill and not seq_group.is_prefill()(L714)。TPOT(Time Per Output Token)在纯 decode 步骤采集(L722)。这两者共同决定了用户感受到的响应速度。
step() 的整体数据流图
下图展示了一次完整的 step() 调用中,数据如何在各组件间流动,以及每个阶段的关键数据结构变换。
数据结构变换总结
| 阶段 | 输入 | 输出 | 关键变换 |
|---|---|---|---|
| add_request | str prompt | SequenceGroup in waiting | tokenize → Sequence → SequenceGroup |
| schedule() | waiting/running/swapped 队列 | metadata_list + SchedulerOutputs | 调度决策 + KV block 分配 |
| execute_model() | ExecuteModelRequest | List[SamplerOutput] | GPU forward pass + sampling |
| create_output_by_sequence_group() | List[SamplerOutput] (step×group) | List[List[SequenceGroupOutput]] (group×step) | 维度转置 |
| process_outputs() | SequenceGroup + outputs | Sequence 状态更新(RUNNING/FINISHED) | detokenize + stop check + fork/free |
| RequestOutput.from_seq_group() | SequenceGroup | RequestOutput | 提取文本 + 完成状态 + logprobs |
step() 的设计体现了 vLLM 的核心理念:iteration-level scheduling(迭代级调度)。每次 step() 都是一个完整的决策-执行-反馈循环,调度器在每步重新评估所有请求的状态,动态决定哪些序列参与本轮 forward pass。这使得 vLLM 能够实现:
- 请求的动态抢占(preemption):高优先级请求可以抢占低优先级请求的 KV block
- 连续批处理(continuous batching):新请求无需等待当前批次完成即可加入
- KV Cache 的精细化管理:通过 PagedAttention + block 级 swap,最大化 GPU 显存利用率