核心组件主要是LLMEngine,Executor和Worker,我们从最底层逐渐向上阅读。
学习这部分的初衷是希望捋清楚以后能用ray来分别起多个进程,然后它们再组成model parallel group,我将这部分折叠起来,作为一个修改vLLM代码的参考。
用ray独立起进程的好处是可以使用使用分数的num_gpus,这使得不同的模型能同时放在一个设备上。
这个功能对于我正在实现的一个易用强化学习框架十分关键。
用ray启动进程的好处是能非常好地隔离不同进程之间的资源,尤其是通信组的设置和PyTorch的CUDA runtime状态。
上一篇blog分析了整体的结构,所以我们知道要修改分布式启动方式,只需要提供一个新的Executor类即可。
这个类只需要根据传入的环境变量来初始化模型的通信组,让参数和数据同步正常进行即可。
vllm/distributed
1
2
3
4
5
6
7
8
9
| vllm/distributed/
├── 📁 device_communicators/
├── 📁 kv_transfer/
├── 📄 communication_op.py
├── 📄 __init__.py
├── 📄 parallel_state.py
└── 📄 utils.py
3 directories, 4 files
|
vllm/distributed/device_communicators
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
| vllm/distributed/device_communicators/
├── 📄 base_device_communicator.py
├── 📄 cpu_communicator.py
├── 📄 cuda_communicator.py
├── 📄 cuda_wrapper.py
├── 📄 custom_all_reduce.py
├── 📄 custom_all_reduce_utils.py
├── 📄 hpu_communicator.py
├── 📄 __init__.py
├── 📄 neuron_communicator.py
├── 📄 pynccl.py
├── 📄 pynccl_wrapper.py
├── 📄 shm_broadcast.py
├── 📄 tpu_communicator.py
└── 📄 xpu_communicator.py
1 directory, 14 files
|
🔥 以下的文件和Nvidia GPU高度相关:
base_device_communicator.py用torch.distributed的对应函数实现通信卡间通信cuda_communicator.py对于tensor parallel采用自定义的通信算子,其对单个物理节点内的通信做了优化。对于其他情形使用pynccl的通信custom_all_reduce.py实现了自定义通信算子,其对NVLink和CUDA graph有专门的支持pynccl.py实现了基于pynccl的通信,通过 ctypes 调用 NCCL 库,避免了 PyTorch 分布式层的额外开销,同时对CUDA graph的支持也更好
💡 对动态链接库的封装文件我在这里暂时不深入阅读:
cuda_wrapper.pypynccl_wrapper.py
💡 其他的芯片上的通信类主要是载入对应的框架,并使用框架提供的all_reduce等接口:
hpu_communicator.py实现了Habana Gaudi处理器间的通信neuron_communicator.py实现了AWS Inferentia/Trainium处理器间的通信tpu_communicator.py实现了Google TPU处理器之间的通信xpu_communicator.py实现了Intel XPU处理器之间的通信
另外,shm_broadcaster.py实现了利用共享内存做进程间通信的功能,并且在MultiprocExecutor和GroupCooridnator中用于worker进程之间的通信。
vllm/distributed/kv_transfer
这里实现了Prefill-Decode分离,仍然是实验状态,留待以后来看。
vllm/distributed/parallel_state.py
GroupCoordinator类
这个类包装了当前进程的两个通信组,device_group和cpu_group,使得不同的通信操作可以调用不同的组。
这是为了解决device_group在某些情况下的不稳定,而cpu_group在很多场景下性能不够的问题。
此类记录了当前通讯组里的所有进程的rank,相关的变量是
rank,当前进程的global rankranks,当前组中所有进程的global ranklocal_rank,在实际的物理节点上的rankrank_in_group,在本组中的rank
代码中的例子:
| Process | Node | Rank | Local Rank | Rank in Group |
|---|
| 0 | 0 | 0 | 0 | 0 |
| 1 | 0 | 1 | 1 | 1 |
| 2 | 1 | 2 | 0 | 2 |
| 3 | 1 | 3 | 1 | 3 |
这个类有三种类型的方法,分别是取组里的某个rank(首尾和相邻的),通信用方法(如all_gather等),和graph_capture。
通信方法中比较有意思的地方是:
all_gather, gather、all_reduce、send、recv用的是device_communicator的对应方法,而broadcast、send_object、recv_object用的都是PyTorch原生的。send_object、recv_object和barrier都是走的cpu_group。尤其是barrier:
NOTE: Don’t use device_group here! barrier in NCCL is
terrible because it is internally a broadcast operation with
secretly created GPU tensors. It is easy to mess up the current
device. Use the CPU group instead.
使用ray起需要注意的问题
这段代码会有问题:
1
2
3
4
| if current_platform.is_cuda_alike():
self.device = torch.device(f"cuda:{local_rank}")
else:
self.device = torch.device("cpu")
|
因为ray起的进程只能看到一个GPU,因此这里local_rank不能大于0。
同时,通信的时候容易因为多个rank都是cuda:0而报错(例如pynccl)。
其他还有一些初始化通信组的函数,目的都是创建一个GroupCoordinator
init_world_group函数
所有vllm进程参与创建一个GroupCoordinator,赋值给一个全局的变量WORLD_,其中device_communicator=False。
init_model_parallel_group函数
返回一个GroupCoordinator类,使用device_communicator=True,并且支持使用message_queue_broadcaster。
Tensor parallel、pipeline parallel和data parallel组分别是_TP,_PP,_DP全局变量。
init_distributed_backend函数
从全局的配置parallel_config中取rank等信息,用torch.distributed.init_process_group初始化PyTorch的分布式通信组,最后初始化WORLD_。
initialize_model_parallel函数
初始化_TP、_PP和_DP。这里rank的划分确实挺简洁的,将不同的parallelism用4D的rank tensor来划分:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
| # First dimension in this 4D tensor is for `verl` compatibility
all_ranks = torch.arange(world_size).reshape(
-1,
data_parallel_size,
pipeline_model_parallel_size,
tensor_model_parallel_size
)
# For _TP
group_ranks = all_ranks.view(-1, tensor_model_parallel_size).unbind(0)
# For _PP
group_ranks = all_ranks.transpose(2, 3).reshape(
-1, pipeline_model_parallel_size
).unbind(0)
# For _DP
group_ranks = all_ranks.transpose(1, 3).reshape(
-1, data_parallel_size
).unbind(0)
|
vllm/model_executor
1
2
3
4
5
6
7
8
9
10
11
12
13
| vllm/model_executor/
├── 📁 guided_decoding/
├── 📁 layers/
├── 📁 model_loader/
├── 📁 models/
├── 📄 custom_op.py
├── 📄 __init__.py
├── 📄 parameter.py
├── 📄 pooling_metadata.py
├── 📄 sampling_metadata.py
└── 📄 utils.py
5 directories, 6 files
|
vllm/model_executor/guided_decoding
这个文件夹下包含了4个结构化输出引擎支持的logits processor,即llguidance、lmformatenforcer、outlines和xgrammar。
一个logits processor就是一个callable,输入是input_ids和logits,输出是处理后的logits。
1
2
3
4
5
6
7
8
9
10
11
12
13
14
| vllm/model_executor/guided_decoding/
├── 📁 reasoner/
│ └── 📄 __init__.py
├── 📄 guidance_decoding.py
├── 📄 guidance_logits_processors.py
├── 📄 guided_fields.py
├── 📄 __init__.py
├── 📄 lm_format_enforcer_decoding.py
├── 📄 outlines_decoding.py
├── 📄 outlines_logits_processors.py
├── 📄 utils.py
└── 📄 xgrammar_decoding.py
2 directories, 10 files
|
vllm/model_executor/model_loader
1
2
3
4
5
6
7
8
9
| vllm/model_executor/model_loader/
├── 📄 __init__.py
├── 📄 loader.py
├── 📄 neuron.py
├── 📄 tensorizer.py
├── 📄 utils.py
└── 📄 weight_utils.py
1 directory, 6 files
|
此目录下核心文件为loader.py,实现了各种模型的加载器。
模型加载器只有两个接口:download_model和load_model。
vllm实现了7种模型加载器(但是同时也支持自定义的加载器,传给load_config中的load_format即可),分别是
DummyModelLoader随机初始化模型参数。这里我看实现用了with torch.device(device)上下文,需要注意这在device="meta"时会有buffer没有被初始化的问题,需要显式调用init_buffers方法。TensorizerLoader利用tensorizer库对模型做序列化和反序列化,load_model就是反序列化一个序列化后的模型。ShardedStateLoader是用于加载分片模型状态的加载器,支持大模型的并行加载。BitsAndBytesModelLoader支持 BitAndBytes 量化模型的加载器。GGUFModelLoader专门用于加载 GGUF 格式模型的加载器。RunaiModelStreamerLoader是用于流式加载模型的加载器。DefaultModelLoader最后fallback的加载器,每个进程独立加载完整的参数,支持各种存储的格式,例如safetensors等。
vllm/model_executor/models
完整的文件列表太长因此折叠起来
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
| vllm/model_executor/models/
├── 📄 adapters.py
├── 📄 arctic.py
├── 📄 aria.py
├── 📄 baichuan.py
├── 📄 bamba.py
├── 📄 bart.py
├── 📄 bert.py
├── 📄 blip2.py
├── 📄 blip.py
├── 📄 bloom.py
├── 📄 chameleon.py
├── 📄 chatglm.py
├── 📄 clip.py
├── 📄 commandr.py
├── 📄 dbrx.py
├── 📄 decilm.py
├── 📄 deepseek_mtp.py
├── 📄 deepseek.py
├── 📄 deepseek_v2.py
├── 📄 deepseek_vl2.py
├── 📄 eagle.py
├── 📄 exaone.py
├── 📄 fairseq2_llama.py
├── 📄 falcon.py
├── 📄 florence2.py
├── 📄 fuyu.py
├── 📄 gemma2.py
├── 📄 gemma3_mm.py
├── 📄 gemma3.py
├── 📄 gemma.py
├── 📄 glm4v.py
├── 📄 glm.py
├── 📄 gpt2.py
├── 📄 gpt_bigcode.py
├── 📄 gpt_j.py
├── 📄 gpt_neox.py
├── 📄 granitemoe.py
├── 📄 granitemoeshared.py
├── 📄 granite.py
├── 📄 gritlm.py
├── 📄 grok1.py
├── 📄 h2ovl.py
├── 📄 idefics2_vision_model.py
├── 📄 idefics3.py
├── 📄 __init__.py
├── 📄 interfaces_base.py
├── 📄 interfaces.py
├── 📄 internlm2.py
├── 📄 internlm2_ve.py
├── 📄 intern_vit.py
├── 📄 internvl.py
├── 📄 jais.py
├── 📄 jamba.py
├── 📄 llama.py
├── 📄 llava_next.py
├── 📄 llava_next_video.py
├── 📄 llava_onevision.py
├── 📄 llava.py
├── 📄 mamba2.py
├── 📄 mamba_cache.py
├── 📄 mamba.py
├── 📄 medusa.py
├── 📄 minicpm3.py
├── 📄 minicpmo.py
├── 📄 minicpm.py
├── 📄 minicpmv.py
├── 📄 mixtral.py
├── 📄 mixtral_quant.py
├── 📄 mllama.py
├── 📄 mlp_speculator.py
├── 📄 module_mapping.py
├── 📄 molmo.py
├── 📄 mpt.py
├── 📄 nemotron.py
├── 📄 nvlm_d.py
├── 📄 olmo2.py
├── 📄 olmoe.py
├── 📄 olmo.py
├── 📄 opt.py
├── 📄 orion.py
├── 📄 paligemma.py
├── 📄 persimmon.py
├── 📄 phi3.py
├── 📄 phi3_small.py
├── 📄 phi3v.py
├── 📄 phi4mm_audio.py
├── 📄 phi4mm.py
├── 📄 phi4mm_utils.py
├── 📄 phimoe.py
├── 📄 phi.py
├── 📄 pixtral.py
├── 📄 prithvi_geospatial_mae.py
├── 📄 qwen2_5_vl.py
├── 📄 qwen2_audio.py
├── 📄 qwen2_moe.py
├── 📄 qwen2.py
├── 📄 qwen2_rm.py
├── 📄 qwen2_vl.py
├── 📄 qwen.py
├── 📄 qwen_vl.py
├── 📄 registry.py
├── 📄 roberta.py
├── 📄 siglip.py
├── 📄 solar.py
├── 📄 stablelm.py
├── 📄 starcoder2.py
├── 📄 telechat2.py
├── 📄 teleflm.py
├── 📄 transformers.py
├── 📄 ultravox.py
├── 📄 utils.py
├── 📄 vision.py
├── 📄 whisper.py
└── 📄 zamba2.py
|
vllm/model_executor/models中实现了几大类模型,这些模型的接口遵循interface_base.py中的定义。
具体来说是三大类模型:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
| class VllmModel(Protocol[T_co]):
"""The interface required for all models in vLLM."""
def __init__(
self,
vllm_config: "VllmConfig",
prefix: str = "",
) -> None:
...
def forward(
self,
input_ids: torch.Tensor,
positions: torch.Tensor,
) -> T_co:
...
class VllmModelForTextGeneration(VllmModel[T], Protocol[T]):
"""The interface required for all generative models in vLLM."""
def compute_logits(
self,
hidden_states: T,
sampling_metadata: "SamplingMetadata",
) -> Optional[T]:
"""Return `None` if TP rank > 0."""
...
def sample(
self,
logits: T,
sampling_metadata: "SamplingMetadata",
) -> "SamplerOutput":
"""Only called on TP rank 0."""
...
class VllmModelForPooling(VllmModel[T], Protocol[T]):
"""The interface required for all pooling models in vLLM."""
def pooler(
self,
hidden_states: T,
pooling_metadata: "PoolingMetadata",
) -> "PoolerOutput":
"""Only called on TP rank 0."""
...
|
其中 Pooling 模型是一类特殊的模型,用于处理非自回归(non-autoregressive)的任务。
主要包括 embedding 模型、重排序(reranking)模型和奖励(reward)模型。
具体的模型实现的例子:
- 基础模型
llama.py - LLaMA 模型实现gpt2.py - GPT-2 模型实现bloom.py - BLOOM 模型实现opt.py - OPT 模型实现falcon.py - Falcon 模型实现
- 多模态模型
llava.py - LLaVA 视觉语言模型llava_next.py - LLaVA 的下一代版本qwen_vl.py 和 qwen2_vl.py - 通义千问的视觉语言模型phi4mm.py - Phi-4 多模态模型minicpmv.py - MiniCPM 视觉模型
- MoE (Mixture of Experts) 模型
mixtral.py - Mixtral 模型qwen2_moe.py - 通义千问的 MoE 版本granitemoe.py - Granite MoE 模型phimoe.py - Phi MoE 模型
- 线性注意力模型
mamba.py 和 mamba2.py - Mamba 架构模型jamba.py - Jamba 模型zamba2.py - Zamba2 模型
- 特殊功能模型:
whisper.py - 语音识别模型vision.py - 视觉模型phi4mm_audio.py - 音频处理模型
- 量化相关:
mixtral_quant.py - Mixtral 量化版本
还有更多的模型实现见具体的文件。
这里以LlaMA为例,其实现很简洁,仅包含LlamaMLP、LlamaAttention、LLamaDecoderLayer、LlamaModel和LlamaForCausalLM。
其中LlamaModel是只实现了forward的VllmModel实例,而LlamaForCausalLM是实现了forward、compute_logits和sample的VllmModelForTextGeneration类型。
注意这里的模型都已经使用的是并行的线性层。
例如LlamaMLP中,gate_up_proj是MergedColumnParallelLinear,LlamaAttention中,qkv_proj是QKVParallelLinear等。
NOTE: 对于任意一个模型都可以实现一个load_weights(self, weights: Iterable[Tuple[str, torch.Tensor]])接口以自定义加载权重的行为。
NOTE: transformers.py建议逐行阅读,里面用到了很多HF transformers的good practice,例如如何初始化模型在meta上(除了用init_empty_weights,这会引入accelerate库)。
vllm/model_executor/layers
这个目录提供了更具体的神经网络层实现,之后再回头研究。
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
| vllm/model_executor/layers/
├── 📁 fused_moe/
├── 📁 mamba/
├── 📁 quantization/
├── 📄 activation.py
├── 📄 __init__.py
├── 📄 layernorm.py
├── 📄 linear.py
├── 📄 logits_processor.py
├── 📄 pooler.py
├── 📄 rejection_sampler.py
├── 📄 resampler.py
├── 📄 rotary_embedding.py
├── 📄 sampler.py
├── 📄 spec_decode_base_sampler.py
├── 📄 typical_acceptance_sampler.py
├── 📄 utils.py
└── 📄 vocab_parallel_embedding.py
4 directories, 14 files
|
vllm/worker