Skip to content

Worker: hält Modell und KV-Cache, exposes execute_model für den Executor

源码版本v0.25.1

Verantwortung

Der Executor (executor) leitet nur die RPCs weiter; die eigentliche Arbeit übernimmt der Worker (worker). Die v1-Worker-Abstraktion liegt in WorkerBase(worker_base.py:39-96), die GPU-Implementierung ist Worker(gpu_worker.py:126-175): er hält das Modell (model_runner.model), hält den KV-Cache (KV cache), hält die verteilten Kommunikationsgruppen (TP/PP/EP/DP) und exposes execute_model, sample_tokens, profile, sleep/wake_up, initialize_from_config, determine_available_memory als Methoden, die über collective_rpc erreichbar sind.

Zwischen Executor und Worker liegt zudem eine Schicht WorkerWrapperBase(worker_base.py:187-294): der Executor hält den Wrapper, der über init_worker erst die eigentliche Worker-Instanz erzeugt (möglicherweise nachdem Ray bereits CUDA_VISIBLE_DEVICES gesetzt hat) und dann per __getattr__ alle Methoden an self.worker weiterleitet(worker_base.py:333-334). Diese verzögerte Instantiierung ist für das Ray-Backend essenziell: zum Zeitpunkt der Actor-Erzeugung steht noch nicht fest, welche GPU ihm gehört.

Worker.__init__ macht hauptsächlich drei Dinge: Konfiguration speichern, ElasticEPScalingExecutor anhängen und Platzhalter für Sleep/Wake-Up/Gewichtstransfer/Profiler vorsehen(gpu_worker.py:127-175). Das eigentliche CUDA-Init, das Laden der Gewichte und die Allokation des KV-Cache erfolgen später in den drei Schritten init_device / load_model / initialize_from_config.

Entwurfsmotivation

  • Abstraktionsschicht WorkerBase:worker_base.py:39-96 legt die hardwareunabhängigen Methodensignaturen fest, Worker implementiert nur den GPU-Teil, CPU/XPU/TPU erhalten jeweils eigene Unterklassen, sodass der Executor von Hardwareunterschieden nichts merkt.
  • Verzögerte Instantiierung im Wrapper:WorkerWrapperBase.init_worker(worker_base.py:229-294) wird erst nach update_environment_variables aufgerufen, sodass CUDA_VISIBLE_DEVICES bereits gesetzt ist; zudem wird worker_extension_cls für dynamische Mehrfachvererbung unterstützt(worker_base.py:261-287), um den Worker per Mixin zu patchen.
  • Forward an model_runner delegiert:Worker.execute_model führt das Modell nicht selbst aus, sondern erledigt nur PP-spezifisches irecv/isend_tensor_dict und Profiler-Anmerkungen; der eigentliche Forward geht an self.model_runner.execute_model(gpu_worker.py:1014-1017), womit PP-Kommunikation und Forward entkoppelt sind.
  • Memory Profiling als separate Methode:determine_available_memory führt einen Dummy-Forward aus, um die Spitzen-VRAM-Belegung zu messen(gpu_worker.py:429-493); nach Abzug von Gewichten und Aktivierungen bleibt das für den KV-Cache Verfügbare, das der EngineCore anschließend zur Bestimmung von num_gpu_blocks nutzt.
  • Sleep/Wake über Plattform-Backend:sleep/wake_up laufen über SleepModeBackendFactory.create_backend(gpu_worker.py:176-185); bei level=2 werden Buffer auf die CPU kopiert und beim Aufwachen per copy_ zurückgespielt(gpu_worker.py:187-229).
  • KV-Cache konfigurationsgetrieben allokieren:initialize_from_config(kv_cache_config)(gpu_worker.py:695-721) ruft model_runner.initialize_kv_cache mit den vom Scheduler berechneten num_blocks auf; der Worker entscheidet nicht selbst über die KV-Cache-Größe.

