Skip to content

Scheduler.schedule : décider à chaque étape qui s'exécute

源码版本v0.25.1

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 de RequestQueue, voir RequestQueue.

Fichiers clés

  • Scheduler 类:68 — définition de la classe, hérite de SchedulerInterface ; 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, signature schedule(throttle_prefills=False) -> SchedulerOutput.
  • token_budget:416token_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 : calcule num_new_tokens de 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 — quand num_running >= self.max_num_running_reqs, plus aucune nouvelle request n'est acceptée.
  • SchedulerOutput:1092-1109 — assemble le SchedulerOutput final 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, avance num_computed_tokens et 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 :

python
# 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.
            break

Ce 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_budget est forcé à 0(PAUSED_ALL:417-419) ; en PAUSED_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:637 garantit 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 fait continue plutôt que break(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 == 1 et qu'il n'y a pas de prefill, on pad en 1 + num_spec_tokens pour 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_prefill est 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=True et 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_KVS sont 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