上一篇讲到路径 D 的 build_async_engine_client 一路 spawn 出 EngineCore 子进程,EngineCore.init 里 self.model_executor = executor_class(vllm_config) 实例化了 MultiprocExecutor。
本篇从这里开始,拆开引擎侧的进程/类关系,
重点分析 vLLM 0.25 默认启用的 ModelRunner V2(MRV2)。
代码根目录:/usr/local/lib/python3.12/dist-packages/vllm
示例场景:Qwen3-32B,TP=4,--distributed-executor-backend mp。
1. 涉及的核心类与职责
进程拓扑(TP=4, DP=1):
主进程 APIServer(pid=5147)└─ EngineCore 子进程(pid=5386)└─ MultiprocExecutor├─ Worker 子进程(pid=5596, TP rank0, driver)├─ Worker 子进程(pid=5597, TP rank1)├─ Worker 子进程(pid=5598, TP rank2)└─ Worker 子进程(pid=5599, TP rank3)
2. MultiprocExecutor._init_executor(multiproc_executor.py:110)
def _init_executor(self) -> None:tp, pp, pcp = self._get_parallel_sizes()# TP=4, PP=1, PCP=1assert world_size == tp * pp * pcp# world_size=4set_multiprocessing_worker_envs()# 1) 建广播消息队列:把 SchedulerOutput 广播给所有 Worker(共享内存 MessageQueue)self.rpc_broadcast_mq = MessageQueue(world_size, local_world_size, ...)scheduler_output_handle = self.rpc_broadcast_mq.export_handle()# 2) 逐个 rank spawn Worker 进程context = get_mp_context()# spawnfor local_rank in range(self.local_world_size):# 0..3global_rank = global_start_rank + local_rankunready = WorkerProc.make_worker_process(vllm_config, local_rank, global_rank,distributed_init_method,# tcp://127.0.0.1:portinput_shm_handle=scheduler_output_handle,# 传广播队列句柄is_driver_worker=self._is_driver_worker(global_rank),)unready_workers.append(unready)# 3) 等所有 Worker 就绪# (worker.init_device 里有 device sync,必须先全部 spawn 再统一等,否则死锁)self.workers = WorkerProc.wait_for_ready(unready_workers)# 4) 起后台线程监控 Worker 健康self.start_worker_monitor()# 5) 收集每个 Worker 的 response 队列(收 ModelRunnerOutput)self.response_mqs = [...]
要点:
**进程间通信基于共享内存 **MessageQueue:rpc_broadcast_mq(executor → workers 广播 SchedulerOutput)+ 每个 worker 的 worker_response_mq(worker → executor 回传输出)。 TP=4 → local_world_size=4 → spawn 4 个 Worker 进程(日志里的 pid 5596~5599)。 rank 0 为 driver worker。 先全部 spawn 再统一 wait_for_ready,是因为 init_device() 里有跨卡 device sync,必须并行起来否则互相等待死锁。
3. WorkerProc.make_worker_process:spawn 进程(multiproc_executor.py:662)
ready_reader, ready_writer = context.Pipe(duplex=False)# 子->父 就绪信号death_reader, death_writer = context.Pipe(duplex=False)# 检测父进程退出proc = context.Process(target=WorkerProc.worker_main,# 子进程入口kwargs={...},)proc.start()# ★ spawn Worker 子进程
spawn 出来的子进程执行 WorkerProc.worker_main(:810)→ worker = WorkerProc(*args) → 进入 WorkerProc.init。
4. WorkerProc.init:worker对应哪个类,里面做什么(multiproc_executor.py:597)
关键三步:
# ① 实例化真正的 Workerwrapper = WorkerWrapperBase(rpc_rank=local_rank, global_rank=rank)wrapper.init_worker(all_kwargs)# 按 worker_cls 动态 import 并实例化self.worker = wrapper# ② 绑 GPU + 初始化 NCCL + 构建 model_runnerself.worker.init_device()# ③ 加载模型权重self.worker.load_model()# 建收/发消息队列(收 SchedulerOutput / 发 ModelRunnerOutput)self._init_message_queues(input_shm_handle, vllm_config)
worker 对应的类:init_worker 根据 parallel_config.worker_cls 动态实例化。CUDA 平台在 platforms/cuda.py:312 把它设为:
parallel_config.worker_cls = ”vllm.v1.worker.gpu_worker.Worker”所以worker 类就是vllm.v1.worker.gpu_worker.Worker(继承WorkerBase)。WorkerWrapperBase 只是延迟实例化的壳。
5. Worker.init_device:分布式环境 + 构建 Model Runner(gpu_worker.py:297)
self.device = torch.device(f”cuda:{visible_device_index}”)torch.accelerator.set_device_index(self.device)# 初始化分布式(NCCL)——日志 ”world_size=4 rank=0 ... backend=nccl”init_worker_distributed_environment(vllm_config, rank, distributed_init_method, local_rank, ”nccl”)if self.use_v2_model_runner:logger.info_once(”Using V2 Model Runner”)# 日志里那行的出处# ★ 构建 model runner,V2 / V1 二选一if self.use_v2_model_runner:from vllm.v1.worker.gpu.model_runner import GPUModelRunner as GPUModelRunnerV2self.model_runner = GPUModelRunnerV2(vllm_config, device)else:from vllm.v1.worker.gpu_model_runner import GPUModelRunner as GPUModelRunnerV1self.model_runner = GPUModelRunnerV1(vllm_config, device)
即:先建 CUDA device 和 NCCL 通信组,再根据 use_v2_model_runner 决定实例化哪个 Model Runner。
6. 是否使用MRV2:use_v2_model_runner 判定(config/vllm.py:550)
@propertydef use_v2_model_runner(self) -> bool:if envs.VLLM_USE_V2_MODEL_RUNNER is not None:# 环境变量强制优先return envs.VLLM_USE_V2_MODEL_RUNNERif spec_method == ”dspark” or dflash_multi_kv or is_diffusion:return True# 这几类只有 V2 支持if not self._is_default_v2_model_runner_model():# generate + 非hybrid + 非attn_free +return False# (arch 在白名单 或 非 MoE)if not HAS_TRITON:# V2 依赖 Tritonreturn Falseif self._get_v2_model_runner_unsupported_features():# 有不支持特性则回退 V1return Falsereturn True
_is_default_v2_model_runner_model(config/vllm.py:605):runner_type == "generate" 且非 hybrid、非 attention_free,且(架构在 DEFAULT_V2_MODEL_RUNNER_ARCHITECTURES 白名单里或非 MoE)。
本场景结论:Qwen3-32B 是 dense(非 MoE)、generate、非 hybrid,Triton 可用,无不支持特性 → use_v2_model_runner = True,所以日志打印 Using V2 Model Runner。可用 VLLM_USE_V2_MODEL_RUNNER=0 强制回退 V1。
7. 为什么 vLLM 0.25 引入 Model Runner V2
v1/worker/gpu/README.md 直言这是 "Model Runner V2, under active development"。V2 的核心思想:把原来gpu_model_runner.py(V1,单文件近 7000 行的巨型 class)拆成一个极简、稳定、所有模型共享的主 runner + 一批职责单一的工具模块。
gpu/model_runner.py 文件头的编码约束(原文大意):
这个 model runner 被所有模型共享(文本/多模态、生成/embedding、公开/私有),所以只能放通用代码; 对改动要"偏执",对新增行要"更偏执",保持最小、稳定; 特性逻辑(哪怕是并行模式这种共享特性)一律外置到工具函数,越不常用越要藏起来。
模块拆分(v1/worker/gpu/):
对比 V1 的单体 gpu_model_runner.py,V2 更易维护、更易按模型类型扩展。
8. GPUModelRunner(V2) 的构建与 load_model
8.1 init(gpu/model_runner.py:122)预建通用组件
从 vllm_config 摊平各 sub-config,并预建:
RequestState(states.py):请求状态(max_num_reqs / max_model_len 等) InputBuffers(input_batch.py):输入缓冲池 PP 时的 PPHandler EPLBController 多模态时的 EncoderCache 投机解码时的 speculator(init_speculator)
采样相关(Sampler / cudagraph_manager 等)延迟到load_model之后,因为它们依赖 model_state。
8.2 load_model(gpu/model_runner.py:276)
Worker.load_model(gpu_worker.py:424)会调用它:
model_loader = get_model_loader(load_config)self.model = model_loader.load_model(vllm_config, model_config)# 加载 safetensors 权重# 日志: ”Model loading took 15.52 GiB and 7.10 seconds”# 按模型类型建状态(generate / pooling / 多模态等各自的 model_state)self.model_state = init_model_state(vllm_config, self.model, encoder_cache, device)# 末位 PP rank 上建采样相关组件if self.is_last_pp_rank and not self.is_pooling_model:self.sampler = Sampler(...)if speculative_config:self.rejection_sampler = RejectionSampler(...)self.prompt_logprobs_worker = PromptLogprobsWorker(...)self.structured_outputs_worker = StructuredOutputsWorker(...)
8.3 后续阶段(后文展开)
load_model 之后由 EngineCore 侧驱动:
get_kv_cache_spec(:408)→ 汇报 KV cache 规格 initialize_kv_cache(:411)→ 分配 KV cache 显存 profile_run(:673)→ dummy forward 测峰值显存,算出可用 KV cache capture_model(:710)→ 捕获 CUDA Graph(PIECEWISE / FULL) 推理期:execute_model(:1151)+ sample_tokens(:1395)
9. 一句话链路总结
EngineCore.__init__: self.model_executor = MultiprocExecutor(vllm_config)└ MultiprocExecutor._init_executor (multiproc_executor.py:110)建 rpc_broadcast_mq(广播 SchedulerOutput)×4 WorkerProc.make_worker_process (:662) → spawn Worker 子进程└ WorkerProc.worker_main (:810)└ WorkerProc.__init__ (:597)① WorkerWrapperBase.init_worker → Worker(gpu_worker.py:130)② Worker.init_device (gpu_worker.py:297)NCCL init + 判 use_v2_model_runner→ GPUModelRunner V2 (gpu/model_runner.py:121)③ Worker.load_model → model_runner.load_model (gpu/model_runner.py:276)加载权重 + init_model_state + SamplerWorkerProc.wait_for_ready + start_worker_monitor
10. 流程图

夜雨聆风