parallel_state und GroupCoordinator: TP/PP/DP-Kommunikationsgruppen
Verantwortung
vLLM führt verteilte Inferenz über PyTorch torch.distributed und dessen ProcessGroup aus, aber die rohe ProcessGroup ist umständlich und schwer zu vereinheitlichen: CPU-Kommunikation und Device-Kommunikation müssen in zwei getrennte Groups aufgeteilt werden, benutzerdefinierte All-Reduce läuft über Custom Ops, und das Broadcasten von Objekten vs. Tensoren sind zwei verschiedene APIs. GroupCoordinator(parallel_state.py:358) ist die Schicht, die vLLM über die ProcessGroup legt; sie verpackt die CPU-Group, die Device-Group, den device_communicator und den mq_broadcaster einer Prozessgruppe in einem Objekt und stellt eine einheitliche Schnittstelle all_reduce / all_gather / broadcast / send_object / recv_tensor_dict / barrier bereit.
Das parallel_state-Modul verwaltet mehrere globale Singletons: _TP (Tensor Parallelism), _PP (Pipeline Parallelism), _DP (Data Parallelism), _EP (Expert Parallelism), _DCP (Decode Context Parallelism), _PCP (Prefill Context Parallelism), _EPLB (Expert Load Balancing). Jede Gruppe wird über get_tp_group / get_pp_group / get_dp_group etc. ausgelesen(parallel_state.py:1368-1424); auf Modellebene greifen ColumnParallelLinear / RowParallelLinear über diese Getter auf die Kommunikationsgruppe zu. Die Initialisierung der Gruppen erfolgt durch initialize_model_parallel(parallel_state.py:1713), das beim Start einmal aufgerufen wird und den globalen Rang nach tensor_model_parallel_size * pipeline_model_parallel_size * ... in Segmente aufteilt und new_group aufruft.
GroupCoordinator übernimmt noch eine zweite Aufgabe: all_reduce / all_gather werden als torch.ops.vllm.* Custom Ops registriert, sodass Dynamo beim Kompilieren die Kommunikations-Ops als gewöhnliche Ops in den Graph fusionieren kann(parallel_state.py:641-663), anstatt ein Python-Objekt wie self in den Kompilations-Graphen zu pressen. Bei einer World-Size (world size) von 1 kehren alle Kommunikations-Ops direkt mit der Eingabe zurück und umgehen NCCL.
Entwurfsmotivation
- Eine Gruppe, zwei PGs:
cpu_group(Gloo-Backend) unddevice_group(nccl/gloo/Plattform-Backend) sind getrennt(parallel_state.py:435-452); auf der CPU finden Objekt-Broadcast und Koordination zwischen Koordinatoren statt, auf dem Device die Tensor-Kommunikation. - Drei Rang-Konzepte sauber trennen:
rank(global),local_rank(auf der Maschine),rank_in_group(in der Gruppe)(parallel_state.py:369-380) — bei Multi-Node sind diese drei verschieden; die Modellebene mussrank_in_groupverwenden, um korrekt zu slicen. - Rank-Layout feste Reihenfolge:
ExternalDP × DP × PP × PCP × TP(parallel_state.py:1779-1794); nachtorch.arange(world_size).reshape(...)wird entlang einer Dimension unbindet, um die Rang-Liste der jeweiligen Gruppe zu erhalten. - Single-GPU-Bypass: Methoden wie
all_reducekehren beiworld_size == 1direkt zurück(parallel_state.py:657-658), damit auch eine einzelne GPU nicht NCCL starten muss. - Custom Op für Dynamo: Bei
use_custom_op_call=Trueläufttorch.ops.vllm.all_reduce(input_, group_name=self.unique_name)(parallel_state.py:660-661); die Gruppe wird als String nachgeschlagen, damit Dynamo nicht über ein Python-Objekt stolpert. - DeviceCommunicator austauschbar: Das Feld
device_communicatorerlaubt CUDA / XPU / CPU / eigenen All-Reduce-Implementierungen jeweils ihre eigene(parallel_state.py:463-478);is_cuda_alikedbindet ancuda:N, XPU anxpu:N. - TP nutzt Message Queue für Broadcast:
init_model_parallel_groupaktiviert in der TP-Gruppeuse_message_queue_broadcaster=True(parallel_state.py:1805-1811), sodass Worker (worker) Tensor-Dictionaries teilen können, ohne über NCCL zu gehen. - StatelessGroupCoordinator als Fallback: Im elastischen EP-Szenario ohne initialisiertes torch.distributed baut
_init_stateless_group(parallel_state.py:1316-1338) über TCP Store + StatelessProcessGroup eine äquivalente Gruppe auf.
Schlüsseldateien
GroupCoordinator-Klasse:358— Prozessgruppen-Wrapper; hält rank/ranks/world_size/local_rank/rank_in_group/cpu_group/device_group.GroupCoordinator.__init__:387-462— Wählt je nachVLLM_DISTRIBUTED_USE_SPLIT_GROUPsplit_group oder new_group und bindet das Device.GroupCoordinator.all_reduce:641— Geht übertorch.ops.vllm.all_reduceoder_all_reduce_out_place; Bypass bei world_size=1.GroupCoordinator.all_gather:670— Ebenfalls Custom Op; prüft dim und bypass.broadcast:745— Broadcastet einen Tensor auf der Device-Group, mit src-Parameter.broadcast_object:760— Nutztbroadcast_object_listauf der cpu_group, um Konfigurationen oder den Sampling-Generator-State zu synchronisieren.broadcast_tensor_dict:864— Broadcastet ein gesamtes Dict; beim Worker-Start zur Synchronisation von KV-Cache-Metadaten genutzt.barrier:1179— Doppelte Barrier auf Device + CPU, damit GPU-Operationen und CPU-Koordination beidealigned sind.initialize_model_parallel:1713— Teilt Ränge nachExternalDP × DP × PP × PCP × TPund initialisiert alle Gruppen.TP-Gruppen-Konstruktion:1796-1811—group_ranks = all_ranks.view(-1, tensor_model_parallel_size).unbind(0).init_model_parallel_group:1298— Fabrikmethode, delegiert anGroupCoordinator(...).get_tp_group:1368— Globaler Getter; die Modellebene holt darüber die TP-Gruppe.get_tensor_model_parallel_world_size:2031— Bequemlichkeits-Query; internget_tp_group().world_size._register_group:126— Registriert jeden GroupCoordinator in einem globalen Dict, damit Custom Ops per group_name nachschlagen können.
Datenfluss
Beim Start ruft der Worker in seinem Prozess initialize_model_parallel(tp_size, pp_size, ...) auf; diese Funktion reshapt den globalen Rang in der Reihenfolge ExternalDP × DP × PP × PCP × TP und unbindet dann entlang jeder Dimension, um die Rang-Liste jeder Gruppe zu erhalten. Im Folgenden die Konstruktion der TP-Gruppe:
# vllm/distributed/parallel_state.py L1788-L1811
all_ranks = torch.arange(world_size).reshape(
-1,
data_parallel_size,
pipeline_model_parallel_size,
prefill_context_model_parallel_size,
tensor_model_parallel_size,
) # noqa
# Build the tensor model-parallel groups.
global _TP
assert _TP is None, "tensor model parallel group is already initialized"
group_ranks = all_ranks.view(-1, tensor_model_parallel_size).unbind(0)
group_ranks = [x.tolist() for x in group_ranks]
if enable_elastic_ep:
group_ranks = local_all_ranks.view(-1, tensor_model_parallel_size).unbind(0)
group_ranks = [x.tolist() for x in group_ranks]
# message queue broadcaster is only used in tensor model parallel group
_TP = init_model_parallel_group(
group_ranks,
get_world_group().local_rank,
backend,
use_message_queue_broadcaster=True,
group_name="tp",
)Wenn die Modellebene kommuniziert, berührt sie nicht direkt die ProcessGroup, sondern ruft get_tp_group().all_reduce(tensor) auf; intern entscheidet all_reduce je nach use_custom_op_call, ob der Custom Op oder direkt device_communicator.all_reduce verwendet wird:
# vllm/distributed/parallel_state.py L641-L668
def all_reduce(self, input_: torch.Tensor) -> torch.Tensor:
"""
User-facing all-reduce function before we actually call the
all-reduce operation.
"""
# Bypass the function if we are using only 1 GPU.
if self.world_size == 1:
return input_
if self.use_custom_op_call:
return torch.ops.vllm.all_reduce(input_, group_name=self.unique_name)
else:
return self._all_reduce_out_place(input_)
def _all_reduce_out_place(self, input_: torch.Tensor) -> torch.Tensor:
if self.device_communicator is None:
raise ValueError("No device communicator found")
return self.device_communicator.all_reduce(input_)Der Forward von ColumnParallelLinear ruft nach der Berechnung des lokalen Anteils all_reduce auf, um die Shards zusammenzuführen; RowParallelLinear geht umgekehrt vor: erst all-gather, dann lokal rechnen, zuletzt reduce_scatter. Diese Ops sind alle in der GroupCoordinator-Schicht einheitlich gekapselt.
Grenzen und Fehler
- rank_in_group muss verwendet werden: Die Modellebene muss zum Slicen
rank_in_groupverwenden, nicht den globalenrank, sonst wird in Multi-Node-Szenarien das falsche Shard geschnitten(parallel_state.py:380). - Doppelte TP-Initialisierung schlägt fehl:
initialize_model_parallelprüft für jede Gruppeassert _TP is None(parallel_state.py:1798); ein erneuter Aufruf führt zu AssertionError und erfordert vorherdestroy_model_parallel. - DCP darf TP nicht überschreiten:
dcp_size <= tp_size, weil DCP die GPU der TP-Gruppe mitverwendet(parallel_state.py:1816-1820). - Single-GPU-Bypass:
all_reducekehrt beiworld_size == 1direkt mit der Eingabe zurück(parallel_state.py:657-658) und benötigt keinen NCCL-Prozess; ist aberdevice_communicator is Noneund kein Bypass-Szenario, wirdValueErrorgeworfen. - Fehlender DeviceCommunicator wirft Fehler:
_all_reduce_out_placeprüftdevice_communicator is Noneund wirft(parallel_state.py:665-668), um Aufrufe auf einem nicht initialisierten Backend zu verhindern. _replace_active_groupsmuss kollektiv aufgerufen werden:parallel_state.py:1341-1362verlangt, dass alle Ränge gleichzeitig aufrufen, sonst droht Deadlock; beim elastischen EP-Skalieren übernimmt der Supervisor die einheitliche Steuerung.
Zusammenfassung
parallel_state ist das Fundament von vLLMs Multi-GPU-Inferenz: GroupCoordinator kapselt PyTorchs ProcessGroup und fasst CPU-/Device-Double-Group, Custom Ops, device_communicator und Message-Queue-Broadcaster in einem Objekt zusammen; die Modellebene holt über get_tp_group() / get_pp_group() die entsprechende Gruppe und ruft dann all_reduce etc. auf. Wie die TP-Linearschichten diese Kommunikation nutzen, siehe /model-loading/tp-layers; wie der Worker diese Gruppen startet, siehe /worker/worker; wie der Engine-Kern (EngineCore) diese Worker verdrahtet, siehe /engine/engine-core.