Skip to content

parallel_state und GroupCoordinator: TP/PP/DP-Kommunikationsgruppen

源码版本v0.25.1

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) und device_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 muss rank_in_group verwenden, um korrekt zu slicen.
  • Rank-Layout feste Reihenfolge: ExternalDP × DP × PP × PCP × TP(parallel_state.py:1779-1794); nach torch.arange(world_size).reshape(...) wird entlang einer Dimension unbindet, um die Rang-Liste der jeweiligen Gruppe zu erhalten.
  • Single-GPU-Bypass: Methoden wie all_reduce kehren bei world_size == 1 direkt 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=True läuft torch.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_communicator erlaubt CUDA / XPU / CPU / eigenen All-Reduce-Implementierungen jeweils ihre eigene(parallel_state.py:463-478); is_cuda_aliked bindet an cuda:N, XPU an xpu:N.
  • TP nutzt Message Queue für Broadcast: init_model_parallel_group aktiviert in der TP-Gruppe use_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

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:

python
# 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:

python
# 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_group verwenden, nicht den globalen rank, sonst wird in Multi-Node-Szenarien das falsche Shard geschnitten(parallel_state.py:380).
  • Doppelte TP-Initialisierung schlägt fehl: initialize_model_parallel prüft für jede Gruppe assert _TP is None(parallel_state.py:1798); ein erneuter Aufruf führt zu AssertionError und erfordert vorher destroy_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_reduce kehrt bei world_size == 1 direkt mit der Eingabe zurück(parallel_state.py:657-658) und benötigt keinen NCCL-Prozess; ist aber device_communicator is None und kein Bypass-Szenario, wird ValueError geworfen.
  • Fehlender DeviceCommunicator wirft Fehler: _all_reduce_out_place prüft device_communicator is None und wirft(parallel_state.py:665-668), um Aufrufe auf einem nicht initialisierten Backend zu verhindern.
  • _replace_active_groups muss kollektiv aufgerufen werden: parallel_state.py:1341-1362 verlangt, 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.

Siehe offizielle Dokumentation: vLLM 文档 · README.