Schlüsseldateien

  • WorkerBase:39-96 — Abstrakte Basisklasse: speichert Konfiguration, definiert Signaturen für init_device/load_model/execute_model/sample_tokens/check_health/shutdown, definiert is_driver_worker, rank, local_rank, distributed_init_method.
  • WorkerWrapperBase.init_worker:187-294 — holt für rpc_rank die passenden kwargs, findet via resolve_obj_by_qualname(parallel_config.workercls) die Klasse, wendet optional dynamische Vererbung mit worker_extension_cls an und instantiiert danach den echten Worker.
  • WorkerWrapperBase.init_device + __getattr__:327-334init_device wickelt set_current_vllm_config ein, danach leitet __getattr__ alle nicht getroffenen Attribute an self.worker weiter.
  • WorkerWrapperBase.execute_model:346-351 — lässt scheduler_output vorher durch _apply_mm_cache laufen und reicht die Anfrage dann an self.worker.execute_model weiter.
  • Worker.__init__:126-175 — speichert Konfiguration, torch.set_float32_matmul_precision, hängt ElasticEPScalingExecutor an, reserviert Platzhalter für _sleep_saved_buffers/weight_transfer_engine/profiler/_pp_send_work.
  • Worker.init_device:279-340 — löscht NCCL_ASYNC_ERROR_HANDLING, berechnet local_rank aus DP×TP×PP, ruft set_assigned_physical_gpu_ids auf, dann init_distributed_environment + ensure_model_parallel_initialized.
  • Worker.load_model:406-421 — umschließt model_runner.load_model mit _maybe_get_memory_pool_context(tag="weights") und legt optional einen WeightTransferEngine an.
  • Worker.determine_available_memory:429-540memory_profiling + model_runner.profile_run() ermitteln die Spitze, ziehen Nicht-KV-Anteile ab und liefern available_kv_cache_memory_bytes.
  • Worker.initialize_from_config:695-721 — schreibt num_gpu_blocks, ruft ensure_kv_transfer_initialized auf und ruft model_runner.initialize_kv_cache innerhalb von _maybe_get_memory_pool_context(tag="kv_cache") auf.
  • Worker.compile_or_warm_up_model:723-790 — führt für die compile_sizes _dummy_run aus, dann kernel_warmup(self) und schließlich model_runner.capture_model() zur Aufnahme der cudagraph-Graphen.
  • Worker.sample_tokens:948-951 — einzeilige Weiterleitung an model_runner.sample_tokens.
  • Worker.execute_model:955-1037@torch.inference_mode() + @with_gpu_sync_check, behandelt PP irecv/isend, in der Mitte model_runner.execute_model; mittlere PP-Ranks liefern IntermediateTensors zurück.
  • Worker.profile:1048-1099 — wählt je nach profiler_config.profiler entweder TorchProfilerWrapper oder CudaProfilerWrapper, mit Start/Stop in beide Richtungen.
  • Worker.shutdown:1221-1247ensure_kv_transfer_shutdown + ensure_ec_transfer_shutdown + profiler.shutdown() + weight_transfer_engine.shutdown() + model_runner.shutdown() + CuMemAllocator.release_pools().

Datenfluss

In einem einzelnen Step ruft der Executor Worker.execute_model(scheduler_output) auf. Zunächst werden PP-Verbindungen abgearbeitet (nicht-erste Ranks empfangen asynchron irecv_tensor_dict für IntermediateTensors), dann folgt ein annotate_profile-Wrapper, und schließlich läuft der Forward über model_runner.execute_model:

python
# vllm/v1/worker/gpu_worker.py L1000-L1017
if forward_pass and not get_pp_group().is_first_rank:
    tensor_dict, comm_handles, comm_postprocess = (
        get_pp_group().irecv_tensor_dict(
            all_gather_group=get_tp_group(),
            all_gather_tensors=all_gather_tensors,
        )
    )
    assert tensor_dict is not None
    intermediate_tensors = AsyncIntermediateTensors(
        tensor_dict,
        comm_handles=comm_handles,
        comm_postprocess=comm_postprocess,
    )

