Skip to content
Closed
Changes from 1 commit
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
46 changes: 35 additions & 11 deletions vllm/v1/core/sched/request_queue.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@
from __future__ import annotations

import heapq
import time
from abc import ABC, abstractmethod
from collections import deque
from collections.abc import Iterable, Iterator
Expand Down Expand Up @@ -136,36 +137,57 @@ def __reversed__(self) -> Iterator[Request]:
return super().__reversed__()


class PrioritizedItem:

def __init__(self, request: Request, aging_factor: float = 0.1):
self.request = request
self.aging_factor = aging_factor
self.insert_time = request.arrival_time

def __lt__(self, other: PrioritizedItem) -> bool:
now = time.time()
eff_self = self.request.priority - self.aging_factor * (
now - self.insert_time)
eff_other = other.request.priority - other.aging_factor * (
now - other.insert_time)
Comment thread
chaunceyjiang marked this conversation as resolved.
Outdated

if eff_self != eff_other:
return eff_self < eff_other
return self.insert_time < other.insert_time
Comment thread
chaunceyjiang marked this conversation as resolved.
Outdated


class PriorityRequestQueue(RequestQueue):
"""
A priority queue that supports heap operations.

Requests with a smaller value of `priority` are processed first.
If multiple requests have the same priority, the one with the earlier
`arrival_time` is processed first.
Requests are aged over time based on the `aging_factor`, which
reduces their effective priority as time passes.
"""

def __init__(self) -> None:
self._heap: list[tuple[int, float, Request]] = []
def __init__(self, aging_factor: float = 0.1) -> None:
self._heap: list[PrioritizedItem] = []
self.aging_factor = aging_factor

def add_request(self, request: Request) -> None:
"""Add a request to the queue according to priority policy."""
heapq.heappush(self._heap,
(request.priority, request.arrival_time, request))
item = PrioritizedItem(request, self.aging_factor)
heapq.heappush(self._heap, item)

def pop_request(self) -> Request:
"""Pop a request from the queue according to priority policy."""
if not self._heap:
raise IndexError("pop from empty heap")
_, _, request = heapq.heappop(self._heap)
request = heapq.heappop(self._heap).request
return request

def peek_request(self) -> Request:
"""Peek at the next request in the queue without removing it."""
if not self._heap:
raise IndexError("peek from empty heap")
_, _, request = self._heap[0]
return request
return self._heap[0].request

def prepend_request(self, request: Request) -> None:
"""Add a request to the queue according to priority policy.
Expand All @@ -184,14 +206,16 @@ def prepend_requests(self, requests: RequestQueue) -> None:

def remove_request(self, request: Request) -> None:
"""Remove a specific request from the queue."""
self._heap = [(p, t, r) for p, t, r in self._heap if r != request]
self._heap = [item for item in self._heap if item.request != request]
heapq.heapify(self._heap)

def remove_requests(self, requests: Iterable[Request]) -> None:
"""Remove multiple specific requests from the queue."""
requests_to_remove = set(requests)
self._heap = [(p, t, r) for p, t, r in self._heap
if r not in requests_to_remove]
self._heap = [
item for item in self._heap
if item.request not in requests_to_remove
]
heapq.heapify(self._heap)

def __bool__(self) -> bool:
Expand All @@ -206,7 +230,7 @@ def __iter__(self) -> Iterator[Request]:
"""Iterate over the queue according to priority policy."""
heap_copy = self._heap[:]
while heap_copy:
_, _, request = heapq.heappop(heap_copy)
request = heapq.heappop(heap_copy).request
yield request

def __reversed__(self) -> Iterator[Request]:
Expand Down