Skip to content

RequestQueue : FCFS et file prioritaire

源码版本v0.25.1

Responsabilités

Scheduler maintient en interne trois files : waiting, skipped_waiting et running(waiting/running:181-184). Les deux premières sont des instances de RequestQueue, chargées de mettre en file les requests tant qu'elles n'ont pas encore reçu de blocs KV. RequestQueue(RequestQueue:20) est une classe de base abstraite qui définit un jeu d'interfaces communes (add_request, pop_request, peek_request, prepend_request, remove_request, etc.) ; l'ordre concret est fixé par deux implémentations : FCFSRequestQueue(FCFSRequestQueue:75) et PriorityRequestQueue(PriorityRequestQueue:131).

L'intérêt de la file est d'offrir à l'ordonnanceur une entrée unifiée pour « récupérer la prochaine request à planifier » — peu importe qu'en dessous ce soit une deque ou un heap, la boucle principale ne fait qu'appeler peek_request() / pop_request()(peek_request:650). La stratégie de mise en file est décidée par scheduler_config.policy, injectée dans Scheduler.__init__ via create_request_queue(self.policy)(create_request_queue:181).

Motivation de conception

Pourquoi abstraire la file plutôt que d'utiliser directement list[Request] ?

  • Stratégie interchangeable : l'enum SchedulingPolicy(SchedulingPolicy:13) ne liste que FCFS et PRIORITY, mais la classe de base relègue les détails de tri dans l'implémentation. L'ordonnanceur ne parle qu'à l'interface RequestQueue ; ajouter une nouvelle stratégie (par exemple un ordonnancement équitable) se résume à écrire une nouvelle sous-classe.
  • FCFS en deque : FCFSRequestQueue(deque[Request], RequestQueue)(FCFSRequestQueue:75) réutilise directement le collections.deque de Python, dont les append/popleft en O(1) des deux côtés sont l'implémentation optimale de FCFS. add_request correspond à append, le pop à popleft(add/pop:78-84).
  • priority en heap : PriorityRequestQueue utilise heapq(heappush:144). L'ordre est fixé par le __lt__ du Request lui-même — d'abord priority ascendant, puis arrival_time ascendant(docstring:131-139) — afin que l'utilisateur puisse attribuer une priorité à chaque request.
  • Sémantique de prepend unifiée : après une préemption ou la levée d'une dépendance asynchrone, il faut remettre une request en tête de file pour une re-planification prioritaire. Les deux files fournissent prepend_request, mais dans une file priority il n'existe pas de notion de « tête », donc prepend_request dégénère en add_request(PriorityRequestQueue.prepend:160-165), comportement que la docstring explicite.
  • skipped_waiting réutilise le même type : les requests sautées en raison d'un dépassement LoRA, d'un chargement KV asynchrone, etc., sont placées dans skipped_waiting(step_skipped_waiting:638). C'est la même sorte de file que waiting, et la boucle principale de l'ordonnanceur choisit entre les deux via _select_waiting_queue_for_scheduling() en prenant celle dont la tête est plus ancienne(_select_waiting_queue_for_scheduling:1867-1877).

Fichiers clés

  • SchedulingPolicy:13Enum, uniquement FCFS et PRIORITY.
  • RequestQueue ABC:20 — classe abstraite, définit add_request / pop_request / peek_request / prepend_request / remove_request / __bool__ / __len__ / __iter__.
  • FCFSRequestQueue:75 — hérite de deque[Request] + RequestQueue ; toutes les opérations sont des appels natifs deque.
  • FCFS add/pop:78-84add_request via append, pop_request via popleft, implémentation standard FCFS.
  • FCFS prepend:92-94prepend_request via appendleft, de sorte qu'une request réinjectée par préemption soit re-planifiée en priorité au step suivant.
  • FCFS prepend_requests:96-103prepend_requests via extendleft ; la docstring signale que l'ordre des éléments prépended est inversé.
  • FCFS remove_requests:109-116remove_requests : comme deque ne supporte pas le filter en place, on enchaîne clear() + extend().
  • PriorityRequestQueue:131 — maintient _heap: list[Request] via heapq.
  • Priority add/pop:144-152heappush à l'entrée, heappop à la sortie ; l'ordre est fixé par Request.__lt__.
  • Priority remove:175-184remove_request fait list.remove puis heapify, O(n) mais sémantiquement correct.
  • Priority __iter__:194-198 — à l'itération, copie le heap puis heappop un par un, afin de parcourir dans l'ordre priority sans détruire le tas original.
  • create_request_queue:201 — fonction factory, instancie la file correspondant au SchedulingPolicy.

