From 646c4201ac3279713800cd84113d6136e69c5208 Mon Sep 17 00:00:00 2001 From: fangfengbin <869218239a@zju.edu.cn> Date: Wed, 23 Dec 2020 11:25:01 +0800 Subject: [PATCH] [GCS]Decouple gcs resource manager and gcs node manager (#13012) --- src/ray/gcs/gcs_server/gcs_node_manager.cc | 38 ++++++++----------- src/ray/gcs/gcs_server/gcs_node_manager.h | 24 +++++++----- src/ray/gcs/gcs_server/gcs_server.cc | 12 +++++- .../test/gcs_actor_scheduler_test.cc | 4 +- .../gcs_server/test/gcs_node_manager_test.cc | 6 +-- .../test/gcs_object_manager_test.cc | 4 +- .../test/gcs_placement_group_manager_test.cc | 4 +- .../gcs_placement_group_scheduler_test.cc | 17 +++++---- 8 files changed, 58 insertions(+), 51 deletions(-) diff --git a/src/ray/gcs/gcs_server/gcs_node_manager.cc b/src/ray/gcs/gcs_server/gcs_node_manager.cc index 57f878d60..322b0349f 100644 --- a/src/ray/gcs/gcs_server/gcs_node_manager.cc +++ b/src/ray/gcs/gcs_server/gcs_node_manager.cc @@ -23,16 +23,14 @@ namespace ray { namespace gcs { ////////////////////////////////////////////////////////////////////////////////////////// -GcsNodeManager::GcsNodeManager( - boost::asio::io_service &main_io_service, std::shared_ptr gcs_pub_sub, - std::shared_ptr gcs_table_storage, - std::shared_ptr gcs_resource_manager) +GcsNodeManager::GcsNodeManager(boost::asio::io_service &main_io_service, + std::shared_ptr gcs_pub_sub, + std::shared_ptr gcs_table_storage) : resource_timer_(main_io_service), light_report_resource_usage_enabled_( RayConfig::instance().light_report_resource_usage_enabled()), gcs_pub_sub_(gcs_pub_sub), - gcs_table_storage_(gcs_table_storage), - gcs_resource_manager_(gcs_resource_manager) { + gcs_table_storage_(gcs_table_storage) { SendBatchedResourceUsage(); } @@ -104,10 +102,19 @@ void GcsNodeManager::HandleReportResourceUsage( auto resources_data = std::make_shared(); resources_data->CopyFrom(request.resources()); - UpdateNodeResourceUsage(node_id, request); + // We use `node_resource_usages_` to filter out the nodes that report resource + // information for the first time. `UpdateNodeResourceUsage` will modify + // `node_resource_usages_`, so we need to do it before `UpdateNodeResourceUsage`. + if (!light_report_resource_usage_enabled_ || + node_resource_usages_.count(node_id) == 0 || + resources_data->resources_available_changed()) { + const auto &resource_changed = MapFromProtobuf(resources_data->resources_available()); + for (auto &listener : node_resource_changed_listeners_) { + listener(node_id, resource_changed); + } + } - // Update node realtime resources. - UpdateNodeRealtimeResources(node_id, *resources_data); + UpdateNodeResourceUsage(node_id, request); if (!light_report_resource_usage_enabled_ || resources_data->should_global_gc() || resources_data->resources_total_size() > 0 || @@ -240,7 +247,6 @@ void GcsNodeManager::AddNode(std::shared_ptr node) { for (auto &listener : node_added_listeners_) { listener(node); } - gcs_resource_manager_->OnNodeAdd(*node); } } @@ -255,8 +261,6 @@ std::shared_ptr GcsNodeManager::RemoveNode( stats::NodeFailureTotal.Record(1); // Remove from alive nodes. alive_nodes_.erase(iter); - // Remove from cluster resources. - gcs_resource_manager_->OnNodeDead(node_id); resources_buffer_.erase(node_id); node_resource_usages_.erase(node_id); if (!is_intended) { @@ -313,16 +317,6 @@ void GcsNodeManager::Initialize(const GcsInitData &gcs_init_data) { const std::pair &right) { return left.second < right.second; }); } -void GcsNodeManager::UpdateNodeRealtimeResources( - const NodeID &node_id, const rpc::ResourcesData &resource_data) { - if (!light_report_resource_usage_enabled_ || - gcs_resource_manager_->GetClusterResources().count(node_id) == 0 || - resource_data.resources_available_changed()) { - gcs_resource_manager_->SetAvailableResources( - node_id, ResourceSet(MapFromProtobuf(resource_data.resources_available()))); - } -} - void GcsNodeManager::UpdatePlacementGroupLoad( const std::shared_ptr placement_group_load) { placement_group_load_ = absl::make_optional(placement_group_load); diff --git a/src/ray/gcs/gcs_server/gcs_node_manager.h b/src/ray/gcs/gcs_server/gcs_node_manager.h index 3c62af3a7..8b99eaa13 100644 --- a/src/ray/gcs/gcs_server/gcs_node_manager.h +++ b/src/ray/gcs/gcs_server/gcs_node_manager.h @@ -39,11 +39,9 @@ class GcsNodeManager : public rpc::NodeInfoHandler { /// \param main_io_service The main event loop. /// \param gcs_pub_sub GCS message publisher. /// \param gcs_table_storage GCS table external storage accessor. - /// \param gcs_resource_manager GCS resource manager. explicit GcsNodeManager(boost::asio::io_service &main_io_service, std::shared_ptr gcs_pub_sub, - std::shared_ptr gcs_table_storage, - std::shared_ptr gcs_resource_manager); + std::shared_ptr gcs_table_storage); /// Handle register rpc request come from raylet. void HandleRegisterNode(const rpc::RegisterNodeRequest &request, @@ -135,16 +133,22 @@ class GcsNodeManager : public rpc::NodeInfoHandler { node_added_listeners_.emplace_back(std::move(listener)); } + /// Add listener to monitor the resource change of nodes. + /// + /// \param listener The handler which process the resource change of nodes. + void AddNodeResourceChangedListener( + std::function &)> + listener) { + RAY_CHECK(listener); + node_resource_changed_listeners_.emplace_back(std::move(listener)); + } + /// Initialize with the gcs tables data synchronously. /// This should be called when GCS server restarts after a failure. /// /// \param gcs_init_data. void Initialize(const GcsInitData &gcs_init_data); - // Update node realtime resources. - void UpdateNodeRealtimeResources(const NodeID &node_id, - const rpc::ResourcesData &heartbeat); - /// Update the placement group load information so that it will be reported through /// heartbeat. /// @@ -185,12 +189,14 @@ class GcsNodeManager : public rpc::NodeInfoHandler { /// Listeners which monitors the removal of nodes. std::vector)>> node_removed_listeners_; + /// Listeners which monitors the resource change of nodes. + std::vector &)>> + node_resource_changed_listeners_; /// A publisher for publishing gcs messages. std::shared_ptr gcs_pub_sub_; /// Storage for GCS tables. std::shared_ptr gcs_table_storage_; - /// Gcs resource manager. - std::shared_ptr gcs_resource_manager_; /// Placement group load information that is used for autoscaler. absl::optional> placement_group_load_; diff --git a/src/ray/gcs/gcs_server/gcs_server.cc b/src/ray/gcs/gcs_server/gcs_server.cc index 23a12f6ec..71e2a6d81 100644 --- a/src/ray/gcs/gcs_server/gcs_server.cc +++ b/src/ray/gcs/gcs_server/gcs_server.cc @@ -133,8 +133,8 @@ void GcsServer::Stop() { void GcsServer::InitGcsNodeManager(const GcsInitData &gcs_init_data) { RAY_CHECK(redis_gcs_client_ && gcs_table_storage_ && gcs_pub_sub_); - gcs_node_manager_ = std::make_shared( - main_service_, gcs_pub_sub_, gcs_table_storage_, gcs_resource_manager_); + gcs_node_manager_ = + std::make_shared(main_service_, gcs_pub_sub_, gcs_table_storage_); // Initialize by gcs tables data. gcs_node_manager_->Initialize(gcs_init_data); // Register service. @@ -292,6 +292,7 @@ void GcsServer::InstallEventListeners() { gcs_node_manager_->AddNodeAddedListener([this](std::shared_ptr node) { // Because a new node has been added, we need to try to schedule the pending // placement groups and the pending actors. + gcs_resource_manager_->OnNodeAdd(*node); gcs_placement_group_manager_->SchedulePendingPlacementGroups(); gcs_actor_manager_->SchedulePendingActors(); gcs_heartbeat_manager_->AddNode(NodeID::FromBinary(node->node_id())); @@ -301,10 +302,17 @@ void GcsServer::InstallEventListeners() { auto node_id = NodeID::FromBinary(node->node_id()); // All of the related placement groups and actors should be reconstructed when a // node is removed from the GCS. + gcs_resource_manager_->OnNodeDead(node_id); gcs_placement_group_manager_->OnNodeDead(node_id); gcs_actor_manager_->OnNodeDead(node_id); raylet_client_pool_->Disconnect(NodeID::FromBinary(node->node_id())); }); + gcs_node_manager_->AddNodeResourceChangedListener( + [this](const NodeID &node_id, + const std::unordered_map &resource_changed) { + gcs_resource_manager_->SetAvailableResources(node_id, + ResourceSet(resource_changed)); + }); // Install worker event listener. gcs_worker_manager_->AddWorkerDeadListener( 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 4ddba0627..7bb1ca716 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 @@ -27,8 +27,8 @@ class GcsActorSchedulerTest : public ::testing::Test { gcs_pub_sub_ = std::make_shared(redis_client_); gcs_table_storage_ = std::make_shared(redis_client_); gcs_resource_manager_ = std::make_shared(nullptr, nullptr); - gcs_node_manager_ = std::make_shared( - io_service_, gcs_pub_sub_, gcs_table_storage_, gcs_resource_manager_); + gcs_node_manager_ = std::make_shared(io_service_, gcs_pub_sub_, + gcs_table_storage_); store_client_ = std::make_shared(io_service_); gcs_actor_table_ = std::make_shared(store_client_); diff --git a/src/ray/gcs/gcs_server/test/gcs_node_manager_test.cc b/src/ray/gcs/gcs_server/test/gcs_node_manager_test.cc index 74c4b8fd1..25f80733a 100644 --- a/src/ray/gcs/gcs_server/test/gcs_node_manager_test.cc +++ b/src/ray/gcs/gcs_server/test/gcs_node_manager_test.cc @@ -35,8 +35,7 @@ class GcsNodeManagerTest : public ::testing::Test { TEST_F(GcsNodeManagerTest, TestManagement) { boost::asio::io_service io_service; - gcs::GcsNodeManager node_manager(io_service, gcs_pub_sub_, gcs_table_storage_, - gcs_resource_manager_); + gcs::GcsNodeManager node_manager(io_service, gcs_pub_sub_, gcs_table_storage_); // Test Add/Get/Remove functionality. auto node = Mocker::GenNodeInfo(); auto node_id = NodeID::FromBinary(node->node_id()); @@ -82,8 +81,7 @@ TEST_F(GcsNodeManagerTest, TestManagement) { TEST_F(GcsNodeManagerTest, TestListener) { boost::asio::io_service io_service; - gcs::GcsNodeManager node_manager(io_service, gcs_pub_sub_, gcs_table_storage_, - gcs_resource_manager_); + gcs::GcsNodeManager node_manager(io_service, gcs_pub_sub_, gcs_table_storage_); // Test AddNodeAddedListener. int node_count = 1000; std::vector> added_nodes; diff --git a/src/ray/gcs/gcs_server/test/gcs_object_manager_test.cc b/src/ray/gcs/gcs_server/test/gcs_object_manager_test.cc index 15f96a6a8..700fdfc10 100644 --- a/src/ray/gcs/gcs_server/test/gcs_object_manager_test.cc +++ b/src/ray/gcs/gcs_server/test/gcs_object_manager_test.cc @@ -55,8 +55,8 @@ class GcsObjectManagerTest : public ::testing::Test { void SetUp() override { gcs_table_storage_ = std::make_shared(io_service_); gcs_resource_manager_ = std::make_shared(nullptr, nullptr); - gcs_node_manager_ = std::make_shared( - io_service_, gcs_pub_sub_, gcs_table_storage_, gcs_resource_manager_); + gcs_node_manager_ = std::make_shared(io_service_, gcs_pub_sub_, + gcs_table_storage_); gcs_object_manager_ = std::make_shared( gcs_table_storage_, gcs_pub_sub_, *gcs_node_manager_); GenTestData(); diff --git a/src/ray/gcs/gcs_server/test/gcs_placement_group_manager_test.cc b/src/ray/gcs/gcs_server/test/gcs_placement_group_manager_test.cc index e74b5fe1b..70bfdce31 100644 --- a/src/ray/gcs/gcs_server/test/gcs_placement_group_manager_test.cc +++ b/src/ray/gcs/gcs_server/test/gcs_placement_group_manager_test.cc @@ -69,8 +69,8 @@ class GcsPlacementGroupManagerTest : public ::testing::Test { gcs_pub_sub_ = std::make_shared(redis_client_); gcs_table_storage_ = std::make_shared(io_service_); gcs_resource_manager_ = std::make_shared(nullptr, nullptr); - gcs_node_manager_ = std::make_shared( - io_service_, gcs_pub_sub_, gcs_table_storage_, gcs_resource_manager_); + gcs_node_manager_ = std::make_shared(io_service_, gcs_pub_sub_, + gcs_table_storage_); gcs_placement_group_manager_.reset( new gcs::GcsPlacementGroupManager(io_service_, mock_placement_group_scheduler_, gcs_table_storage_, *gcs_node_manager_)); diff --git a/src/ray/gcs/gcs_server/test/gcs_placement_group_scheduler_test.cc b/src/ray/gcs/gcs_server/test/gcs_placement_group_scheduler_test.cc index ef81f8887..6a0a5839b 100644 --- a/src/ray/gcs/gcs_server/test/gcs_placement_group_scheduler_test.cc +++ b/src/ray/gcs/gcs_server/test/gcs_placement_group_scheduler_test.cc @@ -40,8 +40,8 @@ class GcsPlacementGroupSchedulerTest : public ::testing::Test { gcs_table_storage_ = std::make_shared(io_service_); gcs_pub_sub_ = std::make_shared(redis_client_); gcs_resource_manager_ = std::make_shared(nullptr, nullptr); - gcs_node_manager_ = std::make_shared( - io_service_, gcs_pub_sub_, gcs_table_storage_, gcs_resource_manager_); + gcs_node_manager_ = std::make_shared(io_service_, gcs_pub_sub_, + gcs_table_storage_); gcs_table_storage_ = std::make_shared(io_service_); store_client_ = std::make_shared(io_service_); raylet_client_pool_ = std::make_shared( @@ -98,12 +98,13 @@ class GcsPlacementGroupSchedulerTest : public ::testing::Test { void AddNode(const std::shared_ptr &node, int cpu_num = 10) { gcs_node_manager_->AddNode(node); - rpc::ResourcesData resource; - resource.set_node_id(node->node_id()); - (*resource.mutable_resources_available())["CPU"] = cpu_num; - resource.set_resources_available_changed(true); - gcs_node_manager_->UpdateNodeRealtimeResources(NodeID::FromBinary(node->node_id()), - resource); + gcs_resource_manager_->OnNodeAdd(*node); + + const auto &node_id = NodeID::FromBinary(node->node_id()); + std::unordered_map resource_map; + resource_map["CPU"] = cpu_num; + ResourceSet resources(resource_map); + gcs_resource_manager_->SetAvailableResources(node_id, resources); } void ScheduleFailedWithZeroNodeTest(rpc::PlacementStrategy strategy) {