mirror of
https://github.com/wassname/ray.git
synced 2026-07-24 13:20:22 +08:00
Fix hang if actor object id is returned from a task that exits (#6885)
This commit is contained in:
@@ -135,6 +135,7 @@ scripts/nodes.txt
|
||||
**/.pytest_cache
|
||||
**/.cache
|
||||
.benchmarks
|
||||
python-driver-*
|
||||
|
||||
# Vscode
|
||||
.vscode/
|
||||
|
||||
@@ -16,6 +16,25 @@ from ray.test_utils import run_string_as_driver
|
||||
from ray.experimental.internal_kv import _internal_kv_get, _internal_kv_put
|
||||
|
||||
|
||||
def test_actor_exit_from_task(ray_start_regular):
|
||||
@ray.remote
|
||||
class Actor:
|
||||
def __init__(self):
|
||||
print("Actor created")
|
||||
|
||||
def f(self):
|
||||
return 0
|
||||
|
||||
@ray.remote
|
||||
def f():
|
||||
a = Actor.remote()
|
||||
x_id = a.f.remote()
|
||||
return [x_id]
|
||||
|
||||
x_id = ray.get(f.remote())[0]
|
||||
print(ray.get(x_id)) # This should not hang.
|
||||
|
||||
|
||||
def test_actor_init_error_propagated(ray_start_regular):
|
||||
@ray.remote
|
||||
class Actor:
|
||||
|
||||
@@ -288,13 +288,14 @@ void CoreWorker::SetCurrentTaskId(const TaskID &task_id) {
|
||||
absl::MutexLock lock(&mutex_);
|
||||
not_actor_task = actor_id_.IsNil();
|
||||
}
|
||||
// Clear all actor handles at the end of each non-actor task.
|
||||
if (not_actor_task && task_id.IsNil()) {
|
||||
absl::MutexLock lock(&actor_handles_mutex_);
|
||||
// Reset the seqnos so that for the next task it start off at 0.
|
||||
for (const auto &handle : actor_handles_) {
|
||||
RAY_CHECK_OK(gcs_client_->Actors().AsyncUnsubscribe(handle.first, nullptr));
|
||||
handle.second->Reset();
|
||||
}
|
||||
actor_handles_.clear();
|
||||
// TODO(ekl) we can't unsubscribe to actor notifications here due to
|
||||
// https://github.com/ray-project/ray/pull/6885
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user