Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
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
3 changes: 2 additions & 1 deletion doc/source/serve/configure-serve-deployment.md
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,8 @@ You can also refer to the [API reference](../serve/api/doc/ray.serve.deployment_
- `ray_actor_options` - Options to pass to the Ray Actor decorator, such as resource requirements. Valid options are: `accelerator_type`, `memory`, `num_cpus`, `num_gpus`, `object_store_memory`, `resources`, and `runtime_env` For more details - [Resource management in Serve](serve-cpus-gpus)
- `max_ongoing_requests` - Maximum number of queries that are sent to a replica of this deployment without receiving a response. Defaults to 5 (note the default changed from 100 to 5 in Ray 2.32.0). This may be an important parameter to configure for [performance tuning](serve-perf-tuning).
- `autoscaling_config` - Parameters to configure autoscaling behavior. If this is set, you can't set `num_replicas` to a number. For more details on configurable parameters for autoscaling, see [Ray Serve Autoscaling](serve-autoscaling).
- `max_queued_requests` - Maximum number of requests to this deployment that will be queued at each caller (proxy or DeploymentHandle). Once this limit is reached, subsequent requests will raise a BackPressureError (for handles) or return an HTTP 503 status code (for HTTP requests). Defaults to -1 (no limit).
- `max_queued_requests` - Maximum number of requests to this deployment that will be queued at each caller (proxy or DeploymentHandle). Once this limit is reached, subsequent requests will raise a BackPressureError (for handles) or return an HTTP 503 status code by default (configurable via `backpressure_config`) for HTTP requests. Defaults to -1 (no limit).
- `request_router_config` - Advanced configuration for the request router. Accepts a `RequestRouterConfig` object with options including `max_request_retries` (maximum number of times the router retries a request when replicas reject it; defaults to -1 for unlimited retries; set to a non-negative integer to bound retries and prevent OOM under sustained load).
- `user_config` - Config to pass to the reconfigure method of the deployment. This can be updated dynamically without restarting the replicas of the deployment. The user_config must be fully JSON-serializable. For more details, see [Serve User Config](serve-user-config).
- `health_check_period_s` - Duration between health check calls for the replica. Defaults to 10s. The health check is by default a no-op Actor call to the replica, but you can define your own health check using the "check_health" method in your deployment that raises an exception when unhealthy.
- `health_check_timeout_s` - Duration in seconds, that replicas wait for a health check method to return before considering it as failed. Defaults to 30s.
Expand Down
6 changes: 6 additions & 0 deletions python/ray/serve/_private/constants.py
Original file line number Diff line number Diff line change
Expand Up @@ -536,6 +536,12 @@
"RAY_SERVE_ROUTER_RETRY_MAX_BACKOFF_S", 0.5
)

# Maximum number of times the router retries routing a request after a replica
# rejects it. -1 means unlimited retries (default, for backwards compatibility).
RAY_SERVE_ROUTER_MAX_REQUEST_RETRIES = get_env_int(
"RAY_SERVE_ROUTER_MAX_REQUEST_RETRIES", -1
)

