RequestQueue : FCFS et file prioritaire
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 queFCFSetPRIORITY, mais la classe de base relègue les détails de tri dans l'implémentation. L'ordonnanceur ne parle qu'à l'interfaceRequestQueue; 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 lecollections.dequede Python, dont les append/popleft en O(1) des deux côtés sont l'implémentation optimale de FCFS.add_requestcorrespond àappend, le pop àpopleft(add/pop:78-84). - priority en heap :
PriorityRequestQueueutiliseheapq(heappush:144). L'ordre est fixé par le__lt__duRequestlui-même — d'abordpriorityascendant, puisarrival_timeascendant(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 », doncprepend_requestdégénère enadd_request(PriorityRequestQueue.prepend:160-165), comportement que la docstring explicite. skipped_waitingré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 dansskipped_waiting(step_skipped_waiting:638). C'est la même sorte de file quewaiting, 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:13—Enum, uniquementFCFSetPRIORITY.RequestQueue ABC:20— classe abstraite, définitadd_request/pop_request/peek_request/prepend_request/remove_request/__bool__/__len__/__iter__.FCFSRequestQueue:75— hérite dedeque[Request]+RequestQueue; toutes les opérations sont des appels natifs deque.FCFS add/pop:78-84—add_requestviaappend,pop_requestviapopleft, implémentation standard FCFS.FCFS prepend:92-94—prepend_requestviaappendleft, de sorte qu'une request réinjectée par préemption soit re-planifiée en priorité au step suivant.FCFS prepend_requests:96-103—prepend_requestsviaextendleft; la docstring signale que l'ordre des éléments prépended est inversé.FCFS remove_requests:109-116—remove_requests: comme deque ne supporte pas le filter en place, on enchaîneclear()+extend().PriorityRequestQueue:131— maintient_heap: list[Request]viaheapq.Priority add/pop:144-152—heappushà l'entrée,heappopà la sortie ; l'ordre est fixé parRequest.__lt__.Priority remove:175-184—remove_requestfaitlist.removepuisheapify, 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 auSchedulingPolicy.
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() :
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)
continueCe 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_requestlève IndexError sur file vide :FCFSRequestQueue.peek_requestfaitraise IndexError("peek from an empty queue")(FCFS peek:88-90), etPriorityRequestQueuepareil(Priority peek:156-158). L'ordonnanceur utilise_select_waiting_queue_for_schedulingpour ne jamais peek sur une file vide, mais l'implémentation elle-même ne se défend pas.remove_requestest en O(n) : FCFS viadeque.remove(O(n)), priority vialist.remove+heapify(O(n) + O(n) de re-heapify)(Priority remove:175-178) ; coûteux sur les longues files mais correct sémantiquement.remove_requestsne 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_removedoit être un set pour rester O(1). Passer une liste fonctionne mais devient O(n*m).prepend_requestsinverse l'ordre : FCFS utiliseextendleft(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 enadd_requestun 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) faitraise ValueErrorpour unSchedulingPolicyinconnu, 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