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。