with self.annotate_profile(scheduler_output):
    output = self.model_runner.execute_model(
        scheduler_output, intermediate_tensors
    )

Ist der Forward abgeschlossen und output ein IntermediateTensors-Objekt (mittlerer PP-Rank), sendet der Worker es asynchron per isend_tensor_dict an den nächsten PP-Rank(gpu_worker.py:1036-1037); die zugehörigen Handles werden zu Beginn des nächsten Steps per wait abgearbeitet. Der finale Rank (is_last_rank) liefert ModelRunnerOutput/AsyncModelRunnerOutput zurück, das der Executor einsammelt und dem Future des EngineCore übergibt.

Die Startphase ist ein zweiter Strang: init_deviceload_model → (Executor-Aufruf) determine_available_memory → (Executor-Aufruf) initialize_from_config → (Executor-Aufruf) compile_or_warm_up_model. Erst nach diesen fünf Schritten ist der Worker bereit.

Grenzen und Fehler

  • Falle CUDA_VISIBLE_DEVICES-Reihenfolge:init_device berechnet bei nicht-ray/external_launcher und nnodes_within_dp == 1 den local_rank neu als dp_local_rank * tp_pp_world_size + tp_local_rank(gpu_worker.py:283-301); Ray und external_launcher überspringen diesen Abschnitt, da sie die Geräte selbst bereits angeordnet haben.
  • assigned_physical_gpu_ids-Validierung:Multi-Node Ray/external_launcher überspringen die Prüfung local_world_size <= len(assigned_physical_gpu_ids), weil nnodes bei diesen Backends immer 1 ist(gpu_worker.py:314-328).
  • PP-Sendevorgang nicht abgeschlossen:Am Anfang von execute_model wird _pp_send_work des vorherigen Steps per wait() abgewartet(gpu_worker.py:958-962), da ein vorheriges isend abgeschlossen sein muss, bevor der nächste Batch betreten wird.
  • Ungültige Profiler-Konfiguration:In __init__ wird geprüft, dass profiler_config.profiler in ("torch", "cuda", None)(gpu_worker.py:166-167); beim Aufruf von profile() ohne Konfiguration wird direkt raise RuntimeError mit dem Hinweis auf --profiler-config geworfen(gpu_worker.py:1050-1056).
  • sleep level=2: Buffer zurückspielen:In wake_up werden die in _sleep_saved_buffers abgelegten Buffer per name-Weiser copy_ zurückgespielt(gpu_worker.py:220-229); je nach Tags wird optional post_kv_cache_wake_up ausgeführt.
  • Reihenfolge beim Shutdown:gc.unfreeze → kv transfer → ec transfer → profiler → weight transfer engine → model_runner → cumem pools(gpu_worker.py:1221-1247); jeder Schritt prüft auf None, mehrfacher Aufruf ist sicher.
  • check_health liefert immer:Worker.check_health ist direkt ein return(gpu_worker.py:1117-1119) — solange der Prozess lebt, gilt er als gesund; bei echtem Absturz muss die RPC-Ausnahme nach oben propagiert werden.

Zusammenfassung

Der Worker (worker) ist die Stelle in v1, an der der Forward tatsächlich stattfindet. WorkerBase definiert die hardwareunabhängige Schnittstelle, Worker implementiert den GPU-Teil, und WorkerWrapperBase fügt eine verzögerte Instantiierung hinzu, damit das Ray-Backend die GPU erst nach der Actor-Erzeugung binden kann. Der Worker hält das Modell, den KV-Cache und die verteilten Gruppen und exposes die RPC-Methoden wie execute_model gegenüber dem Executor; den eigentlichen Forward delegiert er an GPUModelRunner und kümmert sich selbst nur um PP-Kommunikation, Profiler, Sleep/Wake-Up und Gewichtstransfer. Wie der Executor die Worker zusammenführt, ist in /executor/executor und /executor/ray beschrieben; der eigentliche Forward in /worker/gpu-model-runner; cudagraph und ubatch in /worker/ubatch-cudagraph.

Siehe offizielle Dokumentation: vLLM 文档 · README.