diff --git a/src/ray/gcs/gcs_server/gcs_actor_manager.cc b/src/ray/gcs/gcs_server/gcs_actor_manager.cc index d2b15dd82..b75e4ef78 100644 --- a/src/ray/gcs/gcs_server/gcs_actor_manager.cc +++ b/src/ray/gcs/gcs_server/gcs_actor_manager.cc @@ -90,7 +90,9 @@ GcsActorManager::GcsActorManager(std::shared_ptr sch : gcs_actor_scheduler_(std::move(scheduler)), gcs_table_storage_(std::move(gcs_table_storage)), gcs_pub_sub_(std::move(gcs_pub_sub)), - worker_client_factory_(worker_client_factory) {} + worker_client_factory_(worker_client_factory) { + RAY_CHECK(worker_client_factory_); +} void GcsActorManager::HandleRegisterActor(const rpc::RegisterActorRequest &request, rpc::RegisterActorReply *reply, @@ -99,15 +101,18 @@ void GcsActorManager::HandleRegisterActor(const rpc::RegisterActorRequest &reque auto actor_id = ActorID::FromBinary(request.task_spec().actor_creation_task_spec().actor_id()); - RAY_LOG(INFO) << "Registering actor, actor id = " << actor_id; + RAY_LOG(INFO) << "Registering actor, job id = " << actor_id.JobId() + << ", actor id = " << actor_id; Status status = RegisterActor(request, [reply, send_reply_callback, actor_id](const std::shared_ptr &actor) { - RAY_LOG(INFO) << "Registered actor, actor id = " << actor_id; + RAY_LOG(INFO) << "Registered actor, job id = " << actor_id.JobId() + << ", actor id = " << actor_id; GCS_RPC_SEND_REPLY(send_reply_callback, reply, Status::OK()); }); if (!status.ok()) { - RAY_LOG(ERROR) << "Failed to register actor: " << status.ToString(); + RAY_LOG(ERROR) << "Failed to register actor: " << status.ToString() + << ", job id = " << actor_id.JobId() << ", actor id = " << actor_id; GCS_RPC_SEND_REPLY(send_reply_callback, reply, status); } } @@ -119,15 +124,17 @@ void GcsActorManager::HandleCreateActor(const rpc::CreateActorRequest &request, auto actor_id = ActorID::FromBinary(request.task_spec().actor_creation_task_spec().actor_id()); - RAY_LOG(INFO) << "Creating actor, actor id = " << actor_id; + RAY_LOG(INFO) << "Creating actor, job id = " << actor_id.JobId() + << ", actor id = " << actor_id; Status status = CreateActor(request, [reply, send_reply_callback, actor_id]( const std::shared_ptr &actor) { - RAY_LOG(INFO) << "Created actor, actor id = " << actor_id; + RAY_LOG(INFO) << "Finished creating actor, job id = " << actor_id.JobId() + << ", actor id = " << actor_id; GCS_RPC_SEND_REPLY(send_reply_callback, reply, Status::OK()); }); if (!status.ok()) { - RAY_LOG(ERROR) << "Failed to create actor, actor id = " << actor_id - << " status: " << status.ToString(); + RAY_LOG(ERROR) << "Failed to create actor, job id = " << actor_id.JobId() + << ", actor id = " << actor_id << ", status: " << status.ToString(); GCS_RPC_SEND_REPLY(send_reply_callback, reply, status); } } @@ -181,8 +188,7 @@ void GcsActorManager::HandleGetNamedActorInfo( const rpc::GetNamedActorInfoRequest &request, rpc::GetNamedActorInfoReply *reply, rpc::SendReplyCallback send_reply_callback) { const std::string &name = request.name(); - RAY_LOG(DEBUG) << "Getting actor info" - << ", name = " << name; + RAY_LOG(DEBUG) << "Getting actor info, name = " << name; auto on_done = [name, reply, send_reply_callback]( const Status &status, @@ -426,7 +432,7 @@ Status GcsActorManager::RegisterActor( auto worker_id = WorkerID::FromBinary(owner_address.worker_id()); RAY_CHECK(unresolved_actors_[node_id][worker_id].emplace(actor->GetActorID()).second); - if (!actor->IsDetached() && worker_client_factory_) { + if (!actor->IsDetached()) { // This actor is owned. Send a long polling request to the actor's // owner to determine when the actor should be removed. PollOwnerForActorOutOfScope(actor); @@ -483,20 +489,8 @@ Status GcsActorManager::CreateActor(const ray::rpc::CreateActorRequest &request, // Remove the actor from the unresolved actor map. auto actor = std::make_shared(request.task_spec()); actor->GetMutableActorTableData()->set_state(rpc::ActorTableData::PENDING_CREATION); - const auto &owner_address = actor->GetOwnerAddress(); - auto node_id = ClientID::FromBinary(owner_address.raylet_id()); - auto worker_id = WorkerID::FromBinary(owner_address.worker_id()); - auto it = unresolved_actors_.find(node_id); - RAY_CHECK(it != unresolved_actors_.end()); - auto worker_to_actors_it = it->second.find(worker_id); - RAY_CHECK(worker_to_actors_it != it->second.end()); - RAY_CHECK(worker_to_actors_it->second.erase(actor_id) != 0); - if (worker_to_actors_it->second.empty()) { - it->second.erase(worker_to_actors_it); - if (it->second.empty()) { - unresolved_actors_.erase(it); - } - } + RemoveUnresolvedActor(actor); + // Update the registered actor as its creation task specification may have changed due // to resolved dependencies. registered_actors_[actor_id] = actor; @@ -537,7 +531,10 @@ void GcsActorManager::PollOwnerForActorOutOfScope( wait_request, [this, owner_node_id, owner_id, actor_id]( Status status, const rpc::WaitForActorOutOfScopeReply &reply) { if (!status.ok()) { - RAY_LOG(INFO) << "Worker " << owner_id << " failed, destroying actor child"; + RAY_LOG(WARNING) << "Worker " << owner_id << " failed, destroying actor child."; + } else { + RAY_LOG(WARNING) << "Actor " << actor_id + << " is out of scope,, destroying actor child."; } auto node_it = owners_.find(owner_node_id); @@ -550,7 +547,7 @@ void GcsActorManager::PollOwnerForActorOutOfScope( } void GcsActorManager::DestroyActor(const ActorID &actor_id) { - RAY_LOG(DEBUG) << "Destroying actor " << actor_id; + RAY_LOG(INFO) << "Destroying actor, actor id = " << actor_id; actor_to_register_callbacks_.erase(actor_id); actor_to_create_callbacks_.erase(actor_id); auto it = registered_actors_.find(actor_id); @@ -561,21 +558,7 @@ void GcsActorManager::DestroyActor(const ActorID &actor_id) { // Clean up the client to the actor's owner, if necessary. if (!actor->IsDetached()) { - const auto &owner_node_id = actor->GetOwnerNodeID(); - const auto &owner_id = actor->GetOwnerID(); - RAY_LOG(DEBUG) << "Erasing actor " << actor_id << " owned by " << owner_id; - - auto &node = owners_[owner_node_id]; - auto worker_it = node.find(owner_id); - RAY_CHECK(worker_it != node.end()); - auto &owner = worker_it->second; - RAY_CHECK(owner.children_actor_ids.erase(actor_id)); - if (owner.children_actor_ids.empty()) { - node.erase(worker_it); - if (node.empty()) { - owners_.erase(owner_node_id); - } - } + RemoveActorFromOwner(actor); } // The actor is already dead, most likely due to process or node failure. @@ -586,21 +569,7 @@ void GcsActorManager::DestroyActor(const ActorID &actor_id) { if (actor->GetState() == rpc::ActorTableData::DEPENDENCIES_UNREADY) { // The actor creation task still has unresolved dependencies. Remove from the // unresolved actors map. - const auto &owner_address = actor->GetOwnerAddress(); - auto node_id = ClientID::FromBinary(owner_address.raylet_id()); - auto worker_id = WorkerID::FromBinary(owner_address.worker_id()); - auto iter = unresolved_actors_.find(node_id); - if (iter != unresolved_actors_.end()) { - auto it = iter->second.find(worker_id); - RAY_CHECK(it != iter->second.end()); - RAY_CHECK(it->second.erase(actor_id) != 0); - if (it->second.empty()) { - iter->second.erase(it); - if (iter->second.empty()) { - unresolved_actors_.erase(iter); - } - } - } + RemoveUnresolvedActor(actor); } else { // The actor is still alive or pending creation. Clean up all remaining // state. @@ -610,13 +579,7 @@ void GcsActorManager::DestroyActor(const ActorID &actor_id) { if (node_it != created_actors_.end() && node_it->second.count(worker_id)) { // The actor has already been created. Destroy the process by force-killing // it. - auto actor_client = worker_client_factory_(actor->GetAddress()); - rpc::KillActorRequest request; - request.set_intended_actor_id(actor_id.Binary()); - request.set_force_kill(true); - request.set_no_restart(true); - RAY_UNUSED(actor_client->KillActor(request, nullptr)); - + KillActor(actor); RAY_CHECK(node_it->second.erase(actor->GetWorkerID())); if (node_it->second.empty()) { created_actors_.erase(node_it); @@ -697,6 +660,13 @@ absl::flat_hash_set GcsActorManager::GetUnresolvedActorsByOwnerWorker( void GcsActorManager::OnWorkerDead(const ray::ClientID &node_id, const ray::WorkerID &worker_id, bool intentional_exit) { + if (intentional_exit) { + RAY_LOG(INFO) << "Worker " << worker_id << " on node " << node_id + << " intentional exit."; + } else { + RAY_LOG(WARNING) << "Worker " << worker_id << " on node " << node_id + << " failed and exited abnormally."; + } // Destroy all actors that are owned by this worker. const auto it = owners_.find(node_id); if (it != owners_.end() && it->second.count(worker_id)) { @@ -738,9 +708,6 @@ void GcsActorManager::OnWorkerDead(const ray::ClientID &node_id, // Otherwise, try to reconstruct the actor that was already created or in the creation // process. - RAY_LOG(WARNING) << "Worker " << worker_id << " on node " << node_id - << " failed, restarting actor " << actor_id - << ", intentional exit: " << intentional_exit; ReconstructActor(actor_id, /*need_reschedule=*/!intentional_exit); } @@ -871,7 +838,7 @@ void GcsActorManager::OnActorCreationFailed(std::shared_ptr actor) { void GcsActorManager::OnActorCreationSuccess(const std::shared_ptr &actor) { auto actor_id = actor->GetActorID(); - RAY_LOG(DEBUG) << "Actor created successfully, actor id = " << actor_id; + RAY_LOG(INFO) << "Actor created successfully, actor id = " << actor_id; // NOTE: If an actor is deleted immediately after the user creates the actor, reference // counter may return a reply to the request of WaitForActorOutOfScope to GCS server, // and GCS server will destroy the actor. The actor creation is asynchronous, it may be @@ -982,8 +949,8 @@ void GcsActorManager::LoadInitialData(const EmptyCallback &done) { // We should not reschedule actors in state of `ALIVE`. // We could not reschedule actors in state of `DEPENDENCIES_UNREADY` because the // dependencies of them may not have been resolved yet. - RAY_LOG(DEBUG) << "Rescheduling a non-alive actor, actor id = " - << actor->GetActorID() << ", state = " << actor->GetState(); + RAY_LOG(INFO) << "Rescheduling a non-alive actor, actor id = " + << actor->GetActorID() << ", state = " << actor->GetState(); gcs_actor_scheduler_->Reschedule(actor); } } @@ -1056,5 +1023,51 @@ const absl::flat_hash_map> return actor_to_register_callbacks_; } +void GcsActorManager::RemoveUnresolvedActor(const std::shared_ptr &actor) { + const auto &owner_address = actor->GetOwnerAddress(); + auto node_id = ClientID::FromBinary(owner_address.raylet_id()); + auto worker_id = WorkerID::FromBinary(owner_address.worker_id()); + auto iter = unresolved_actors_.find(node_id); + if (iter != unresolved_actors_.end()) { + auto it = iter->second.find(worker_id); + RAY_CHECK(it != iter->second.end()); + RAY_CHECK(it->second.erase(actor->GetActorID()) != 0); + if (it->second.empty()) { + iter->second.erase(it); + if (iter->second.empty()) { + unresolved_actors_.erase(iter); + } + } + } +} + +void GcsActorManager::RemoveActorFromOwner(const std::shared_ptr &actor) { + const auto &actor_id = actor->GetActorID(); + const auto &owner_id = actor->GetOwnerID(); + RAY_LOG(INFO) << "Erasing actor " << actor_id << " owned by " << owner_id; + + const auto &owner_node_id = actor->GetOwnerNodeID(); + auto &node = owners_[owner_node_id]; + auto worker_it = node.find(owner_id); + RAY_CHECK(worker_it != node.end()); + auto &owner = worker_it->second; + RAY_CHECK(owner.children_actor_ids.erase(actor_id)); + if (owner.children_actor_ids.empty()) { + node.erase(worker_it); + if (node.empty()) { + owners_.erase(owner_node_id); + } + } +} + +void GcsActorManager::KillActor(const std::shared_ptr &actor) { + auto actor_client = worker_client_factory_(actor->GetAddress()); + rpc::KillActorRequest request; + request.set_intended_actor_id(actor->GetActorID().Binary()); + request.set_force_kill(true); + request.set_no_restart(true); + RAY_UNUSED(actor_client->KillActor(request, nullptr)); +} + } // namespace gcs } // namespace ray diff --git a/src/ray/gcs/gcs_server/gcs_actor_manager.h b/src/ray/gcs/gcs_server/gcs_actor_manager.h index 4b9d21825..6ce4f2a2a 100644 --- a/src/ray/gcs/gcs_server/gcs_actor_manager.h +++ b/src/ray/gcs/gcs_server/gcs_actor_manager.h @@ -337,6 +337,21 @@ class GcsActorManager : public rpc::ActorInfoHandler { /// again. void ReconstructActor(const ActorID &actor_id, bool need_reschedule = true); + /// Remove the specified actor from `unresolved_actors_`. + /// + /// \param actor The actor to be removed. + void RemoveUnresolvedActor(const std::shared_ptr &actor); + + /// Remove the specified actor from owner. + /// + /// \param actor The actor to be removed. + void RemoveActorFromOwner(const std::shared_ptr &actor); + + /// Kill the specified actor. + /// + /// \param actor The actor to be killed. + void KillActor(const std::shared_ptr &actor); + /// Callbacks of pending `RegisterActor` requests. /// Maps actor ID to actor registration callbacks, which is used to filter duplicated /// messages from a driver/worker caused by some network problems. diff --git a/src/ray/gcs/redis_context.cc b/src/ray/gcs/redis_context.cc index ad8f2f542..bb1230c06 100644 --- a/src/ray/gcs/redis_context.cc +++ b/src/ray/gcs/redis_context.cc @@ -137,8 +137,6 @@ void CallbackReply::ParseAsStringArrayOrScanArray(redisReply *redis_reply) { RAY_CHECK(REDIS_REPLY_STRING == entry->type) << "Unexcepted type: " << entry->type; string_array_reply_.push_back(std::string(entry->str, entry->len)); - RAY_LOG(DEBUG) << "Element index is " << i << ", element length is " - << entry->len; } return; }