# The default autoscaling policy to use if none is specified.
DEFAULT_AUTOSCALING_POLICY_NAME = (
"ray.serve.autoscaling_policy:default_autoscaling_policy"
Expand Down
27 changes: 27 additions & 0 deletions python/ray/serve/_private/router.py
Original file line number Diff line number Diff line change
Expand Up @@ -673,6 +673,8 @@ def __init__(
self._initial_backoff_s: Optional[float] = None
self._backoff_multiplier: Optional[float] = None
self._max_backoff_s: Optional[float] = None
self._max_request_retries: int = -1
Comment thread
HrushiYadav marked this conversation as resolved.
self._max_queued_requests: int = -1

# Initializing `self._metrics_manager` before `self.long_poll_client` is
# necessary to avoid race condition where `self.update_deployment_config()`
Expand Down Expand Up @@ -848,6 +850,10 @@ 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._max_request_retries = (
deployment_config.request_router_config.max_request_retries
)
Comment thread
HrushiYadav marked this conversation as resolved.
self._max_queued_requests = deployment_config.max_queued_requests

if self._request_router:
self._request_router.update_backoff_params(
Expand Down Expand Up @@ -1117,6 +1123,15 @@ async def _route_and_send_request_once(

return None

def _check_retry_limit(self, num_retries: int) -> None:
"""Raise BackPressureError if the retry limit has been exceeded."""
max_retries = self._max_request_retries
if max_retries >= 0 and num_retries > max_retries:
raise BackPressureError(
num_queued_requests=self._metrics_manager.num_queued_requests,
max_queued_requests=self._max_queued_requests,
)
Comment thread
HrushiYadav marked this conversation as resolved.
Comment thread
cursor[bot] marked this conversation as resolved.

async def route_and_send_request(
self,
pr: PendingRequest,
Expand All @@ -1126,11 +1141,13 @@ async def route_and_send_request(

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.
If max_request_retries is configured (>= 0), retries are bounded.
"""
# Wait for the router to be initialized before sending the request.
await self._request_router_initialized.wait()

is_retry = False
num_retries = 0
while True:
result = await self._route_and_send_request_once(
pr,
Expand All @@ -1145,6 +1162,8 @@ async def route_and_send_request(
# TODO(edoakes): this retry procedure is not perfect because it'll reset the
# process of choosing candidates replicas (i.e., for locality-awareness).
is_retry = True
num_retries += 1
self._check_retry_limit(num_retries)

@tracing_decorator_factory(
trace_name="route_to_replica",
Expand Down Expand Up @@ -1473,8 +1492,10 @@ async def _pick_and_reserve_replica(
capacity rejection, replica actor death, or transient
unavailability. Updates the queue-length cache from the replica's
reported count.
If max_request_retries is configured (>= 0), retries are bounded.
"""
is_retry = False
num_retries = 0
while True:
num_curr_replicas = len(self.request_router.curr_replicas)
with self._metrics_manager.wrap_queued_request(
Expand All @@ -1490,11 +1511,15 @@ async def _pick_and_reserve_replica(
replica.replica_id, replica.actor_id, e
)
is_retry = True
num_retries += 1
self._check_retry_limit(num_retries)
continue
except ActorUnavailableError:
self.request_router.on_replica_actor_unavailable(replica.replica_id)
logger.warning(f"{replica.replica_id} is temporarily unavailable.")
is_retry = True
num_retries += 1
self._check_retry_limit(num_retries)
continue

self.request_router.on_new_queue_len_info(
Expand All @@ -1504,6 +1529,8 @@ async def _pick_and_reserve_replica(
return replica, slot_token

is_retry = True
num_retries += 1
self._check_retry_limit(num_retries)

def _register_completion_callback(
self, result: ReplicaResult, replica: RunningReplica, pr: PendingRequest
Expand Down
23 changes: 23 additions & 0 deletions python/ray/serve/config.py
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,7 @@
DEFAULT_REQUEST_ROUTING_STATS_TIMEOUT_S,
DEFAULT_TARGET_ONGOING_REQUESTS,
DEFAULT_UVICORN_KEEP_ALIVE_TIMEOUT_S,
RAY_SERVE_ROUTER_MAX_REQUEST_RETRIES,
RAY_SERVE_ROUTER_RETRY_BACKOFF_MULTIPLIER,
RAY_SERVE_ROUTER_RETRY_INITIAL_BACKOFF_S,
RAY_SERVE_ROUTER_RETRY_MAX_BACKOFF_S,
Expand Down Expand Up @@ -279,6 +280,26 @@ class RequestRouterConfig(BaseModel):
),
)

max_request_retries: int = Field(
default=RAY_SERVE_ROUTER_MAX_REQUEST_RETRIES,
description=(
"Maximum number of times the router retries routing a request "
"after a replica rejects it. -1 means unlimited retries "
"(the default, for backwards compatibility). When the limit is "
"exceeded, the request is dropped with a BackPressureError (503)."
),
)
Comment thread
cursor[bot] marked this conversation as resolved.

@field_validator("max_request_retries")
@classmethod
def validate_max_request_retries(cls, v):
if v < -1:
raise ValueError(
"max_request_retries must be -1 (unlimited) "
"or a non-negative integer."
)
return v

@field_validator("request_router_kwargs")
@classmethod
def request_router_kwargs_json_serializable(cls, v):
Expand Down Expand Up @@ -313,6 +334,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.max_request_retries == other.max_request_retries
)

def __hash__(self):
Expand All @@ -334,6 +356,7 @@ def __hash__(self):
self.initial_backoff_s,
self.backoff_multiplier,
self.max_backoff_s,
self.max_request_retries,
)
)

Expand Down
1 change: 1 addition & 0 deletions python/ray/serve/tests/test_controller.py
Original file line number Diff line number Diff line change
Expand Up @@ -210,6 +210,7 @@ def autoscaling_app():
"initial_backoff_s": 0.025,
"backoff_multiplier": 2.0,
"max_backoff_s": 0.5,
"max_request_retries": -1,
},
"rolling_update_percentage": 0.2,
},
Expand Down
1 change: 1 addition & 0 deletions python/ray/serve/tests/test_direct_ingress.py
Original file line number Diff line number Diff line change
Expand Up @@ -2570,6 +2570,7 @@ def autoscaling_app():
"initial_backoff_s": 0.025,
"backoff_multiplier": 2.0,
"max_backoff_s": 0.5,
"max_request_retries": -1,
},
"rolling_update_percentage": 0.2,
},
Expand Down
29 changes: 29 additions & 0 deletions python/ray/serve/tests/unit/test_config.py
Original file line number Diff line number Diff line change
Expand Up @@ -1666,6 +1666,35 @@ def test_optional_field(self):
assert "initial_replicas" not in result


def test_max_request_retries_config_validation():
"""max_request_retries accepts -1 and non-negative ints, rejects < -1."""
RequestRouterConfig(max_request_retries=-1)
RequestRouterConfig(max_request_retries=0)
RequestRouterConfig(max_request_retries=10)

with pytest.raises(ValueError):
RequestRouterConfig(max_request_retries=-2)


@pytest.mark.parametrize("max_retries", [-1, 0, 5])
def test_max_request_retries_proto_roundtrip(max_retries):
"""Ensure max_request_retries survives to_proto/from_proto."""
config = DeploymentConfig(
request_router_config=RequestRouterConfig(max_request_retries=max_retries)
)
roundtripped = DeploymentConfig.from_proto_bytes(config.to_proto_bytes())
assert roundtripped.request_router_config.max_request_retries == max_retries


def test_max_request_retries_proto_absent_field():
"""Absent max_request_retries (older controller) falls back to -1 default."""
proto = DeploymentConfig().to_proto()
proto.request_router_config.ClearField("max_request_retries")
assert not proto.request_router_config.HasField("max_request_retries")
deserialized = DeploymentConfig.from_proto(proto)
assert deserialized.request_router_config.max_request_retries == -1


if __name__ == "__main__":
import sys

Expand Down
35 changes: 35 additions & 0 deletions python/ray/serve/tests/unit/test_router.py
Original file line number Diff line number Diff line change
Expand Up @@ -1857,6 +1857,41 @@ async def test_choose_replica_retries_when_reservation_rejected(

assert fake_request_router.replica_queue_len_cache.get(r2_id) == 0

@pytest.mark.parametrize(
"setup_router",
[{"enable_queue_len_cache": True}],
indirect=True,
)
async def test_choose_replica_bounds_retries_on_repeated_rejection(
self, setup_router: Tuple[AsyncioRouter, FakeRequestRouter]
):
"""When a replica keeps rejecting reservations, the router should raise
BackPressureError after max_request_retries instead of retrying forever."""
router, fake_request_router = setup_router
router._max_request_retries = 2

r1_id = ReplicaID(
unique_id="test-replica-1", deployment_id=DeploymentID(name="test")
)
r1 = FakeReplica(r1_id)
r1._reject_reservation = True
# Both the initial and retry picks return the same always-rejecting replica,
# so no reservation ever succeeds and the retry counter drives the outcome.
fake_request_router.set_replica_to_return(r1)
fake_request_router.set_replica_to_return_on_retry(r1)

request_metadata = RequestMetadata(
request_id="test-request-1",
internal_request_id="test-internal-request-1",
)

async def _enter():
async with router.choose_replica(request_metadata):
pass

with pytest.raises(BackPressureError):
await asyncio.wait_for(_enter(), timeout=5)

@pytest.mark.parametrize(
"setup_router",
[{"enable_queue_len_cache": True}],
Expand Down
6 changes: 6 additions & 0 deletions src/ray/protobuf/serve.proto
Original file line number Diff line number Diff line change
Expand Up @@ -135,6 +135,12 @@ message RequestRouterConfig {

// Maximum backoff time (in seconds) between retries.
double max_backoff_s = 8;

// Maximum number of times the router retries routing a request after a
// replica rejects it. -1 means unlimited retries (the default).
// Declared optional so that unset (e.g. from an older controller) is
// distinguishable from 0 (which means zero retries).
optional int32 max_request_retries = 9;
Comment thread
HrushiYadav marked this conversation as resolved.
}
//[End] ROUTING CONFIG

Expand Down
Loading