diff --git a/src/ray/gcs/gcs_server/gcs_server.cc b/src/ray/gcs/gcs_server/gcs_server.cc index 50266e45b..51243052b 100644 --- a/src/ray/gcs/gcs_server/gcs_server.cc +++ b/src/ray/gcs/gcs_server/gcs_server.cc @@ -183,15 +183,17 @@ void GcsServer::InitGcsActorManager() { // the GCS. gcs_actor_manager_->OnNodeDead(ClientID::FromBinary(node->node_id())); }); - RAY_CHECK_OK(redis_gcs_client_->Workers().AsyncSubscribeToWorkerFailures( - [this](const WorkerID &id, const rpc::WorkerFailureData &worker_failure_data) { - auto &worker_address = worker_failure_data.worker_address(); - WorkerID worker_id = WorkerID::FromBinary(worker_address.worker_id()); - ClientID node_id = ClientID::FromBinary(worker_address.raylet_id()); - gcs_actor_manager_->OnWorkerDead(node_id, worker_id, - worker_failure_data.intentional_disconnect()); - }, - /*done_callback=*/nullptr)); + + auto on_subscribe = [this](const std::string &id, const std::string &data) { + rpc::WorkerFailureData worker_failure_data; + worker_failure_data.ParseFromString(data); + auto &worker_address = worker_failure_data.worker_address(); + WorkerID worker_id = WorkerID::FromBinary(id); + ClientID node_id = ClientID::FromBinary(worker_address.raylet_id()); + gcs_actor_manager_->OnWorkerDead(node_id, worker_id, + worker_failure_data.intentional_disconnect()); + }; + RAY_CHECK_OK(gcs_pub_sub_->SubscribeAll(WORKER_FAILURE_CHANNEL, on_subscribe, nullptr)); } std::unique_ptr GcsServer::InitJobInfoHandler() {