Update named actor API (#8559)

This commit is contained in:
Edward Oakes
2020-05-24 20:08:03 -05:00
committed by GitHub
parent 92c2e41dfd
commit 860eb6f13a
10 changed files with 109 additions and 140 deletions
+1 -2
View File
@@ -107,7 +107,7 @@ def init(cluster_name=None,
global master_actor
master_actor_name = format_actor_name(SERVE_MASTER_NAME, cluster_name)
try:
master_actor = ray.util.get_actor(master_actor_name)
master_actor = ray.get_actor(master_actor_name)
return
except ValueError:
pass
@@ -124,7 +124,6 @@ def init(cluster_name=None,
# in the future.
http_node_id = ray.state.current_node_id()
master_actor = ServeMaster.options(
detached=True,
name=master_actor_name,
max_restarts=-1,
).remote(cluster_name, start_server, http_node_id, http_host, http_port,
+7 -13
View File
@@ -127,11 +127,10 @@ class ServeMaster:
"""
router_name = format_actor_name(SERVE_ROUTER_NAME, self.cluster_name)
try:
self.router = ray.util.get_actor(router_name)
self.router = ray.get_actor(router_name)
except ValueError:
logger.info("Starting router with name '{}'".format(router_name))
self.router = async_retryable(ray.remote(Router)).options(
detached=True,
name=router_name,
max_concurrency=ASYNC_CONCURRENCY,
max_restarts=-1,
@@ -148,13 +147,12 @@ class ServeMaster:
"""
proxy_name = format_actor_name(SERVE_PROXY_NAME, self.cluster_name)
try:
self.http_proxy = ray.util.get_actor(proxy_name)
self.http_proxy = ray.get_actor(proxy_name)
except ValueError:
logger.info(
"Starting HTTP proxy with name '{}' on node '{}'".format(
proxy_name, node_id))
self.http_proxy = async_retryable(HTTPProxyActor).options(
detached=True,
name=proxy_name,
max_concurrency=ASYNC_CONCURRENCY,
max_restarts=-1,
@@ -180,12 +178,11 @@ class ServeMaster:
metric_sink_name = format_actor_name(SERVE_METRIC_SINK_NAME,
self.cluster_name)
try:
self.metric_exporter = ray.util.get_actor(metric_sink_name)
self.metric_exporter = ray.get_actor(metric_sink_name)
except ValueError:
logger.info("Starting metric exporter with name '{}'".format(
metric_sink_name))
self.metric_exporter = MetricExporterActor.options(
detached=True,
name=metric_sink_name).remote(metric_exporter_class)
def get_metric_exporter(self):
@@ -246,7 +243,7 @@ class ServeMaster:
for replica_tag in replica_tags:
replica_name = format_actor_name(replica_tag,
self.cluster_name)
self.workers[backend_tag][replica_tag] = ray.util.get_actor(
self.workers[backend_tag][replica_tag] = ray.get_actor(
replica_name)
# Push configuration state to the router.
@@ -311,7 +308,6 @@ class ServeMaster:
replica_name = format_actor_name(replica_tag, self.cluster_name)
worker_handle = async_retryable(ray.remote(backend_worker)).options(
detached=True,
name=replica_name,
max_restarts=-1,
**replica_config.ray_actor_options).remote(
@@ -328,7 +324,7 @@ class ServeMaster:
# failed after creating them but before writing a
# checkpoint.
try:
worker_handle = ray.util.get_actor(replica_tag)
worker_handle = ray.get_actor(replica_tag)
except ValueError:
worker_handle = await self._start_backend_worker(
backend_tag, replica_tag)
@@ -371,7 +367,7 @@ class ServeMaster:
# NOTE(edoakes): the replicas may already be stopped if we
# failed after stopping them but before writing a checkpoint.
try:
replica = ray.util.get_actor(replica_tag)
replica = ray.get_actor(replica_tag)
except ValueError:
continue
@@ -384,9 +380,7 @@ class ServeMaster:
# use replica.__ray_terminate__, we may send it while the
# replica is being restarted and there's no way to tell if it
# successfully killed the worker or not.
worker = ray.worker.global_worker
# Kill the actor with no_restart=True.
worker.core_worker.kill_actor(replica._ray_actor_id, True)
ray.kill(replica, no_restart=True)
self.replicas_to_stop.clear()
+9 -9
View File
@@ -35,7 +35,7 @@ def test_master_failure(serve_instance):
response = request_with_retries("/master_failure", timeout=30)
assert response.text == "hello1"
ray.kill(serve.api._get_master_actor())
ray.kill(serve.api._get_master_actor(), no_restart=False)
for _ in range(10):
response = request_with_retries("/master_failure", timeout=30)
@@ -44,7 +44,7 @@ def test_master_failure(serve_instance):
def function():
return "hello2"
ray.kill(serve.api._get_master_actor())
ray.kill(serve.api._get_master_actor(), no_restart=False)
serve.create_backend("master_failure:v2", function)
serve.set_traffic("master_failure", {"master_failure:v2": 1.0})
@@ -56,11 +56,11 @@ def test_master_failure(serve_instance):
def function():
return "hello3"
ray.kill(serve.api._get_master_actor())
ray.kill(serve.api._get_master_actor(), no_restart=False)
serve.create_endpoint("master_failure_2", "/master_failure_2")
ray.kill(serve.api._get_master_actor())
ray.kill(serve.api._get_master_actor(), no_restart=False)
serve.create_backend("master_failure_2", function)
ray.kill(serve.api._get_master_actor())
ray.kill(serve.api._get_master_actor(), no_restart=False)
serve.set_traffic("master_failure_2", {"master_failure_2": 1.0})
for _ in range(10):
@@ -73,7 +73,7 @@ def test_master_failure(serve_instance):
def _kill_http_proxy():
[http_proxy] = ray.get(
serve.api._get_master_actor().get_http_proxy.remote())
ray.kill(http_proxy)
ray.kill(http_proxy, no_restart=False)
def test_http_proxy_failure(serve_instance):
@@ -107,7 +107,7 @@ def test_http_proxy_failure(serve_instance):
def _kill_router():
[router] = ray.get(serve.api._get_master_actor().get_router.remote())
ray.kill(router)
ray.kill(router, no_restart=False)
def test_router_failure(serve_instance):
@@ -169,7 +169,7 @@ def test_worker_restart(serve_instance):
# Kill the worker.
handles = _get_worker_handles("worker_failure:v1")
assert len(handles) == 1
ray.kill(handles[0])
ray.kill(handles[0], no_restart=False)
# Wait until the worker is killed and a one is started.
start = time.time()
@@ -227,7 +227,7 @@ def test_worker_replica_failure(serve_instance):
# Kill one of the replicas.
handles = _get_worker_handles("replica_failure")
assert len(handles) == 2
ray.kill(handles[0])
ray.kill(handles[0], no_restart=False)
# Check that the other replica still serves requests.
for _ in range(10):