Skip to content

RequestQueue: FCFS und Priority-Warteschlange

源码版本v0.25.1

Verantwortung

Scheduler verwaltet intern drei Warteschlangen: waiting, skipped_waiting, running (waiting/running:181-184). Die ersten beiden sind Instanzen von RequestQueue und sind zuständig, Requests einzuordnen, bevor sie KV-Blöcke erhalten. RequestQueue (RequestQueue:20) selbst ist eine abstrakte Basisklasse, die einen allgemeinen Satz an Methoden definiert: add_request, pop_request, peek_request, prepend_request, remove_request und weitere; die konkrete Reihenfolge wird durch zwei Implementierungen bestimmt: FCFSRequestQueue (FCFSRequestQueue:75) und PriorityRequestQueue (PriorityRequestQueue:131).

Die Existenzberechtigung der Warteschlange ist es, dem Scheduler einen einheitlichen Eingang „nimm den nächsten einzuplanenden Request" zu bieten — egal ob darunter ein deque oder ein heap liegt, die Hauptschleife des Schedulers ruft nur peek_request() / pop_request() auf (peek_request:650); die Einreihungsstrategie wird durch scheduler_config.policy bestimmt und in Scheduler.__init__ über create_request_queue(self.policy) (create_request_queue:181) injiziert.

Entwurfsmotivation

Warum die Warteschlange abstrahieren, statt direkt list[Request] zu verwenden?

  • Strategie austauschbar: Das Enum SchedulingPolicy (SchedulingPolicy:13) listet nur FCFS und PRIORITY, aber die abstrakte Basisklasse verlagert die Sortierdetails in die Implementierung. Der Scheduler interagiert nur mit dem RequestQueue-Interface; eine neue Strategie (z. B. Fair Scheduling) erfordert nur eine neue Unterklasse.
  • FCFS via deque: FCFSRequestQueue(deque[Request], RequestQueue) (FCFSRequestQueue:75) verwendet Pythons collections.deque direkt; append/popleft in O(1) auf beiden Seiten ist die optimale Implementierung für FCFS. add_request ist append, pop ist popleft (add/pop:78-84).
  • Priority via heap: PriorityRequestQueue nutzt heapq (heappush:144) — die Reihenfolge wird durch Request.__lt__ bestimmt: zuerst aufsteigend nach priority, dann aufsteigend nach arrival_time (docstring:131-139) — sodass Nutzer jedem Request eine Priorität geben können.
  • Eindeutige prepend-Semantik: Nach einer Preemption oder der Auflösung einer asynchronen Abhängigkeit muss ein Request vorne wieder eingereiht werden, um als Nächstes neu eingeplant zu werden. Beide Warteschlangen bieten prepend_request, aber in der Priority-Warteschlange gibt es keinen Begriff eines „Kopfes", sodass prepend_request dort auf add_request zurückfällt (PriorityRequestQueue.prepend:160-165); das wird im docstring klargestellt.
  • skipped_waiting verwendet denselben Typ: Requests, die wegen LoRA-Überschreitung oder asynchronem KV-Laden übersprungen wurden, werden in skipped_waiting (step_skipped_waiting:638) gepackt; das ist dieselbe Art Warteschlange wie waiting, und die Hauptschleife wählt über _select_waiting_queue_for_scheduling() diejenige mit dem früheren Kopf aus (_select_waiting_queue_for_scheduling:1867-1877).

Schlüsseldateien

  • SchedulingPolicy:13Enum, nur FCFS und PRIORITY.
  • RequestQueue ABC:20 — Abstrakte Basisklasse; definiert die Schnittstellen add_request / pop_request / peek_request / prepend_request / remove_request / __bool__ / __len__ / __iter__.
  • FCFSRequestQueue:75 — Erbt von deque[Request] + RequestQueue; alle Operationen sind native deque-Aufrufe.
  • FCFS add/pop:78-84add_request via append, pop_request via popleft; die Standardimplementierung für FCFS.
  • FCFS prepend:92-94prepend_request via appendleft, damit ein durch Preemption zurückgekehrter Request im nächsten Schritt vorrangig neu eingeplant wird.
  • FCFS prepend_requests:96-103prepend_requests via extendleft; beachten Sie, dass das docstring darauf hinweist, dass die Reihenfolge der eingefügten Elemente umgekehrt ist.
  • FCFS remove_requests:109-116remove_requests macht clear() + extend(), da deque kein in-place filter unterstützt.
  • PriorityRequestQueue:131 — Verwaltet _heap: list[Request] via heapq.
  • Priority add/pop:144-152heappush zum Einreihen, heappop zum Entnehmen; die Reihenfolge wird durch Request.__lt__ bestimmt.
  • Priority remove:175-184remove_request macht erst list.remove, dann heapify; O(n), aber semantisch korrekt.
  • Priority __iter__:194-198 — Beim Iterieren wird eine Kopie des Heaps angelegt und dann elementweise heappop ausgeführt, damit die Iteration nach Priority sortiert erfolgt, ohne den Original-Heap zu zerstören.
  • create_request_queue:201 — Factory-Funktion; instanziiert je nach SchedulingPolicy die entsprechende Warteschlange.

