From ecff19724203bc923f186a805e66cb2850f9cfb7 Mon Sep 17 00:00:00 2001 From: codechrl Date: Sun, 27 Sep 2026 00:29:05 +0700 Subject: [PATCH 1/2] [serve] Add request_routing_timeout_s to bound replica assignment _choose_replica_for_request awaits the pending request's future with no bound. A request waits with no error while every replica is at capacity. It also waits while the deployment has no replicas, for example when a scaled-to-zero deployment can't schedule a new replica. HTTPOptions.request_timeout_s applies to the HTTP proxy only and also counts processing time. Add request_routing_timeout_s to RequestRouterConfig. When set, a request that isn't assigned a replica within the timeout raises TimeoutError. The timeout does not cover the time the replica takes to handle the request. The default of None keeps the current behavior. RequestRouterConfig equality includes the new field, so a timeout-only config change is broadcast to routers. Closes #40995 Signed-off-by: codechrl --- .../serve/advanced-guides/performance.md | 17 +++++++++ .../_private/request_router/request_router.py | 26 +++++++++++++- python/ray/serve/_private/router.py | 14 +++++++- python/ray/serve/config.py | 12 +++++++ python/ray/serve/tests/unit/test_config.py | 19 ++++++++++ .../tests/unit/test_pow_2_request_router.py | 36 +++++++++++++++++++ python/ray/serve/tests/unit/test_router.py | 2 ++ src/ray/protobuf/serve.proto | 3 ++ 8 files changed, 127 insertions(+), 2 deletions(-) diff --git a/doc/source/serve/advanced-guides/performance.md b/doc/source/serve/advanced-guides/performance.md index d5d0757e75c5..0789054a746e 100644 --- a/doc/source/serve/advanced-guides/performance.md +++ b/doc/source/serve/advanced-guides/performance.md @@ -98,6 +98,23 @@ Ray Serve allows you to fine-tune the backoff behavior of the request router, wh - `RAY_SERVE_ROUTER_RETRY_BACKOFF_MULTIPLIER`: The multiplier applied to the backoff time after each retry. Default is `2`. - `RAY_SERVE_ROUTER_RETRY_MAX_BACKOFF_S`: The maximum backoff time (in seconds) between retries. Default is `0.5`. +### Set a timeout for choosing a replica + +By default, a request waits in the router until a replica accepts it. If no replica becomes available, for example because a deployment scaled to zero can't schedule a new replica, the request waits indefinitely. Set `request_routing_timeout_s` in the deployment's `request_router_config` to bound this wait: + +```python +from ray import serve +from ray.serve.config import RequestRouterConfig + +@serve.deployment( + request_router_config=RequestRouterConfig(request_routing_timeout_s=10), +) +class Model: + ... +``` + +A request that isn't assigned to a replica within the timeout fails with a `TimeoutError`. The timeout doesn't cover the time the replica takes to process the request. Unlike `request_timeout_s`, it also applies to `DeploymentHandle` calls. + ### Set timeouts while probing replicas for queue length Ray Serve's request router probes replicas for their queue lengths to make intelligent load balancing decisions. You can tune the following environment variables to optimize this behavior for your workload: diff --git a/python/ray/serve/_private/request_router/request_router.py b/python/ray/serve/_private/request_router/request_router.py index 4489d0d19d9d..1889d65f16d2 100644 --- a/python/ray/serve/_private/request_router/request_router.py +++ b/python/ray/serve/_private/request_router/request_router.py @@ -495,6 +495,7 @@ def __init__( initial_backoff_s: float = RAY_SERVE_ROUTER_RETRY_INITIAL_BACKOFF_S, backoff_multiplier: float = RAY_SERVE_ROUTER_RETRY_BACKOFF_MULTIPLIER, max_backoff_s: float = RAY_SERVE_ROUTER_RETRY_MAX_BACKOFF_S, + request_routing_timeout_s: Optional[float] = None, *args, **kwargs, ): @@ -510,6 +511,8 @@ def __init__( self.backoff_multiplier = backoff_multiplier self.max_backoff_s = max_backoff_s + self.request_routing_timeout_s = request_routing_timeout_s + # Current replicas available to be routed. # Updated via `update_replicas`. self._replica_id_set: Set[ReplicaID] = set() @@ -1405,13 +1408,34 @@ async def _choose_replica_for_request( self._add_pending_request_to_indices(pending_request) self._maybe_start_routing_tasks() - replica = await pending_request.future + if self.request_routing_timeout_s is None: + replica = await pending_request.future + else: + replica = await asyncio.wait_for( + pending_request.future, self.request_routing_timeout_s + ) except asyncio.CancelledError as e: pending_request.future.cancel() self._cancel_routing_task_for_pending_request(pending_request) self._remove_pending_request_from_indices(pending_request) raise e from None + except asyncio.TimeoutError: + self._cancel_routing_task_for_pending_request(pending_request) + self._remove_pending_request_from_indices(pending_request) + # No routing task runs while the deployment has no replicas, so the + # lazy cleanup in `_fulfill_pending_requests` can't drop the request. + for queue in ( + self._pending_requests_to_fulfill, + self._pending_requests_to_route, + ): + while queue and queue[0].future.done(): + queue.popleft() + + raise TimeoutError( + f"Failed to route request to a replica of {self._deployment_id} " + f"within {self.request_routing_timeout_s}s." + ) from None return replica diff --git a/python/ray/serve/_private/router.py b/python/ray/serve/_private/router.py index 3b77915351e9..710df591e293 100644 --- a/python/ray/serve/_private/router.py +++ b/python/ray/serve/_private/router.py @@ -667,6 +667,7 @@ def __init__( self._initial_backoff_s: Optional[float] = None self._backoff_multiplier: Optional[float] = None self._max_backoff_s: Optional[float] = None + self._request_routing_timeout_s: Optional[float] = None # Initializing `self._metrics_manager` before `self.long_poll_client` is # necessary to avoid race condition where `self.update_deployment_config()` @@ -781,6 +782,10 @@ def request_router(self) -> Optional[RequestRouter]: backoff_kwargs["backoff_multiplier"] = self._backoff_multiplier if self._max_backoff_s is not None: backoff_kwargs["max_backoff_s"] = self._max_backoff_s + if self._request_routing_timeout_s is not None: + backoff_kwargs[ + "request_routing_timeout_s" + ] = self._request_routing_timeout_s request_router = self._request_router_class( deployment_id=self.deployment_id, @@ -855,6 +860,9 @@ def update_deployment_config(self, deployment_config: DeploymentConfig): deployment_config.request_router_config.backoff_multiplier ) self._max_backoff_s = deployment_config.request_router_config.max_backoff_s + self._request_routing_timeout_s = ( + deployment_config.request_router_config.request_routing_timeout_s + ) if self._request_router: self._request_router.update_backoff_params( @@ -862,6 +870,9 @@ def update_deployment_config(self, deployment_config: DeploymentConfig): backoff_multiplier=self._backoff_multiplier, max_backoff_s=self._max_backoff_s, ) + self._request_router.request_routing_timeout_s = ( + self._request_routing_timeout_s + ) # Guard against the case where request_router is None (e.g., when # request_router_class is None and lazy initialization has not yet @@ -1139,7 +1150,8 @@ async def route_and_send_request( """Choose a replica for the request and send it. This will block indefinitely if no replicas are available to handle the - request, so it's up to the caller to time out or cancel the request. + request and `request_routing_timeout_s` is unset, so it's up to the + caller to time out or cancel the request. """ # Wait for the router to be initialized before sending the request. await self._request_router_initialized.wait() diff --git a/python/ray/serve/config.py b/python/ray/serve/config.py index 3a91d3ac52d1..6c11163062bb 100644 --- a/python/ray/serve/config.py +++ b/python/ray/serve/config.py @@ -274,6 +274,16 @@ class RequestRouterConfig(BaseModel): ), ) + request_routing_timeout_s: Optional[PositiveFloat] = Field( + default=None, + description=( + "Maximum duration in seconds that a request waits to be assigned a " + "replica before it fails with a TimeoutError. This bounds routing " + "only, not the time the replica takes to handle the request. " + "Defaults to None, meaning a request waits indefinitely." + ), + ) + @field_validator("request_router_kwargs") @classmethod def request_router_kwargs_json_serializable(cls, v): @@ -308,6 +318,7 @@ def __eq__(self, other): and self.initial_backoff_s == other.initial_backoff_s and self.backoff_multiplier == other.backoff_multiplier and self.max_backoff_s == other.max_backoff_s + and self.request_routing_timeout_s == other.request_routing_timeout_s ) def __hash__(self): @@ -329,6 +340,7 @@ def __hash__(self): self.initial_backoff_s, self.backoff_multiplier, self.max_backoff_s, + self.request_routing_timeout_s, ) ) diff --git a/python/ray/serve/tests/unit/test_config.py b/python/ray/serve/tests/unit/test_config.py index 2c3791a38239..3ea9b3674209 100644 --- a/python/ray/serve/tests/unit/test_config.py +++ b/python/ray/serve/tests/unit/test_config.py @@ -411,6 +411,25 @@ def test_backoff_params_declarative_schema(self): assert schema.request_router_config.backoff_multiplier == 3.0 assert schema.request_router_config.max_backoff_s == 2.0 + @pytest.mark.parametrize("request_routing_timeout_s", [None, 1.5]) + def test_request_routing_timeout_proto_round_trip(self, request_routing_timeout_s): + """None must not become 0.0 on the way through the proto.""" + config = DeploymentConfig.from_default( + request_router_config=RequestRouterConfig( + request_routing_timeout_s=request_routing_timeout_s + ) + ) + round_tripped = DeploymentConfig.from_proto_bytes(config.to_proto_bytes()) + assert ( + round_tripped.request_router_config.request_routing_timeout_s + == request_routing_timeout_s + ) + + def test_request_routing_timeout_changes_equality(self): + """A timeout-only change must be broadcast to routers.""" + config = RequestRouterConfig(request_routing_timeout_s=1.5) + assert config != RequestRouterConfig() + def test_deployment_actors_config(self): """Test deployment_actors config and proto roundtrip.""" diff --git a/python/ray/serve/tests/unit/test_pow_2_request_router.py b/python/ray/serve/tests/unit/test_pow_2_request_router.py index 0edb3ab2755f..50898013e8f9 100644 --- a/python/ray/serve/tests/unit/test_pow_2_request_router.py +++ b/python/ray/serve/tests/unit/test_pow_2_request_router.py @@ -2232,5 +2232,41 @@ def test_compute_backoff_s_does_not_overflow(): assert router._compute_backoff_s(2048) == router.max_backoff_s +@pytest.mark.asyncio +async def test_request_routing_timeout_no_replicas(pow_2_router): + """A request that isn't assigned a replica within the timeout fails.""" + s = pow_2_router + s.request_routing_timeout_s = 0.05 + loop = get_or_create_event_loop() + + task = loop.create_task(s._choose_replica_for_request(fake_pending_request())) + with pytest.raises(TimeoutError, match="Failed to route request to a replica"): + await asyncio.wait_for(task, timeout=10) + + assert s.num_pending_requests == 0 + assert len(s._pending_requests_to_route) == 0 + + +@pytest.mark.asyncio +async def test_request_routing_timeout_when_replicas_maxed(pow_2_router): + """A request that times out stops its routing task.""" + s = pow_2_router + s.request_routing_timeout_s = 0.05 + loop = get_or_create_event_loop() + + r1 = FakeRunningReplica("r1") + r1.set_queue_len_response(DEFAULT_MAX_ONGOING_REQUESTS) + s.update_replicas([r1]) + + task = loop.create_task(s._choose_replica_for_request(fake_pending_request())) + with pytest.raises(TimeoutError, match="Failed to route request to a replica"): + await asyncio.wait_for(task, timeout=10) + + await async_wait_for_condition( + lambda: s.curr_num_routing_tasks == 0, retry_interval_ms=1 + ) + assert s.num_pending_requests == 0 + + if __name__ == "__main__": sys.exit(pytest.main(["-v", "-s", __file__])) diff --git a/python/ray/serve/tests/unit/test_router.py b/python/ray/serve/tests/unit/test_router.py index 3528d88715a6..4f639a58fda7 100644 --- a/python/ray/serve/tests/unit/test_router.py +++ b/python/ray/serve/tests/unit/test_router.py @@ -3172,6 +3172,7 @@ async def test_update_deployment_config_sets_backoff_params(self): initial_backoff_s=custom_initial_backoff, backoff_multiplier=custom_multiplier, max_backoff_s=custom_max_backoff, + request_routing_timeout_s=1.5, ) ) @@ -3187,6 +3188,7 @@ async def test_update_deployment_config_sets_backoff_params(self): assert fake_request_router.initial_backoff_s == custom_initial_backoff assert fake_request_router.backoff_multiplier == custom_multiplier assert fake_request_router.max_backoff_s == custom_max_backoff + assert fake_request_router.request_routing_timeout_s == 1.5 class TestOnRequestCompleted: diff --git a/src/ray/protobuf/serve.proto b/src/ray/protobuf/serve.proto index 5e45c0efb548..47744eea5884 100644 --- a/src/ray/protobuf/serve.proto +++ b/src/ray/protobuf/serve.proto @@ -135,6 +135,9 @@ message RequestRouterConfig { // Maximum backoff time (in seconds) between retries. double max_backoff_s = 8; + + // Maximum time (in seconds) a request waits to be assigned a replica. + optional double request_routing_timeout_s = 9; } //[End] ROUTING CONFIG From 373bf496b369d08ec5137d8f2884f9303eb626f7 Mon Sep 17 00:00:00 2001 From: codechrl Date: Mon, 28 Sep 2026 20:10:06 +0700 Subject: [PATCH 2/2] [serve] Drop timed-out requests from the routing queues directly The lazy cleanup only pops a done prefix, so a request that timed out behind a still-pending one stayed queued. Remove it by identity on timeout, cancel its future explicitly, and leave the optional request_routing_timeout_s proto field unset when it is None, matching retry_after_s. Signed-off-by: codechrl --- python/ray/serve/_private/config.py | 3 ++ .../_private/request_router/request_router.py | 13 +++++--- .../tests/unit/test_pow_2_request_router.py | 30 +++++++++++++++++++ 3 files changed, 42 insertions(+), 4 deletions(-) diff --git a/python/ray/serve/_private/config.py b/python/ray/serve/_private/config.py index f3aebc540cce..fa16f7517d5c 100644 --- a/python/ray/serve/_private/config.py +++ b/python/ray/serve/_private/config.py @@ -370,6 +370,9 @@ def to_proto(self): **data["autoscaling_config"] ) if data.get("request_router_config"): + if data["request_router_config"].get("request_routing_timeout_s") is None: + # Leave the `optional` proto field unset rather than passing None. + data["request_router_config"].pop("request_routing_timeout_s", None) router_kwargs = data["request_router_config"].get("request_router_kwargs") if router_kwargs is not None: if not router_kwargs: diff --git a/python/ray/serve/_private/request_router/request_router.py b/python/ray/serve/_private/request_router/request_router.py index 1889d65f16d2..451589ad6936 100644 --- a/python/ray/serve/_private/request_router/request_router.py +++ b/python/ray/serve/_private/request_router/request_router.py @@ -1421,16 +1421,21 @@ async def _choose_replica_for_request( raise e from None except asyncio.TimeoutError: + pending_request.future.cancel() self._cancel_routing_task_for_pending_request(pending_request) self._remove_pending_request_from_indices(pending_request) - # No routing task runs while the deployment has no replicas, so the - # lazy cleanup in `_fulfill_pending_requests` can't drop the request. + # Remove the expired request itself: the lazy cleanup only pops a done + # prefix, so it can't reach a request queued behind a pending one, and + # no routing task runs while the deployment has no replicas. Match by + # identity; dataclass `==` would compare the request args. for queue in ( self._pending_requests_to_fulfill, self._pending_requests_to_route, ): - while queue and queue[0].future.done(): - queue.popleft() + for i, queued in enumerate(queue): + if queued is pending_request: + del queue[i] + break raise TimeoutError( f"Failed to route request to a replica of {self._deployment_id} " diff --git a/python/ray/serve/tests/unit/test_pow_2_request_router.py b/python/ray/serve/tests/unit/test_pow_2_request_router.py index 50898013e8f9..7c45e575ee82 100644 --- a/python/ray/serve/tests/unit/test_pow_2_request_router.py +++ b/python/ray/serve/tests/unit/test_pow_2_request_router.py @@ -2247,6 +2247,36 @@ async def test_request_routing_timeout_no_replicas(pow_2_router): assert len(s._pending_requests_to_route) == 0 +@pytest.mark.asyncio +async def test_request_routing_timeout_behind_pending_request(pow_2_router): + """A timed-out request queued behind a still-pending one is dropped.""" + s = pow_2_router + loop = get_or_create_event_loop() + + # Head request enqueued before the timeout was set: it waits indefinitely. + head = fake_pending_request() + head_task = loop.create_task(s._choose_replica_for_request(head)) + await asyncio.sleep(0) + + s.request_routing_timeout_s = 0.05 + expired = fake_pending_request() + task = loop.create_task(s._choose_replica_for_request(expired)) + with pytest.raises(TimeoutError, match="Failed to route request to a replica"): + await asyncio.wait_for(task, timeout=10) + + assert expired.future.cancelled() + assert list(s._pending_requests_to_fulfill) == [head] + assert list(s._pending_requests_to_route) == [head] + + r1 = FakeRunningReplica("r1") + r1.set_queue_len_response(0) + s.update_replicas([r1]) + + assert (await head_task) == r1 + assert s.curr_num_routing_tasks == 0 + assert s.num_pending_requests == 0 + + @pytest.mark.asyncio async def test_request_routing_timeout_when_replicas_maxed(pow_2_router): """A request that times out stops its routing task."""