Scheduler.schedule : décider à chaque étape qui s'exécute
Responsabilités
Dans l'architecture v1, Scheduler est l'un des trois composants majeurs d'EngineCore(scheduler.py:68), placé entre KVCacheManager et model_executor. Il ne touche pas au GPU et n'exécute pas la passe avant ; il ne répond qu'à une seule question : à cette étape (step), à quels request allouer combien de tokens. Le step() d'EngineCore appelle scheduler.schedule() en tête de chaque boucle(schedule 方法:396), et le SchedulerOutput obtenu sert de base au model runner pour assembler les tenseurs d'entrée. Un step correspond à une passe avant, donc cette méthode est appelée en boucle serrée.
L'ordonnanceur (scheduler) maintient deux files : waiting (requests non encore planifiées) et running (requests déjà pourvues en blocs KV et capables de produire des tokens)(waiting/running:181-184), auxquelles s'ajoute skipped_waiting qui stocke temporairement les requests sautées en raison d'un dépassement LoRA, d'un chargement asynchrone KV connector, etc. Il n'opère pas de séparation du type « tout prefill d'abord, puis tout decode » : la progression de chaque request est représentée comme num_computed_tokens rattrapant num_tokens_with_spec(算法注释:398-407). Le chunked prefill, la cache de préfixe (prefix caching) et le speculative decoding partagent tous cette même sémantique.
Motivation de conception
Pourquoi pas de distinction explicite prefill phase / decode phase ?
- Modèle de progression unifié : chaque request n'a qu'un seul pointeur de progression,
num_computed_tokens. Un prefill qui entre correspond à « calculer de 0 à prompt_len », un decode à « avancer d'un token », un chunked prefill à « avancer d'un segment ». La logique d'ordonnancement partage le même code ; ajouter une optimisation (jump decoding, spec decode) ne nécessite pas de toucher à la boucle principale(算法注释:398-407). - Budget de tokens, pas de requêtes : la contrainte dure de la boucle principale est
max_num_scheduled_tokens(token_budget:416). Chaque request ajoutée déduit du budget ; quand il est épuisé, on s'arrête. Ainsi, un long prefill ne bloque pas tout un batch mélangeant des requêtes courtes et longues. - running prioritaire : la file running bénéficie d'une priorité sur le budget de tokens(
running 循环:442). Une request déjà pourvue de blocs KV gaspille de la mémoire si elle ne progresse pas, donc on sert d'abord running avant de piocher dans waiting. - Filet de sécurité par préemption : quand une request de running ne peut pas obtenir de nouveau bloc, on évince une request running de plus basse priorité(
preempt 循环:534-575). C'est le seul moyen de continuer à ordonnancer quand le KV cache est tendu, voir Préemption et prefill throttle. - Stratégie interchangeable : la file est instanciée via
create_request_queue(self.policy)(create_request_queue:181). FCFS et priority correspondent à deux implémentations deRequestQueue, voir RequestQueue.
Fichiers clés
Scheduler 类:68— définition de la classe, hérite deSchedulerInterface; le constructeur instancie KVCacheManager, connector, files et pause state.waiting/running:181-184— les trois files :waiting,skipped_waiting,running.schedule:396— méthode principale, signatureschedule(throttle_prefills=False) -> SchedulerOutput.token_budget:416—token_budget = self.max_num_scheduled_tokens, nombre total de tokens émettables sur cette étape.defer_prefills:436-438— équilibrage DP prefill : sur l'étape throttle, les chunks prefill en cours sont repoussés à une étape d'alignement.running 循环:442-591— planifie d'abord running : calculenum_new_tokensde chaque request et alloue les blocs KV.allocate_slots:534-535— alloue de nouveaux slots à un running request ; en cas d'échec déclenche la préemption.waiting 循环:637-1011— planifie ensuite waiting : après vérification du hit prefix cache, des entrées encoder, du KV connector, etc., la request est promue dans running.max_num_running:644— quandnum_running >= self.max_num_running_reqs, plus aucune nouvelle request n'est acceptée.SchedulerOutput:1092-1109— assemble leSchedulerOutputfinal et emballe nouvelles requêtes, requêtes en cache, nombre de tokens, preempt id pour le model runner._update_after_schedule:1164-1211— après le step, avancenum_computed_tokenset nettoie les ensembles finished/preempt id.
Flux de données
À l'entrée de schedule(), on bump current_step, puis le travail se fait en deux boucles : « running d'abord, waiting ensuite ». Le cœur de la boucle running consiste à prendre une request, calculer combien de tokens émettre, puis aller chercher les nouveaux blocs auprès de KVCacheManager :
# Schedule newly needed KV blocks for the request.
with record_function_or_nullcontext("schedule: allocate_slots"):
while True:
new_blocks = self.kv_cache_manager.allocate_slots(
request,
num_new_tokens,
num_lookahead_tokens=self.num_lookahead_tokens,
)
if new_blocks is not None:
# The request can be scheduled.
break
# The request cannot be scheduled.
# Preempt the lowest-priority request.
if self.policy == SchedulingPolicy.PRIORITY:
preempted_req = max(
self.running,
key=lambda r: (r.priority, r.arrival_time),
)
self.running.remove(preempted_req)
...
else:
preempted_req = self.running.pop()
self._preempt_request(preempted_req, scheduled_timestamp)
preempted_reqs.append(preempted_req)
if preempted_req == request:
# No more request to preempt. Cannot schedule this request.
breakCe bloc est dans scheduler.py:533-582. Quand allocate_slots renvoie None, c'est que le KV cache est plein : il faut évincer une request de running — la stratégie priority choisit le (priority, arrival_time) maximal, FCFS fait simplement pop() en fin de file. La request évincée passe par _preempt_request(_preempt_request:1140) qui la remet dans waiting. À noter la ligne preempted_req == request : si la préemption finit par s'évincer soi-même, on admet l'échec, on break, et cette request ne s'exécute pas pour ce step.
L'entrée de la boucle waiting est doublement protégée : aucune préemption ce step + pause state à UNPAUSED(waiting 入口:637). Dès qu'une préemption a eu lieu, ce step cesse d'accepter de nouvelles requêtes depuis waiting, afin d'éviter qu'elles ne se disputent les ressources avec les anciennes requests réinjectées dans waiting. Dans la boucle waiting, chaque request enchaîne les contrôles : hit prefix cache, budget encoder, plafond LoRA, chargement asynchrone KV connector ; une fois tous validés, elle peut self.running.append(request)(升入 running:969) et passer à l'état RUNNING.
Enfin, SchedulerOutput(...)(SchedulerOutput 构造:1092-1109) emballe toutes les décisions du step pour le model runner, puis _update_after_schedule(_update_after_schedule:1164) avance le num_computed_tokens de chaque request et nettoie les ensembles finished/preempt id, en attendant le prochain schedule().
Limites et échecs
- PauseState : à l'entrée dans
PAUSED_ALL,token_budgetest forcé à 0(PAUSED_ALL:417-419) ; enPAUSED_NEW, la boucle waiting est sautée(waiting 入口:637), utilisé pour un arrêt global dans les scénarios KV connector. - Un step ayant préempté n'accepte plus de nouvelles requêtes : la ligne
not preempted_reqs:637garantit qu'une request réinjectée dans waiting ne se replanifie qu'au step suivant, évitant de lutter pour le même budget de tokens que les requêtes déjà running. - Traitement de
num_new_tokens == 0: quand une request de running a envoyé tout son prompt sans terminer (PP>1), ou que le budget encoder est épuisé, ou que l'alignement de bloc mamba échoue, on faitcontinueplutôt quebreak(num_new_tokens == 0 注释:514-530), de sorte que les requêtes de plus faible priorité restent éligibles. - Alignement spec decode : quand
num_new_tokens == 1et qu'il n'y a pas de prefill, on pad en1 + num_spec_tokenspour préserver la forme de batch du cuda graph(pad spec decode:814-826). Si le pad ne tient pas, on break sans forcer. - Commutateur chunked prefill : quand
enable_chunked_prefillest désactivé, si un nouveau prefill dépasse le budget restant, on break sans découper(chunked_prefill:834-840), et il attend le step suivant pour s'exécuter en entier. - Équilibrage DP prefill : avec
throttle_prefills=Trueet en l'absence de saturation, les chunks prefill en cours dans running sont différés(defer prefills:467-471) ; dans waiting, les nouveaux prefill font directement break(defer_prefills in waiting:801-804), laissant le decode remplir le step afin de garder les prefill alignés entre DP ranks. - async KV load : en chargement asynchrone, les requests à l'état
WAITING_FOR_REMOTE_KVSsont prepended àskipped_waiting(load_kv_async:946-967) et n'entrent pas dans running tant que le connector n'a pas signalé la fin du chargement. - Assertions de protection :
total_num_scheduled_tokens <= max_num_scheduled_tokens,len(self.running) <= max_num_running_reqs, etc.(断言:1023-1033) sont vérifiées à chaque fin de step ; tout bug d'ordonnancement explose immédiatement au lieu de sur-vendre silencieusement.
Résumé
Scheduler.schedule() est le battement de cœur du moteur v1 : appelé à chaque step, il alloue à chaque request un nombre de tokens et des blocs KV dans l'ordre running d'abord, waiting ensuite ; en cas d'échec d'allocation, il préempte une request running de basse priorité, puis emballe toutes les décisions dans un SchedulerOutput pour le model runner. Pour les détails de la file, voir RequestQueue : FCFS et file prioritaire ; pour les mécanismes de préemption et de prefill throttle, voir Préemption et prefill throttle.
Voir la documentation officielle : Documentation vLLM · README