From 35913566bb403a49da81c49ea5a4c0761776cfb6 Mon Sep 17 00:00:00 2001 From: Jackie2049 Date: Thu, 4 Jun 2026 13:31:23 +0800 Subject: [PATCH] [Feature] Add request overload control with queue timeout eviction Adds a configurable max_waiting_time that evicts requests waiting too long in the scheduler queue, returning HTTP 429 to clients. Changes: - SchedulerConfig: add max_waiting_time field (default: 0 = disabled) - RequestStatus: add FINISHED_OVERLOAD status - FinishReason: add OVERLOAD reason - Scheduler: evict stale requests before each scheduling cycle - HTTP layer: return 429 for overloaded requests via OverloadError exception, registered in chat/completion/responses/streaming paths The feature is off by default and activated by setting --max-waiting-time > 0 in SchedulerConfig. Closes: Jackie2049/vllm#3 Co-authored-by: Boundless --- vllm/config/scheduler.py | 7 +++++++ vllm/entrypoints/openai/api_server.py | 3 ++- .../openai/chat_completion/serving.py | 3 +++ vllm/entrypoints/openai/completion/serving.py | 3 +++ vllm/entrypoints/openai/engine/protocol.py | 9 +++++++++ vllm/entrypoints/openai/engine/serving.py | 19 +++++++++++++++---- vllm/entrypoints/openai/responses/serving.py | 7 +++++++ vllm/v1/core/sched/scheduler.py | 19 +++++++++++++++++++ vllm/v1/engine/__init__.py | 3 ++- vllm/v1/request.py | 2 ++ 10 files changed, 69 insertions(+), 6 deletions(-) diff --git a/vllm/config/scheduler.py b/vllm/config/scheduler.py index 7900c948480b..aa4c7fe6c697 100644 --- a/vllm/config/scheduler.py +++ b/vllm/config/scheduler.py @@ -148,6 +148,13 @@ class SchedulerConfig: avoid gaps in GPU utilization, leading to better latency and throughput. """ + max_waiting_time: float = 0 + """Maximum time (in seconds) a request can wait in the scheduler's waiting + queue before being evicted. When set to 0 (default), the feature is + disabled and requests wait indefinitely. When set to a positive value, + requests exceeding this wait time are evicted with FINISHED_OVERLOAD status + and the client receives an HTTP 429 Too Many Requests response.""" + stream_interval: int = Field(default=1, ge=1) """The interval (or buffer size) for streaming in terms of token length. A smaller value (1) makes streaming smoother by sending each token immediately, diff --git a/vllm/entrypoints/openai/api_server.py b/vllm/entrypoints/openai/api_server.py index 892f9d82d709..320059615f6f 100644 --- a/vllm/entrypoints/openai/api_server.py +++ b/vllm/entrypoints/openai/api_server.py @@ -29,7 +29,7 @@ from vllm.entrypoints.launcher import serve_http from vllm.entrypoints.logger import RequestLogger from vllm.entrypoints.openai.cli_args import make_arg_parser, validate_parsed_serve_args -from vllm.entrypoints.openai.engine.protocol import GenerationError +from vllm.entrypoints.openai.engine.protocol import GenerationError, OverloadError from vllm.entrypoints.openai.models.protocol import BaseModelPath from vllm.entrypoints.openai.models.serving import OpenAIServingModels from vllm.entrypoints.openai.server_utils import ( @@ -250,6 +250,7 @@ def build_app( app.exception_handler(EngineGenerateError)(engine_error_handler) app.exception_handler(EngineDeadError)(engine_error_handler) app.exception_handler(GenerationError)(generation_error_handler) + app.exception_handler(OverloadError)(generation_error_handler) app.exception_handler(Exception)(exception_handler) # Ensure --api-key option from CLI takes precedence over VLLM_API_KEY diff --git a/vllm/entrypoints/openai/chat_completion/serving.py b/vllm/entrypoints/openai/chat_completion/serving.py index a378fb79d3bc..21bbb1ec93e5 100644 --- a/vllm/entrypoints/openai/chat_completion/serving.py +++ b/vllm/entrypoints/openai/chat_completion/serving.py @@ -50,6 +50,7 @@ from vllm.entrypoints.openai.engine.serving import ( GenerationError, OpenAIServing, + OverloadError, clamp_prompt_logprobs, ) from vllm.entrypoints.openai.models.serving import OpenAIServingModels @@ -926,6 +927,8 @@ async def chat_completion_stream_generator( except GenerationError as e: yield f"data: {self._convert_generation_error_to_streaming_response(e)}\n\n" + except OverloadError as e: + yield f"data: {self._convert_generation_error_to_streaming_response(e)}\n\n" except Exception as e: logger.exception("Error in chat completion stream generator.") data = self.create_streaming_error_response(e) diff --git a/vllm/entrypoints/openai/completion/serving.py b/vllm/entrypoints/openai/completion/serving.py index f393954e2a05..bca789c85d03 100644 --- a/vllm/entrypoints/openai/completion/serving.py +++ b/vllm/entrypoints/openai/completion/serving.py @@ -31,6 +31,7 @@ from vllm.entrypoints.openai.engine.serving import ( GenerationError, OpenAIServing, + OverloadError, clamp_prompt_logprobs, ) from vllm.entrypoints.openai.models.serving import OpenAIServingModels @@ -467,6 +468,8 @@ async def completion_stream_generator( except GenerationError as e: yield f"data: {self._convert_generation_error_to_streaming_response(e)}\n\n" + except OverloadError as e: + yield f"data: {self._convert_generation_error_to_streaming_response(e)}\n\n" except Exception as e: logger.exception("Error in completion stream generator.") data = self.create_streaming_error_response(e) diff --git a/vllm/entrypoints/openai/engine/protocol.py b/vllm/entrypoints/openai/engine/protocol.py index 434888df9efa..4f35fe0ee90e 100644 --- a/vllm/entrypoints/openai/engine/protocol.py +++ b/vllm/entrypoints/openai/engine/protocol.py @@ -352,3 +352,12 @@ class GenerationError(Exception): def __init__(self, message: str = "Internal server error"): super().__init__(message) self.status_code = HTTPStatus.INTERNAL_SERVER_ERROR + + +class OverloadError(Exception): + """raised when finish_reason indicates overload (429)""" + + def __init__(self, message: str = "Request exceeded maximum queue waiting " + "time. Please retry later."): + super().__init__(message) + self.status_code = HTTPStatus.TOO_MANY_REQUESTS diff --git a/vllm/entrypoints/openai/engine/serving.py b/vllm/entrypoints/openai/engine/serving.py index 61b2656bac0f..f83261e96a50 100644 --- a/vllm/entrypoints/openai/engine/serving.py +++ b/vllm/entrypoints/openai/engine/serving.py @@ -29,6 +29,7 @@ from vllm.entrypoints.openai.engine.protocol import ( ErrorResponse, GenerationError, + OverloadError, ) from vllm.entrypoints.openai.models.serving import OpenAIServingModels from vllm.entrypoints.openai.responses.protocol import ResponsesRequest @@ -190,21 +191,31 @@ def create_streaming_error_response( return json_str def _raise_if_error(self, finish_reason: str | None, request_id: str) -> None: - """Raise GenerationError if finish_reason indicates an error.""" + """Raise appropriate error if finish_reason indicates an error.""" if finish_reason == "error": logger.error( "Request %s failed with an internal error during generation", request_id, ) raise GenerationError("Internal server error") + if finish_reason == "overload": + logger.warning( + "Request %s evicted from queue due to overload", + request_id, + ) + raise OverloadError() def _convert_generation_error_to_streaming_response( - self, e: GenerationError + self, e: GenerationError | OverloadError ) -> str: - """Convert GenerationError to streaming error response.""" + """Convert GenerationError/OverloadError to streaming error response.""" + if isinstance(e, OverloadError): + err_type = "OverloadError" + else: + err_type = "InternalServerError" return self.create_streaming_error_response( str(e), - err_type="InternalServerError", + err_type=err_type, status_code=e.status_code, ) diff --git a/vllm/entrypoints/openai/responses/serving.py b/vllm/entrypoints/openai/responses/serving.py index eee02707a979..f975c901ec75 100644 --- a/vllm/entrypoints/openai/responses/serving.py +++ b/vllm/entrypoints/openai/responses/serving.py @@ -42,6 +42,7 @@ from vllm.entrypoints.openai.engine.serving import ( GenerationError, OpenAIServing, + OverloadError, ) from vllm.entrypoints.openai.models.serving import OpenAIServingModels from vllm.entrypoints.openai.parser.harmony_utils import ( @@ -1553,6 +1554,12 @@ def _increment_sequence_number_and_return( TypeAdapter(StreamingResponsesResponse).validate_json(error_json) ) return + except OverloadError as e: + error_json = self._convert_generation_error_to_streaming_response(e) + yield _increment_sequence_number_and_return( + TypeAdapter(StreamingResponsesResponse).validate_json(error_json) + ) + return async def empty_async_generator(): # A hack to trick Python to think this is a generator but diff --git a/vllm/v1/core/sched/scheduler.py b/vllm/v1/core/sched/scheduler.py index c39e80c24eb0..718252932a09 100644 --- a/vllm/v1/core/sched/scheduler.py +++ b/vllm/v1/core/sched/scheduler.py @@ -108,6 +108,7 @@ def __init__( else self.scheduler_config.max_num_batched_tokens ) self.max_model_len = vllm_config.model_config.max_model_len + self.max_waiting_time = self.scheduler_config.max_waiting_time self.enable_kv_cache_events = ( self.kv_events_config is not None and self.kv_events_config.enable_kv_cache_events @@ -558,6 +559,24 @@ def schedule(self) -> SchedulerOutput: ) assert len(scheduled_loras) <= self.lora_config.max_loras + # Evict requests that have waited too long in the queue. + if self.max_waiting_time > 0: + now = time.monotonic() + evicted_ids: list[str] = [] + for queue in (self.waiting, self.skipped_waiting): + remaining = create_request_queue(self.policy) + while queue: + req = queue.pop_request() + if (now - req.arrival_time) > self.max_waiting_time: + evicted_ids.append(req.request_id) + else: + remaining.append_request(req) + while remaining: + queue.append_request(remaining.pop_request()) + for req_id in evicted_ids: + self.finish_requests(req_id, + RequestStatus.FINISHED_OVERLOAD) + # Next, schedule the WAITING requests. if not preempted_reqs and self._pause_state == PauseState.UNPAUSED: step_skipped_waiting = create_request_queue(self.policy) diff --git a/vllm/v1/engine/__init__.py b/vllm/v1/engine/__init__.py index 848f530ce334..f0a497b7b8ff 100644 --- a/vllm/v1/engine/__init__.py +++ b/vllm/v1/engine/__init__.py @@ -27,7 +27,7 @@ # These are possible values of RequestOutput.finish_reason, # so form part of the external API. -FINISH_REASON_STRINGS = ("stop", "length", "abort", "error", "repetition") +FINISH_REASON_STRINGS = ("stop", "length", "abort", "error", "repetition", "overload") EEP_NOTIFICATION_CALL_ID = -1 @@ -59,6 +59,7 @@ class FinishReason(enum.IntEnum): ABORT = 2 ERROR = 3 REPETITION = 4 + OVERLOAD = 5 def __str__(self): return FINISH_REASON_STRINGS[self.value] diff --git a/vllm/v1/request.py b/vllm/v1/request.py index 44246e70a8bb..f5422ac0e7a9 100644 --- a/vllm/v1/request.py +++ b/vllm/v1/request.py @@ -333,6 +333,7 @@ class RequestStatus(enum.IntEnum): FINISHED_IGNORED = enum.auto() FINISHED_ERROR = enum.auto() FINISHED_REPETITION = enum.auto() + FINISHED_OVERLOAD = enum.auto() def __str__(self) -> str: return self.name @@ -358,4 +359,5 @@ def get_finished_reason(status: "RequestStatus") -> FinishReason | None: RequestStatus.FINISHED_ERROR: FinishReason.ERROR, RequestStatus.WAITING_FOR_STREAMING_REQ: FinishReason.STOP, RequestStatus.FINISHED_REPETITION: FinishReason.REPETITION, + RequestStatus.FINISHED_OVERLOAD: FinishReason.OVERLOAD, }