Datenfluss

Wenn die Hauptschleife des Schedulers in einem Schritt den nächsten Request aus der waiting-Warteschlange wählen muss, wird peek_request() aufgerufen:

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

Dieser Abschnitt steht in scheduler.py:647-664. Beachte: Es wird nur peek_request ausgeführt und nicht sofort gepoppt — erst wenn der Request die lange Kette aus Prefix-Cache-Hit, Encoder-Budget und LoRA-Limit bestanden hat, wird er unter pop_request:946 mit request_queue.pop_request() entnommen und nach running befördert. Schlägt eine der Prüfungen fein, wird er mit pop_request() + step_skipped_waiting.prepend_request() in die skipped-Warteschlange verschoben, sodass er nachfolgende Requests nicht blockiert.

_select_waiting_queue_for_scheduling (_select_waiting_queue_for_scheduling:1867) entscheidet, ob diese Runde aus waiting oder skipped_waiting genommen wird. Im FCFS-Modus ist es einfach: skipped_waiting or waiting or None — die aus skipped zurückgekehrten Requests werden vorrangig eingeplant. Im Priority-Modus werden die Köpfe beider Warteschlangen verglichen und der mit dem kleineren (priority, arrival_time) gewählt.

Der Eintritt am Ende der Warteschlange erfolgt über Scheduler.add_request (add_request:2012), der _enqueue_waiting_request (_enqueue_waiting_request:1861) aufruft und dann anhand des Request-Status entscheidet, ob die Einordnung in waiting oder skipped_waiting erfolgt. Der gesamte Lebenszyklus: Eintritt in waiting → Beförderung zu running → durch Preemption zurück an waiting.prepend_request → erneute Beförderung zu running; alle Warteschlangenoperationen laufen über die RequestQueue-Schnittstelle.

Grenzen und Fehler

  • peek_request auf leerer Warteschlange wirft IndexError: FCFSRequestQueue.peek_request macht bei leerer Warteschlange raise IndexError("peek from an empty queue") (FCFS peek:88-90); für PriorityRequestQueue gilt Entsprechendes (Priority peek:156-158). Der Scheduler stellt über _select_waiting_queue_for_scheduling sicher, dass nicht auf einer leeren Warteschlange gepeekt wird; die Implementierung selbst ist nicht defensiv.
  • remove_request ist O(n): FCFS geht über deque.remove (O(n)), Priority über list.remove + heapify (O(n) + O(n) reorder) (Priority remove:175-178) — für lange Warteschlangen mit Kosten verbunden, aber semantisch korrekt.
  • remove_requests dedupliziert nicht: In der FCFS-Implementierung ist filtered_requests = [req for req in self if req not in requests_to_remove] (FCFS remove_requests:109-116); requests_to_remove muss ein set sein, damit die Prüfung O(1) ist — übergibt der Aufrufer eine list, wird daraus O(n*m).
  • prepend_requests mit umgekehrter Reihenfolge: FCFS verwendet extendleft (das docstring weist darauf hin, dass die eingefügten Elemente in umgekehrter Reihenfolge erscheinen) (FCFS prepend_requests:96-103); die Priority-Warteschlange fällt auf einzelnes add_request zurück, die Reihenfolge ist bedeutungslos (Priority prepend_requests:167-173).
  • Priority-Warteschlange kopiert bei __iter__: heap_copy = self._heap[:] und dann elementweise heappop (Priority __iter__:194-198), damit die Iteration nach Priority sortiert erfolgt, ohne den Original-Heap zu zerstören — der Preis ist eine O(n) Speicherkopie.
  • Unbekannte Strategie wirft ValueError: create_request_queue (create_request_queue:201-208) macht bei unbekanntem SchedulingPolicy direkt raise ValueError; im Konstruktor von Scheduler ist zusätzlich ein try/except gewrappt, das einen freundlicheren Fehler weiterreicht (policy Prüfung:174-179).

Zusammenfassung

RequestQueue ist die Abstraktion des Wartebereichs des Schedulers: FCFSRequestQueue basiert auf deque, PriorityRequestQueue auf heap, die Schnittstellen sind identisch. Die Hauptschleife von Scheduler.schedule() interagiert nur mit dieser Schnittstelle; die Einreihungsstrategie wird über scheduler_config.policy bestimmt. Wie schedule diese beiden Warteschlangen verwendet, siehe Scheduler.schedule; wie ein durch Preemption gekickter Request wieder an die Spitze der Warteschlange gelangt, siehe Preemption und Prefill-Throttle.

Siehe offizielle Dokumentation: vLLM-Dokumentation · README.