From 3669c0282131f0c9fe7f1f4fc43072d2c082284c Mon Sep 17 00:00:00 2001 From: fangfengbin <869218239a@zju.edu.cn> Date: Thu, 7 Jan 2021 15:44:33 +0800 Subject: [PATCH] [GCS]Add gcs actor schedule strategy (#13156) --- .../gcs_server/gcs_actor_schedule_strategy.cc | 58 +++++++++++++++ .../gcs_server/gcs_actor_schedule_strategy.h | 71 +++++++++++++++++++ src/ray/gcs/gcs_server/gcs_actor_scheduler.cc | 14 +--- src/ray/gcs/gcs_server/gcs_actor_scheduler.h | 5 ++ src/ray/gcs/gcs_server/gcs_server.cc | 4 +- src/ray/gcs/gcs_server/gcs_server.h | 32 ++++----- .../test/gcs_actor_scheduler_test.cc | 5 +- .../gcs_server/test/gcs_server_test_util.h | 1 + 8 files changed, 161 insertions(+), 29 deletions(-) create mode 100644 src/ray/gcs/gcs_server/gcs_actor_schedule_strategy.cc create mode 100644 src/ray/gcs/gcs_server/gcs_actor_schedule_strategy.h diff --git a/src/ray/gcs/gcs_server/gcs_actor_schedule_strategy.cc b/src/ray/gcs/gcs_server/gcs_actor_schedule_strategy.cc new file mode 100644 index 000000000..fa737ea05 --- /dev/null +++ b/src/ray/gcs/gcs_server/gcs_actor_schedule_strategy.cc @@ -0,0 +1,58 @@ +// Copyright 2017 The Ray Authors. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +#include "ray/gcs/gcs_server/gcs_actor_schedule_strategy.h" +#include "ray/gcs/gcs_server/gcs_actor_manager.h" + +namespace ray { +namespace gcs { + +std::shared_ptr GcsRandomActorScheduleStrategy::Schedule( + std::shared_ptr actor) { + // Select a node to lease worker for the actor. + std::shared_ptr node; + + // If an actor has resource requirements, we will try to schedule it on the same node as + // the owner if possible. + const auto &task_spec = actor->GetCreationTaskSpecification(); + if (!task_spec.GetRequiredResources().IsEmpty()) { + auto maybe_node = gcs_node_manager_->GetAliveNode(actor->GetOwnerNodeID()); + node = maybe_node.has_value() ? maybe_node.value() : SelectNodeRandomly(); + } else { + node = SelectNodeRandomly(); + } + + return node; +} + +std::shared_ptr GcsRandomActorScheduleStrategy::SelectNodeRandomly() + const { + auto &alive_nodes = gcs_node_manager_->GetAllAliveNodes(); + if (alive_nodes.empty()) { + return nullptr; + } + + static std::mt19937_64 gen_( + std::chrono::high_resolution_clock::now().time_since_epoch().count()); + std::uniform_int_distribution distribution(0, alive_nodes.size() - 1); + int key_index = distribution(gen_); + int index = 0; + auto iter = alive_nodes.begin(); + for (; index != key_index && iter != alive_nodes.end(); ++index, ++iter) + ; + return iter->second; +} + +} // namespace gcs +} // namespace ray diff --git a/src/ray/gcs/gcs_server/gcs_actor_schedule_strategy.h b/src/ray/gcs/gcs_server/gcs_actor_schedule_strategy.h new file mode 100644 index 000000000..5b1a627af --- /dev/null +++ b/src/ray/gcs/gcs_server/gcs_actor_schedule_strategy.h @@ -0,0 +1,71 @@ +// Copyright 2017 The Ray Authors. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +#pragma once + +#include "ray/common/id.h" +#include "ray/gcs/gcs_server/gcs_node_manager.h" + +namespace ray { +namespace gcs { + +class GcsActor; + +/// \class GcsActorScheduleStrategyInterface +/// +/// Used for different kinds of actor scheduling strategy. +class GcsActorScheduleStrategyInterface { + public: + virtual ~GcsActorScheduleStrategyInterface() = default; + + /// Select a node to schedule the actor. + /// + /// \param actor The actor to be scheduled. + /// \return The selected node. If the scheduling fails, nullptr is returned. + virtual std::shared_ptr Schedule(std::shared_ptr actor) = 0; +}; + +/// \class GcsRandomActorScheduleStrategy +/// +/// This strategy will select node randomly from the node pool to schedule the actor. +class GcsRandomActorScheduleStrategy : public GcsActorScheduleStrategyInterface { + public: + /// Create a GcsRandomActorScheduleStrategy + /// + /// \param gcs_node_manager Node management of the cluster, which provides interfaces + /// to access the node information. + explicit GcsRandomActorScheduleStrategy( + std::shared_ptr gcs_node_manager) + : gcs_node_manager_(std::move(gcs_node_manager)) {} + + virtual ~GcsRandomActorScheduleStrategy() = default; + + /// Select a node to schedule the actor. + /// + /// \param actor The actor to be scheduled. + /// \return The selected node. If the scheduling fails, nullptr is returned. + std::shared_ptr Schedule(std::shared_ptr actor) override; + + private: + /// Select a node from alive nodes randomly. + /// + /// \return The selected node. If the scheduling fails, `nullptr` is returned. + std::shared_ptr SelectNodeRandomly() const; + + /// The node manager. + std::shared_ptr gcs_node_manager_; +}; + +} // namespace gcs +} // namespace ray diff --git a/src/ray/gcs/gcs_server/gcs_actor_scheduler.cc b/src/ray/gcs/gcs_server/gcs_actor_scheduler.cc index bbdd58856..9c81c8c0e 100644 --- a/src/ray/gcs/gcs_server/gcs_actor_scheduler.cc +++ b/src/ray/gcs/gcs_server/gcs_actor_scheduler.cc @@ -29,6 +29,7 @@ GcsActorScheduler::GcsActorScheduler( std::function)> schedule_failure_handler, std::function)> schedule_success_handler, std::shared_ptr raylet_client_pool, + std::shared_ptr actor_schedule_strategy, rpc::ClientFactoryFn client_factory) : io_context_(io_context), gcs_actor_table_(gcs_actor_table), @@ -38,6 +39,7 @@ GcsActorScheduler::GcsActorScheduler( schedule_success_handler_(std::move(schedule_success_handler)), report_worker_backlog_(RayConfig::instance().report_worker_backlog()), raylet_client_pool_(raylet_client_pool), + actor_schedule_strategy_(actor_schedule_strategy), core_worker_clients_(client_factory) { RAY_CHECK(schedule_failure_handler_ != nullptr && schedule_success_handler_ != nullptr); } @@ -46,17 +48,7 @@ void GcsActorScheduler::Schedule(std::shared_ptr actor) { RAY_CHECK(actor->GetNodeID().IsNil() && actor->GetWorkerID().IsNil()); // Select a node to lease worker for the actor. - std::shared_ptr node; - - // If an actor has resource requirements, we will try to schedule it on the same node as - // the owner if possible. - const auto &task_spec = actor->GetCreationTaskSpecification(); - if (!task_spec.GetRequiredResources().IsEmpty()) { - auto maybe_node = gcs_node_manager_.GetAliveNode(actor->GetOwnerNodeID()); - node = maybe_node.has_value() ? maybe_node.value() : SelectNodeRandomly(); - } else { - node = SelectNodeRandomly(); - } + const auto &node = actor_schedule_strategy_->Schedule(actor); if (node == nullptr) { // There are no available nodes to schedule the actor, so just trigger the failed diff --git a/src/ray/gcs/gcs_server/gcs_actor_scheduler.h b/src/ray/gcs/gcs_server/gcs_actor_scheduler.h index b59c4b2d4..71dd35108 100644 --- a/src/ray/gcs/gcs_server/gcs_actor_scheduler.h +++ b/src/ray/gcs/gcs_server/gcs_actor_scheduler.h @@ -22,6 +22,7 @@ #include "ray/common/task/task_execution_spec.h" #include "ray/common/task/task_spec.h" #include "ray/gcs/accessor.h" +#include "ray/gcs/gcs_server/gcs_actor_schedule_strategy.h" #include "ray/gcs/gcs_server/gcs_node_manager.h" #include "ray/gcs/gcs_server/gcs_table_storage.h" #include "ray/raylet_client/raylet_client.h" @@ -90,6 +91,7 @@ class GcsActorScheduler : public GcsActorSchedulerInterface { /// \param schedule_success_handler Invoked when actors are created on the worker /// successfully. /// \param raylet_client_pool Raylet client pool to construct connections to raylets. + /// \param actor_schedule_strategy Actor schedule strategy. /// \param client_factory Factory to create remote core worker client, default factor /// will be used if not set. explicit GcsActorScheduler( @@ -98,6 +100,7 @@ class GcsActorScheduler : public GcsActorSchedulerInterface { std::function)> schedule_failure_handler, std::function)> schedule_success_handler, std::shared_ptr raylet_client_pool, + std::shared_ptr actor_schedule_strategy, rpc::ClientFactoryFn client_factory = nullptr); virtual ~GcsActorScheduler() = default; @@ -286,6 +289,8 @@ class GcsActorScheduler : public GcsActorSchedulerInterface { absl::flat_hash_set nodes_of_releasing_unused_workers_; /// The cached raylet clients used to communicate with raylet. std::shared_ptr raylet_client_pool_; + /// The actor schedule strategy. + std::shared_ptr actor_schedule_strategy_; /// The cached core worker clients which are used to communicate with leased worker. rpc::CoreWorkerClientPool core_worker_clients_; }; diff --git a/src/ray/gcs/gcs_server/gcs_server.cc b/src/ray/gcs/gcs_server/gcs_server.cc index c6ff254fb..dba927143 100644 --- a/src/ray/gcs/gcs_server/gcs_server.cc +++ b/src/ray/gcs/gcs_server/gcs_server.cc @@ -181,6 +181,8 @@ void GcsServer::InitGcsJobManager() { void GcsServer::InitGcsActorManager(const GcsInitData &gcs_init_data) { RAY_CHECK(gcs_table_storage_ && gcs_pub_sub_ && gcs_node_manager_); + auto actor_schedule_strategy = + std::make_shared(gcs_node_manager_); auto scheduler = std::make_shared( main_service_, gcs_table_storage_->ActorTable(), *gcs_node_manager_, gcs_pub_sub_, /*schedule_failure_handler=*/ @@ -195,7 +197,7 @@ void GcsServer::InitGcsActorManager(const GcsInitData &gcs_init_data) { [this](std::shared_ptr actor) { gcs_actor_manager_->OnActorCreationSuccess(std::move(actor)); }, - raylet_client_pool_, + raylet_client_pool_, actor_schedule_strategy, /*client_factory=*/ [this](const rpc::Address &address) { return std::make_shared(address, client_call_manager_); diff --git a/src/ray/gcs/gcs_server/gcs_server.h b/src/ray/gcs/gcs_server/gcs_server.h index 1527ca7cf..9a34afc3a 100644 --- a/src/ray/gcs/gcs_server/gcs_server.h +++ b/src/ray/gcs/gcs_server/gcs_server.h @@ -124,7 +124,7 @@ class GcsServer { /// Print debug info periodically. void PrintDebugInfo(); - /// Gcs server configuration + /// Gcs server configuration. GcsServerConfig config_; /// The main io service to drive event posted from grpc threads. boost::asio::io_context &main_service_; @@ -135,7 +135,7 @@ class GcsServer { rpc::GrpcServer rpc_server_; /// The `ClientCallManager` object that is shared by all `NodeManagerWorkerClient`s. rpc::ClientCallManager client_call_manager_; - /// Node manager client pool + /// Node manager client pool. std::shared_ptr raylet_client_pool_; /// The gcs resource manager. std::shared_ptr gcs_resource_manager_; @@ -145,37 +145,37 @@ class GcsServer { std::shared_ptr gcs_heartbeat_manager_; /// The gcs redis failure detector. std::shared_ptr gcs_redis_failure_detector_; - /// The gcs actor manager + /// The gcs actor manager. std::shared_ptr gcs_actor_manager_; - /// The gcs placement group manager + /// The gcs placement group manager. std::shared_ptr gcs_placement_group_manager_; - /// Job info handler and service + /// Job info handler and service. std::unique_ptr gcs_job_manager_; std::unique_ptr job_info_service_; - /// Actor info service + /// Actor info service. std::unique_ptr actor_info_service_; - /// Node info handler and service + /// Node info handler and service. std::unique_ptr node_info_service_; - /// Node resource info handler and service + /// Node resource info handler and service. std::unique_ptr node_resource_info_service_; - /// Heartbeat info handler and service + /// Heartbeat info handler and service. std::unique_ptr heartbeat_info_service_; - /// Object info handler and service + /// Object info handler and service. std::unique_ptr gcs_object_manager_; std::unique_ptr object_info_service_; - /// Task info handler and service + /// Task info handler and service. std::unique_ptr task_info_handler_; std::unique_ptr task_info_service_; - /// Stats handler and service + /// Stats handler and service. std::unique_ptr stats_handler_; std::unique_ptr stats_service_; - /// The gcs worker manager + /// The gcs worker manager. std::unique_ptr gcs_worker_manager_; - /// Worker info service + /// Worker info service. std::unique_ptr worker_info_service_; - /// Placement Group info handler and service + /// Placement Group info handler and service. std::unique_ptr placement_group_info_service_; - /// Backend client + /// Backend client. std::shared_ptr redis_client_; /// A publisher for publishing gcs messages. std::shared_ptr gcs_pub_sub_; diff --git a/src/ray/gcs/gcs_server/test/gcs_actor_scheduler_test.cc b/src/ray/gcs/gcs_server/test/gcs_actor_scheduler_test.cc index fffef0db7..d84f99b3f 100644 --- a/src/ray/gcs/gcs_server/test/gcs_actor_scheduler_test.cc +++ b/src/ray/gcs/gcs_server/test/gcs_actor_scheduler_test.cc @@ -28,6 +28,8 @@ class GcsActorSchedulerTest : public ::testing::Test { gcs_table_storage_ = std::make_shared(redis_client_); gcs_node_manager_ = std::make_shared(gcs_pub_sub_, gcs_table_storage_); + gcs_actor_schedule_strategy_ = + std::make_shared(gcs_node_manager_); store_client_ = std::make_shared(io_service_); gcs_actor_table_ = std::make_shared(store_client_); @@ -43,7 +45,7 @@ class GcsActorSchedulerTest : public ::testing::Test { [this](std::shared_ptr actor) { success_actors_.emplace_back(std::move(actor)); }, - raylet_client_pool_, + raylet_client_pool_, gcs_actor_schedule_strategy_, /*client_factory=*/ [this](const rpc::Address &address) { return worker_client_; }); } @@ -55,6 +57,7 @@ class GcsActorSchedulerTest : public ::testing::Test { std::shared_ptr raylet_client_; std::shared_ptr worker_client_; std::shared_ptr gcs_node_manager_; + std::shared_ptr gcs_actor_schedule_strategy_; std::shared_ptr gcs_actor_scheduler_; std::vector> success_actors_; std::vector> failure_actors_; diff --git a/src/ray/gcs/gcs_server/test/gcs_server_test_util.h b/src/ray/gcs/gcs_server/test/gcs_server_test_util.h index ede33b395..ef50d4cb2 100644 --- a/src/ray/gcs/gcs_server/test/gcs_server_test_util.h +++ b/src/ray/gcs/gcs_server/test/gcs_server_test_util.h @@ -21,6 +21,7 @@ #include "ray/common/task/task_util.h" #include "ray/common/test_util.h" #include "ray/gcs/gcs_server/gcs_actor_manager.h" +#include "ray/gcs/gcs_server/gcs_actor_schedule_strategy.h" #include "ray/gcs/gcs_server/gcs_actor_scheduler.h" #include "ray/gcs/gcs_server/gcs_node_manager.h" #include "ray/gcs/gcs_server/gcs_placement_group_manager.h"