Worker: hält Modell und KV-Cache, exposes execute_model für den Executor
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-96legt die hardwareunabhängigen Methodensignaturen fest,Workerimplementiert 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 nachupdate_environment_variablesaufgerufen, sodassCUDA_VISIBLE_DEVICESbereits gesetzt ist; zudem wirdworker_extension_clsfür dynamische Mehrfachvererbung unterstützt(worker_base.py:261-287), um den Worker per Mixin zu patchen. - Forward an model_runner delegiert:
Worker.execute_modelführt das Modell nicht selbst aus, sondern erledigt nur PP-spezifischesirecv/isend_tensor_dictund Profiler-Anmerkungen; der eigentliche Forward geht anself.model_runner.execute_model(gpu_worker.py:1014-1017), womit PP-Kommunikation und Forward entkoppelt sind. - Memory Profiling als separate Methode:
determine_available_memoryfü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 vonnum_gpu_blocksnutzt. - Sleep/Wake über Plattform-Backend:
sleep/wake_uplaufen überSleepModeBackendFactory.create_backend(gpu_worker.py:176-185); bei level=2 werden Buffer auf die CPU kopiert und beim Aufwachen percopy_zurückgespielt(gpu_worker.py:187-229). - KV-Cache konfigurationsgetrieben allokieren:
initialize_from_config(kv_cache_config)(gpu_worker.py:695-721) ruftmodel_runner.initialize_kv_cachemit den vom Scheduler berechnetennum_blocksauf; der Worker entscheidet nicht selbst über die KV-Cache-Größe.
Schlüsseldateien
WorkerBase:39-96— Abstrakte Basisklasse: speichert Konfiguration, definiert Signaturen fürinit_device/load_model/execute_model/sample_tokens/check_health/shutdown, definiertis_driver_worker,rank,local_rank,distributed_init_method.WorkerWrapperBase.init_worker:187-294— holt fürrpc_rankdie passenden kwargs, findet viaresolve_obj_by_qualname(parallel_config.workercls)die Klasse, wendet optional dynamische Vererbung mitworker_extension_clsan und instantiiert danach den echtenWorker.WorkerWrapperBase.init_device + __getattr__:327-334—init_devicewickeltset_current_vllm_configein, danach leitet__getattr__alle nicht getroffenen Attribute anself.workerweiter.WorkerWrapperBase.execute_model:346-351— lässtscheduler_outputvorher durch_apply_mm_cachelaufen und reicht die Anfrage dann anself.worker.execute_modelweiter.Worker.__init__:126-175— speichert Konfiguration,torch.set_float32_matmul_precision, hängtElasticEPScalingExecutoran, reserviert Platzhalter für_sleep_saved_buffers/weight_transfer_engine/profiler/_pp_send_work.Worker.init_device:279-340— löschtNCCL_ASYNC_ERROR_HANDLING, berechnetlocal_rankaus DP×TP×PP, ruftset_assigned_physical_gpu_idsauf, danninit_distributed_environment+ensure_model_parallel_initialized.Worker.load_model:406-421— umschließtmodel_runner.load_modelmit_maybe_get_memory_pool_context(tag="weights")und legt optional einenWeightTransferEnginean.Worker.determine_available_memory:429-540—memory_profiling+model_runner.profile_run()ermitteln die Spitze, ziehen Nicht-KV-Anteile ab und liefernavailable_kv_cache_memory_bytes.Worker.initialize_from_config:695-721— schreibtnum_gpu_blocks, ruftensure_kv_transfer_initializedauf und ruftmodel_runner.initialize_kv_cacheinnerhalb von_maybe_get_memory_pool_context(tag="kv_cache")auf.Worker.compile_or_warm_up_model:723-790— führt für diecompile_sizes_dummy_runaus, dannkernel_warmup(self)und schließlichmodel_runner.capture_model()zur Aufnahme der cudagraph-Graphen.Worker.sample_tokens:948-951— einzeilige Weiterleitung anmodel_runner.sample_tokens.Worker.execute_model:955-1037—@torch.inference_mode()+@with_gpu_sync_check, behandelt PPirecv/isend, in der Mittemodel_runner.execute_model; mittlere PP-Ranks liefernIntermediateTensorszurück.Worker.profile:1048-1099— wählt je nachprofiler_config.profilerentwederTorchProfilerWrapperoderCudaProfilerWrapper, mit Start/Stop in beide Richtungen.Worker.shutdown:1221-1247—ensure_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:
# 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_device → load_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_deviceberechnet bei nicht-ray/external_launcher undnnodes_within_dp == 1denlocal_rankneu alsdp_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üfunglocal_world_size <= len(assigned_physical_gpu_ids), weilnnodesbei diesen Backends immer 1 ist(gpu_worker.py:314-328).- PP-Sendevorgang nicht abgeschlossen:Am Anfang von
execute_modelwird_pp_send_workdes vorherigen Steps perwait()abgewartet(gpu_worker.py:958-962), da ein vorherigesisendabgeschlossen sein muss, bevor der nächste Batch betreten wird. - Ungültige Profiler-Konfiguration:In
__init__wird geprüft, dassprofiler_config.profiler in ("torch", "cuda", None)(gpu_worker.py:166-167); beim Aufruf vonprofile()ohne Konfiguration wird direktraise RuntimeErrormit dem Hinweis auf--profiler-configgeworfen(gpu_worker.py:1050-1056). - sleep level=2: Buffer zurückspielen:In
wake_upwerden die in_sleep_saved_buffersabgelegten Buffer pername-Weisercopy_zurückgespielt(gpu_worker.py:220-229); je nach Tags wird optionalpost_kv_cache_wake_upausgefü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 aufNone, mehrfacher Aufruf ist sicher. check_healthliefert immer:Worker.check_healthist direkt einreturn(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.