RequestQueue:FCFS 與優先級佇列
職責
Scheduler 內部維護三條佇列:waiting、skipped_waiting、running(waiting/running:181-184)。前兩條就是 RequestQueue 的實例,負責在 request 還沒拿到 KV 塊之前給它們排隊。RequestQueue(RequestQueue:20)本身是一個抽象基類,定義了 add_request、pop_request、peek_request、prepend_request、remove_request 等一套通用介面,具體順序由兩個實現決定:FCFSRequestQueue(FCFSRequestQueue:75)和 PriorityRequestQueue(PriorityRequestQueue:131)。
佇列的存在意義是為排程器提供一個統一的"取下一個該排的 request"的入口——不管底層是 deque 還是 heap,排程主循環都只調 peek_request() / pop_request()(peek_request:650),排隊策略由 scheduler_config.policy 決定,在 Scheduler.__init__ 裡通過 create_request_queue(self.policy)(create_request_queue:181)注入。
設計動機
為什麼把佇列抽象出來,而不是直接用 list[Request]?
- 策略可換:
SchedulingPolicy枚舉(SchedulingPolicy:13)只列FCFS和PRIORITY兩種,但抽象基類把排序細節藏到實現裡。排程器只跟RequestQueue介面打交道,加新策略(比如公平排程)只要新寫一個子類。 - FCFS 走 deque:
FCFSRequestQueue(deque[Request], RequestQueue)(FCFSRequestQueue:75)直接複用 Pythoncollections.deque,兩端 O(1) 的 append/popleft 正好是 FCFS 的最優實現。add_request就是append,pop 是popleft(add/pop:78-84)。 - priority 走 heap:
PriorityRequestQueue用heapq(heappush:144),順序由Request自己的__lt__決定——先按priority升序、再按arrival_time升序(docstring:131-139),讓使用者給每個 request 打優先級。 - prepend 語義統一:搶佔或非同步依賴解除後,需要把 request 塞回隊首優先重排。兩種佇列都提供
prepend_request,但 priority 佇列裡沒有"隊首"概念,所以prepend_request退化成add_request(PriorityRequestQueue.prepend:160-165),用 docstring 明確寫出來。 skipped_waiting複用同一類型:被 LoRA 超限、非同步 KV 載入等原因跳過的 request 會被塞進skipped_waiting(step_skipped_waiting:638),它跟waiting是同一種佇列,排程主循環通過_select_waiting_queue_for_scheduling()在兩者之間挑隊頭更早的那個(_select_waiting_queue_for_scheduling:1867-1877)。
關鍵檔案
SchedulingPolicy:13—Enum,只有FCFS和PRIORITY。RequestQueue ABC:20— 抽象基類,定義add_request/pop_request/peek_request/prepend_request/remove_request/__bool__/__len__/__iter__介面。FCFSRequestQueue:75— 繼承deque[Request]+RequestQueue,所有操作都是 deque 原生呼叫。FCFS add/pop:78-84—add_request走append,pop_request走popleft,FCFS 的標準實現。FCFS prepend:92-94—prepend_request用appendleft,讓被搶佔回來的 request 在下一步優先重排。FCFS prepend_requests:96-103—prepend_requests用extendleft,注意 docstring 提到 prepended 順序是反的。FCFS remove_requests:109-116—remove_requests因為 deque 不支持就地 filter,所以clear()+extend()。PriorityRequestQueue:131— 用heapq維護_heap: list[Request]。Priority add/pop:144-152—heappush入隊,heappop出隊,順序由Request.__lt__決定。Priority remove:175-184—remove_request先 list.remove 再heapify,O(n) 但語義正確。Priority __iter__:194-198— 迭代時拷一份 heap 再逐個 heappop,保證按 priority 順序遍歷而不破壞原堆。create_request_queue:201— 工廠函數,按SchedulingPolicy實例化對應佇列。
資料流
排程主循環在每個 step 裡需要從 waiting 佇列挑下一個 request 排程時,調的就是 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)
continue這段在 scheduler.py:647-664。注意這裡只 peek_request 不立刻 pop——只有當 request 真的通過了 prefix cache、encoder 預算、LoRA 上限等一長串檢查後,才在 pop_request:946 處 request_queue.pop_request() 摘下來升入 running。中間任何檢查失敗都用 pop_request() + step_skipped_waiting.prepend_request() 把它挪到 skipped 佇列,不擋後面的人。
_select_waiting_queue_for_scheduling(_select_waiting_queue_for_scheduling:1867)決定這一輪從 waiting 還是 skipped_waiting 取。FCFS 模式簡單:skipped_waiting or waiting or None,優先排 skipped 裡回來的 request。priority 模式要比較兩個隊頭的 (priority, arrival_time),挑更小的那個。
隊尾入隊的路徑在 Scheduler.add_request(add_request:2012),它調 _enqueue_waiting_request(_enqueue_waiting_request:1861),再根據 request 狀態決定走 waiting 還是 skipped_waiting。整個生命週期:進 waiting → 升 running → 被搶佔回 waiting.prepend_request → 重新升 running,所有佇列操作都走 RequestQueue 這層介面。
邊界與失敗
peek_request空佇列拋 IndexError:FCFSRequestQueue.peek_request在空時raise IndexError("peek from an empty queue")(FCFS peek:88-90),PriorityRequestQueue同理(Priority peek:156-158)。排程器用_select_waiting_queue_for_scheduling保證不會在空佇列上 peek,但實現本身不防禦。remove_request是 O(n):FCFS 走deque.remove(O(n)),priority 走list.remove+heapify(O(n) + O(n) 重排)(Priority remove:175-178),對長佇列有開銷但語義正確。remove_requests不去重:FCFS 實現裡filtered_requests = [req for req in self if req not in requests_to_remove](FCFS remove_requests:109-116),requests_to_remove必須是 set 才 O(1);呼叫方傳 list 也行但變成 O(n*m)。prepend_requests順序反:FCFS 用extendleft(docstring 提醒 prepended 順序與原佇列出現順序相反)(FCFS prepend_requests:96-103),priority 佇列直接退化成逐個add_request,順序無意義(Priority prepend_requests:167-173)。- priority 佇列
__iter__拷貝:heap_copy = self._heap[:]再逐個 heappop(Priority __iter__:194-198),保證遍歷順序按 priority 但不破壞原堆,代價是 O(n) 記憶體拷貝。 - 未知策略 raise ValueError:
create_request_queue(create_request_queue:201-208)對未知SchedulingPolicy直接raise ValueError, Scheduler 構造函數里也包了一層 try/except 重拋更友好的錯誤(policy 校验:174-179)。
小結
RequestQueue 是排程器的等待區抽象,FCFSRequestQueue 走 deque、PriorityRequestQueue 走 heap,介面統一。Scheduler.schedule() 的主循環只跟這套介面打交道,排隊策略由 scheduler_config.policy 決定。要看 schedule 怎麼用這兩條佇列,讀 Scheduler.schedule;要看搶佔時怎麼把 request 塞回隊首,讀 搶佔與 prefill throttle。