Flux de données

Quand la boucle principale de l'ordonnanceur, à chaque step, doit piocher la prochaine request depuis la file waiting, elle appelle peek_request() :

python
request_queue = self._select_waiting_queue_for_scheduling()
assert request_queue is not None

request = request_queue.peek_request()
request_id = request.request_id

# try to promote blocked statuses while traversing skipped queue.
if self._is_blocked_waiting_status(
    request.status
) and not self._try_promote_blocked_waiting_request(request):
    if request.status == RequestStatus.WAITING_FOR_REMOTE_KV:
        logger.debug(
            "%s is still in WAITING_FOR_REMOTE_KV state.",
            request_id,
        )
    request_queue.pop_request()
    step_skipped_waiting.prepend_request(request)
    continue

Ce bloc est dans scheduler.py:647-664. À noter : on ne fait que peek_request, sans immédiatement pop — ce n'est qu'une fois que la request a passé toute la chaîne de contrôles (prefix cache, budget encoder, plafond LoRA, etc.) qu'elle est effectivement détachée via request_queue.pop_request()(pop_request:946) et promue dans running. En cas d'échec d'un contrôle intermédiaire, on l'écarte par pop_request() + step_skipped_waiting.prepend_request() vers la file skipped, sans bloquer les suivantes.

_select_waiting_queue_for_scheduling(_select_waiting_queue_for_scheduling:1867) décide si l'on tire ce tour de waiting ou de skipped_waiting. En mode FCFS, c'est simple : skipped_waiting or waiting or None, en priorisant les requests revenues de skipped. En mode priority, on compare les deux têtes de file sur (priority, arrival_time) et l'on prend la plus petite.

L'entrée en fin de file passe par Scheduler.add_request(add_request:2012), qui appelle _enqueue_waiting_request(_enqueue_waiting_request:1861), lequel selon l'état de la request route vers waiting ou skipped_waiting. Cycle de vie complet : entrée dans waiting → promotion dans running → préemption puis réinjection par waiting.prepend_request → re-promotion dans running. Toutes ces opérations de file transitent par l'interface RequestQueue.

Limites et échecs

  • peek_request lève IndexError sur file vide : FCFSRequestQueue.peek_request fait raise IndexError("peek from an empty queue")(FCFS peek:88-90), et PriorityRequestQueue pareil(Priority peek:156-158). L'ordonnanceur utilise _select_waiting_queue_for_scheduling pour ne jamais peek sur une file vide, mais l'implémentation elle-même ne se défend pas.
  • remove_request est en O(n) : FCFS via deque.remove (O(n)), priority via list.remove + heapify (O(n) + O(n) de re-heapify)(Priority remove:175-178) ; coûteux sur les longues files mais correct sémantiquement.
  • remove_requests ne déduplique pas : dans l'implémentation FCFS, filtered_requests = [req for req in self if req not in requests_to_remove](FCFS remove_requests:109-116) ; requests_to_remove doit être un set pour rester O(1). Passer une liste fonctionne mais devient O(n*m).
  • prepend_requests inverse l'ordre : FCFS utilise extendleft (la docstring rappelle que l'ordre des éléments prépended est inversé par rapport à la file d'origine)(FCFS prepend_requests:96-103) ; en file priority, on dégénère en add_request un par un, l'ordre n'a pas de sens(Priority prepend_requests:167-173).
  • __iter__ de la file priority copie : heap_copy = self._heap[:] puis heappop un par un(Priority __iter__:194-198) afin de parcourir dans l'ordre priority sans corrompre le tas, au prix d'une copie mémoire en O(n).
  • Stratégie inconnue lève ValueError : create_request_queue(create_request_queue:201-208) fait raise ValueError pour un SchedulingPolicy inconnu, et le constructeur de Scheduler ajoute une couche try/except pour relancer une erreur plus friendly(policy 校验:174-179).

Résumé

RequestQueue est l'abstraction de la salle d'attente de l'ordonnanceur : FCFSRequestQueue s'appuie sur deque, PriorityRequestQueue sur heap, derrière une interface unifiée. La boucle principale de Scheduler.schedule() ne parle qu'à cette interface, et la stratégie est fixée par scheduler_config.policy. Pour voir comment schedule utilise ces deux files, lire Scheduler.schedule ; pour comprendre comment une request préemptée est remise en tête de file, lire Préemption et prefill throttle.

Voir la documentation officielle : Documentation vLLM · README