From 1e02b28abe9d05c6d95805ffc54741069794efd6 Mon Sep 17 00:00:00 2001 From: fangfengbin <869218239a@zju.edu.cn> Date: Sun, 13 Dec 2020 20:37:34 +0800 Subject: [PATCH] [GCS]Move node resource info to gcs resource manager (#12775) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * add part code * add part code * fix review comments * fix ut bug * rebase master * add part code * fix ut bug * fix ut bug * fix review comments * fix review comment Co-authored-by: 灵洵 --- src/ray/gcs/accessor.h | 111 +++++---- src/ray/gcs/gcs_client.h | 8 + .../gcs/gcs_client/global_state_accessor.cc | 7 +- .../gcs/gcs_client/service_based_accessor.cc | 227 +++++++++--------- .../gcs/gcs_client/service_based_accessor.h | 56 +++-- .../gcs_client/service_based_gcs_client.cc | 2 + .../test/global_state_accessor_test.cc | 4 +- .../test/service_based_gcs_client_test.cc | 23 +- src/ray/gcs/gcs_server/gcs_node_manager.cc | 119 +-------- src/ray/gcs/gcs_server/gcs_node_manager.h | 36 +-- .../gcs/gcs_server/gcs_resource_manager.cc | 159 +++++++++++- src/ray/gcs/gcs_server/gcs_resource_manager.h | 113 ++++++--- src/ray/gcs/gcs_server/gcs_server.cc | 15 +- src/ray/gcs/gcs_server/gcs_server.h | 4 +- .../test/gcs_actor_scheduler_test.cc | 2 +- .../gcs_server/test/gcs_node_manager_test.cc | 1 + .../test/gcs_object_manager_test.cc | 2 +- .../test/gcs_placement_group_manager_test.cc | 2 +- .../gcs_placement_group_scheduler_test.cc | 2 +- .../test/gcs_resource_manager_test.cc | 4 +- .../gcs_server/test/gcs_server_test_util.h | 23 -- src/ray/gcs/redis_accessor.cc | 15 +- src/ray/gcs/redis_accessor.h | 59 +++-- src/ray/gcs/redis_gcs_client.cc | 1 + .../gcs/test/redis_node_info_accessor_test.cc | 53 ++-- src/ray/protobuf/gcs_service.proto | 68 +++--- src/ray/raylet/node_manager.cc | 16 +- src/ray/raylet/raylet.cc | 4 +- src/ray/rpc/gcs_server/gcs_rpc_client.h | 32 ++- src/ray/rpc/gcs_server/gcs_rpc_server.h | 76 ++++-- 30 files changed, 701 insertions(+), 543 deletions(-) diff --git a/src/ray/gcs/accessor.h b/src/ray/gcs/accessor.h index ce932fd59..655c47aa7 100644 --- a/src/ray/gcs/accessor.h +++ b/src/ray/gcs/accessor.h @@ -488,51 +488,6 @@ class NodeInfoAccessor { /// \return Whether the node is removed. virtual bool IsRemoved(const NodeID &node_id) const = 0; - // TODO(micafan) Define ResourceMap in GCS proto. - typedef std::unordered_map> - ResourceMap; - - /// Get node's resources from GCS asynchronously. - /// - /// \param node_id The ID of node to lookup dynamic resources. - /// \param callback Callback that will be called after lookup finishes. - /// \return Status - virtual Status AsyncGetResources(const NodeID &node_id, - const OptionalItemCallback &callback) = 0; - - /// Get available resources of all nodes from GCS asynchronously. - /// - /// \param callback Callback that will be called after lookup finishes. - /// \return Status - virtual Status AsyncGetAllAvailableResources( - const MultiItemCallback &callback) = 0; - - /// Update resources of node in GCS asynchronously. - /// - /// \param node_id The ID of node to update dynamic resources. - /// \param resources The dynamic resources of node to be updated. - /// \param callback Callback that will be called after update finishes. - virtual Status AsyncUpdateResources(const NodeID &node_id, const ResourceMap &resources, - const StatusCallback &callback) = 0; - - /// Delete resources of a node from GCS asynchronously. - /// - /// \param node_id The ID of node to delete resources from GCS. - /// \param resource_names The names of resource to be deleted. - /// \param callback Callback that will be called after delete finishes. - virtual Status AsyncDeleteResources(const NodeID &node_id, - const std::vector &resource_names, - const StatusCallback &callback) = 0; - - /// Subscribe to node resource changes. - /// - /// \param subscribe Callback that will be called when any resource is updated. - /// \param done Callback that will be called when subscription is complete. - /// \return Status - virtual Status AsyncSubscribeToResources( - const ItemCallback &subscribe, - const StatusCallback &done) = 0; - /// Report heartbeat of a node to GCS asynchronously. /// /// \param data_ptr The heartbeat that will be reported to GCS. @@ -612,6 +567,72 @@ class NodeInfoAccessor { std::make_shared(); }; +/// \class NodeResourceInfoAccessor +/// `NodeResourceInfoAccessor` is a sub-interface of `GcsClient`. +/// This class includes all the methods that are related to accessing +/// node resource information in the GCS. +class NodeResourceInfoAccessor { + public: + virtual ~NodeResourceInfoAccessor() = default; + + // TODO(micafan) Define ResourceMap in GCS proto. + typedef std::unordered_map> + ResourceMap; + + /// Get node's resources from GCS asynchronously. + /// + /// \param node_id The ID of node to lookup dynamic resources. + /// \param callback Callback that will be called after lookup finishes. + /// \return Status + virtual Status AsyncGetResources(const NodeID &node_id, + const OptionalItemCallback &callback) = 0; + + /// Get available resources of all nodes from GCS asynchronously. + /// + /// \param callback Callback that will be called after lookup finishes. + /// \return Status + virtual Status AsyncGetAllAvailableResources( + const MultiItemCallback &callback) = 0; + + /// Update resources of node in GCS asynchronously. + /// + /// \param node_id The ID of node to update dynamic resources. + /// \param resources The dynamic resources of node to be updated. + /// \param callback Callback that will be called after update finishes. + virtual Status AsyncUpdateResources(const NodeID &node_id, const ResourceMap &resources, + const StatusCallback &callback) = 0; + + /// Delete resources of a node from GCS asynchronously. + /// + /// \param node_id The ID of node to delete resources from GCS. + /// \param resource_names The names of resource to be deleted. + /// \param callback Callback that will be called after delete finishes. + virtual Status AsyncDeleteResources(const NodeID &node_id, + const std::vector &resource_names, + const StatusCallback &callback) = 0; + + /// Subscribe to node resource changes. + /// + /// \param subscribe Callback that will be called when any resource is updated. + /// \param done Callback that will be called when subscription is complete. + /// \return Status + virtual Status AsyncSubscribeToResources( + const ItemCallback &subscribe, + const StatusCallback &done) = 0; + + /// Reestablish subscription. + /// This should be called when GCS server restarts from a failure. + /// PubSub server restart will cause GCS server restart. In this case, we need to + /// resubscribe from PubSub server, otherwise we only need to fetch data from GCS + /// server. + /// + /// \param is_pubsub_server_restarted Whether pubsub server is restarted. + virtual void AsyncResubscribe(bool is_pubsub_server_restarted) = 0; + + protected: + NodeResourceInfoAccessor() = default; +}; + /// \class ErrorInfoAccessor /// `ErrorInfoAccessor` is a sub-interface of `GcsClient`. /// This class includes all the methods that are related to accessing diff --git a/src/ray/gcs/gcs_client.h b/src/ray/gcs/gcs_client.h index d88389c46..5f804606e 100644 --- a/src/ray/gcs/gcs_client.h +++ b/src/ray/gcs/gcs_client.h @@ -107,6 +107,13 @@ class GcsClient : public std::enable_shared_from_this { return *node_accessor_; } + /// Get the sub-interface for accessing node resource information in GCS. + /// This function is thread safe. + NodeResourceInfoAccessor &NodeResources() { + RAY_CHECK(node_resource_accessor_ != nullptr); + return *node_resource_accessor_; + } + /// Get the sub-interface for accessing task information in GCS. /// This function is thread safe. TaskInfoAccessor &Tasks() { @@ -157,6 +164,7 @@ class GcsClient : public std::enable_shared_from_this { std::unique_ptr job_accessor_; std::unique_ptr object_accessor_; std::unique_ptr node_accessor_; + std::unique_ptr node_resource_accessor_; std::unique_ptr task_accessor_; std::unique_ptr error_accessor_; std::unique_ptr stats_accessor_; diff --git a/src/ray/gcs/gcs_client/global_state_accessor.cc b/src/ray/gcs/gcs_client/global_state_accessor.cc index 8940ef6d0..8d188ba07 100644 --- a/src/ray/gcs/gcs_client/global_state_accessor.cc +++ b/src/ray/gcs/gcs_client/global_state_accessor.cc @@ -128,7 +128,8 @@ std::string GlobalStateAccessor::GetNodeResourceInfo(const NodeID &node_id) { auto on_done = [&node_resource_map, &promise]( const Status &status, - const boost::optional &result) { + const boost::optional + &result) { RAY_CHECK_OK(status); if (result) { auto result_value = result.get(); @@ -138,7 +139,7 @@ std::string GlobalStateAccessor::GetNodeResourceInfo(const NodeID &node_id) { } promise.set_value(); }; - RAY_CHECK_OK(gcs_client_->Nodes().AsyncGetResources(node_id, on_done)); + RAY_CHECK_OK(gcs_client_->NodeResources().AsyncGetResources(node_id, on_done)); promise.get_future().get(); return node_resource_map.SerializeAsString(); } @@ -146,7 +147,7 @@ std::string GlobalStateAccessor::GetNodeResourceInfo(const NodeID &node_id) { std::vector GlobalStateAccessor::GetAllAvailableResources() { std::vector available_resources; std::promise promise; - RAY_CHECK_OK(gcs_client_->Nodes().AsyncGetAllAvailableResources( + RAY_CHECK_OK(gcs_client_->NodeResources().AsyncGetAllAvailableResources( TransformForMultiItemCallback(available_resources, promise))); promise.get_future().get(); diff --git a/src/ray/gcs/gcs_client/service_based_accessor.cc b/src/ray/gcs/gcs_client/service_based_accessor.cc index 66179249b..cc2907c63 100644 --- a/src/ray/gcs/gcs_client/service_based_accessor.cc +++ b/src/ray/gcs/gcs_client/service_based_accessor.cc @@ -501,111 +501,6 @@ bool ServiceBasedNodeInfoAccessor::IsRemoved(const NodeID &node_id) const { return removed_nodes_.count(node_id) == 1; } -Status ServiceBasedNodeInfoAccessor::AsyncGetResources( - const NodeID &node_id, const OptionalItemCallback &callback) { - RAY_LOG(DEBUG) << "Getting node resources, node id = " << node_id; - rpc::GetResourcesRequest request; - request.set_node_id(node_id.Binary()); - client_impl_->GetGcsRpcClient().GetResources( - request, - [node_id, callback](const Status &status, const rpc::GetResourcesReply &reply) { - ResourceMap resource_map; - for (auto resource : reply.resources()) { - resource_map[resource.first] = - std::make_shared(resource.second); - } - callback(status, resource_map); - RAY_LOG(DEBUG) << "Finished getting node resources, status = " << status - << ", node id = " << node_id; - }); - return Status::OK(); -} - -Status ServiceBasedNodeInfoAccessor::AsyncGetAllAvailableResources( - const MultiItemCallback &callback) { - rpc::GetAllAvailableResourcesRequest request; - client_impl_->GetGcsRpcClient().GetAllAvailableResources( - request, - [callback](const Status &status, const rpc::GetAllAvailableResourcesReply &reply) { - std::vector result = - VectorFromProtobuf(reply.resources_list()); - callback(status, result); - RAY_LOG(DEBUG) << "Finished getting available resources of all nodes, status = " - << status; - }); - return Status::OK(); -} - -Status ServiceBasedNodeInfoAccessor::AsyncUpdateResources( - const NodeID &node_id, const ResourceMap &resources, const StatusCallback &callback) { - RAY_LOG(DEBUG) << "Updating node resources, node id = " << node_id; - rpc::UpdateResourcesRequest request; - request.set_node_id(node_id.Binary()); - for (auto &resource : resources) { - (*request.mutable_resources())[resource.first] = *resource.second; - } - - auto operation = [this, request, node_id, - callback](const SequencerDoneCallback &done_callback) { - client_impl_->GetGcsRpcClient().UpdateResources( - request, [node_id, callback, done_callback]( - const Status &status, const rpc::UpdateResourcesReply &reply) { - if (callback) { - callback(status); - } - RAY_LOG(DEBUG) << "Finished updating node resources, status = " << status - << ", node id = " << node_id; - done_callback(); - }); - }; - - sequencer_.Post(node_id, operation); - return Status::OK(); -} - -Status ServiceBasedNodeInfoAccessor::AsyncDeleteResources( - const NodeID &node_id, const std::vector &resource_names, - const StatusCallback &callback) { - RAY_LOG(DEBUG) << "Deleting node resources, node id = " << node_id; - rpc::DeleteResourcesRequest request; - request.set_node_id(node_id.Binary()); - for (auto &resource_name : resource_names) { - request.add_resource_name_list(resource_name); - } - - auto operation = [this, request, node_id, - callback](const SequencerDoneCallback &done_callback) { - client_impl_->GetGcsRpcClient().DeleteResources( - request, [node_id, callback, done_callback]( - const Status &status, const rpc::DeleteResourcesReply &reply) { - if (callback) { - callback(status); - } - RAY_LOG(DEBUG) << "Finished deleting node resources, status = " << status - << ", node id = " << node_id; - done_callback(); - }); - }; - - sequencer_.Post(node_id, operation); - return Status::OK(); -} - -Status ServiceBasedNodeInfoAccessor::AsyncSubscribeToResources( - const ItemCallback &subscribe, const StatusCallback &done) { - RAY_CHECK(subscribe != nullptr); - subscribe_resource_operation_ = [this, subscribe](const StatusCallback &done) { - auto on_subscribe = [subscribe](const std::string &id, const std::string &data) { - rpc::NodeResourceChange node_resource_change; - node_resource_change.ParseFromString(data); - subscribe(node_resource_change); - }; - return client_impl_->GetGcsPubSub().SubscribeAll(NODE_RESOURCE_CHANNEL, on_subscribe, - done); - }; - return subscribe_resource_operation_(done); -} - Status ServiceBasedNodeInfoAccessor::AsyncReportHeartbeat( const std::shared_ptr &data_ptr, const StatusCallback &callback) { @@ -764,9 +659,6 @@ void ServiceBasedNodeInfoAccessor::AsyncResubscribe(bool is_pubsub_server_restar RAY_CHECK_OK(subscribe_node_operation_( [this](const Status &status) { fetch_node_data_operation_(nullptr); })); } - if (subscribe_resource_operation_ != nullptr) { - RAY_CHECK_OK(subscribe_resource_operation_(nullptr)); - } if (subscribe_batch_resource_usage_operation_ != nullptr) { RAY_CHECK_OK(subscribe_batch_resource_usage_operation_(nullptr)); } @@ -814,6 +706,125 @@ Status ServiceBasedNodeInfoAccessor::AsyncGetInternalConfig( return Status::OK(); } +ServiceBasedNodeResourceInfoAccessor::ServiceBasedNodeResourceInfoAccessor( + ServiceBasedGcsClient *client_impl) + : client_impl_(client_impl) {} + +Status ServiceBasedNodeResourceInfoAccessor::AsyncGetResources( + const NodeID &node_id, const OptionalItemCallback &callback) { + RAY_LOG(DEBUG) << "Getting node resources, node id = " << node_id; + rpc::GetResourcesRequest request; + request.set_node_id(node_id.Binary()); + client_impl_->GetGcsRpcClient().GetResources( + request, + [node_id, callback](const Status &status, const rpc::GetResourcesReply &reply) { + ResourceMap resource_map; + for (auto resource : reply.resources()) { + resource_map[resource.first] = + std::make_shared(resource.second); + } + callback(status, resource_map); + RAY_LOG(DEBUG) << "Finished getting node resources, status = " << status + << ", node id = " << node_id; + }); + return Status::OK(); +} + +Status ServiceBasedNodeResourceInfoAccessor::AsyncGetAllAvailableResources( + const MultiItemCallback &callback) { + rpc::GetAllAvailableResourcesRequest request; + client_impl_->GetGcsRpcClient().GetAllAvailableResources( + request, + [callback](const Status &status, const rpc::GetAllAvailableResourcesReply &reply) { + std::vector result = + VectorFromProtobuf(reply.resources_list()); + callback(status, result); + RAY_LOG(DEBUG) << "Finished getting available resources of all nodes, status = " + << status; + }); + return Status::OK(); +} + +Status ServiceBasedNodeResourceInfoAccessor::AsyncUpdateResources( + const NodeID &node_id, const ResourceMap &resources, const StatusCallback &callback) { + RAY_LOG(DEBUG) << "Updating node resources, node id = " << node_id; + rpc::UpdateResourcesRequest request; + request.set_node_id(node_id.Binary()); + for (auto &resource : resources) { + (*request.mutable_resources())[resource.first] = *resource.second; + } + + auto operation = [this, request, node_id, + callback](const SequencerDoneCallback &done_callback) { + client_impl_->GetGcsRpcClient().UpdateResources( + request, [node_id, callback, done_callback]( + const Status &status, const rpc::UpdateResourcesReply &reply) { + if (callback) { + callback(status); + } + RAY_LOG(DEBUG) << "Finished updating node resources, status = " << status + << ", node id = " << node_id; + done_callback(); + }); + }; + + sequencer_.Post(node_id, operation); + return Status::OK(); +} + +Status ServiceBasedNodeResourceInfoAccessor::AsyncDeleteResources( + const NodeID &node_id, const std::vector &resource_names, + const StatusCallback &callback) { + RAY_LOG(DEBUG) << "Deleting node resources, node id = " << node_id; + rpc::DeleteResourcesRequest request; + request.set_node_id(node_id.Binary()); + for (auto &resource_name : resource_names) { + request.add_resource_name_list(resource_name); + } + + auto operation = [this, request, node_id, + callback](const SequencerDoneCallback &done_callback) { + client_impl_->GetGcsRpcClient().DeleteResources( + request, [node_id, callback, done_callback]( + const Status &status, const rpc::DeleteResourcesReply &reply) { + if (callback) { + callback(status); + } + RAY_LOG(DEBUG) << "Finished deleting node resources, status = " << status + << ", node id = " << node_id; + done_callback(); + }); + }; + + sequencer_.Post(node_id, operation); + return Status::OK(); +} + +Status ServiceBasedNodeResourceInfoAccessor::AsyncSubscribeToResources( + const ItemCallback &subscribe, const StatusCallback &done) { + RAY_CHECK(subscribe != nullptr); + subscribe_resource_operation_ = [this, subscribe](const StatusCallback &done) { + auto on_subscribe = [subscribe](const std::string &id, const std::string &data) { + rpc::NodeResourceChange node_resource_change; + node_resource_change.ParseFromString(data); + subscribe(node_resource_change); + }; + return client_impl_->GetGcsPubSub().SubscribeAll(NODE_RESOURCE_CHANNEL, on_subscribe, + done); + }; + return subscribe_resource_operation_(done); +} + +void ServiceBasedNodeResourceInfoAccessor::AsyncResubscribe( + bool is_pubsub_server_restarted) { + RAY_LOG(DEBUG) << "Reestablishing subscription for node resource info."; + // If the pub-sub server has also restarted, we need to resubscribe to the pub-sub + // server. + if (is_pubsub_server_restarted && subscribe_resource_operation_ != nullptr) { + RAY_CHECK_OK(subscribe_resource_operation_(nullptr)); + } +} + ServiceBasedTaskInfoAccessor::ServiceBasedTaskInfoAccessor( ServiceBasedGcsClient *client_impl) : client_impl_(client_impl) {} diff --git a/src/ray/gcs/gcs_client/service_based_accessor.h b/src/ray/gcs/gcs_client/service_based_accessor.h index 188850d06..f0e1f45bc 100644 --- a/src/ray/gcs/gcs_client/service_based_accessor.h +++ b/src/ray/gcs/gcs_client/service_based_accessor.h @@ -163,22 +163,6 @@ class ServiceBasedNodeInfoAccessor : public NodeInfoAccessor { bool IsRemoved(const NodeID &node_id) const override; - Status AsyncGetResources(const NodeID &node_id, - const OptionalItemCallback &callback) override; - - Status AsyncGetAllAvailableResources( - const MultiItemCallback &callback) override; - - Status AsyncUpdateResources(const NodeID &node_id, const ResourceMap &resources, - const StatusCallback &callback) override; - - Status AsyncDeleteResources(const NodeID &node_id, - const std::vector &resource_names, - const StatusCallback &callback) override; - - Status AsyncSubscribeToResources(const ItemCallback &subscribe, - const StatusCallback &done) override; - Status AsyncReportHeartbeat(const std::shared_ptr &data_ptr, const StatusCallback &callback) override; @@ -210,7 +194,6 @@ class ServiceBasedNodeInfoAccessor : public NodeInfoAccessor { /// Save the subscribe operation in this function, so we can call it again when PubSub /// server restarts from a failure. SubscribeOperation subscribe_node_operation_; - SubscribeOperation subscribe_resource_operation_; SubscribeOperation subscribe_batch_resource_usage_operation_; /// Save the fetch data operation in this function, so we can call it again when GCS @@ -234,8 +217,6 @@ class ServiceBasedNodeInfoAccessor : public NodeInfoAccessor { GcsNodeInfo local_node_info_; NodeID local_node_id_; - Sequencer sequencer_; - /// The callback to call when a new node is added or a node is removed. NodeChangeCallback node_change_callback_{nullptr}; @@ -245,6 +226,43 @@ class ServiceBasedNodeInfoAccessor : public NodeInfoAccessor { std::unordered_set removed_nodes_; }; +/// \class ServiceBasedNodeResourceInfoAccessor +/// ServiceBasedNodeResourceInfoAccessor is an implementation of +/// `NodeResourceInfoAccessor` that uses GCS Service as the backend. +class ServiceBasedNodeResourceInfoAccessor : public NodeResourceInfoAccessor { + public: + explicit ServiceBasedNodeResourceInfoAccessor(ServiceBasedGcsClient *client_impl); + + virtual ~ServiceBasedNodeResourceInfoAccessor() = default; + + Status AsyncGetResources(const NodeID &node_id, + const OptionalItemCallback &callback) override; + + Status AsyncGetAllAvailableResources( + const MultiItemCallback &callback) override; + + Status AsyncUpdateResources(const NodeID &node_id, const ResourceMap &resources, + const StatusCallback &callback) override; + + Status AsyncDeleteResources(const NodeID &node_id, + const std::vector &resource_names, + const StatusCallback &callback) override; + + Status AsyncSubscribeToResources(const ItemCallback &subscribe, + const StatusCallback &done) override; + + void AsyncResubscribe(bool is_pubsub_server_restarted) override; + + private: + /// Save the subscribe operation in this function, so we can call it again when PubSub + /// server restarts from a failure. + SubscribeOperation subscribe_resource_operation_; + + ServiceBasedGcsClient *client_impl_; + + Sequencer sequencer_; +}; + /// \class ServiceBasedTaskInfoAccessor /// ServiceBasedTaskInfoAccessor is an implementation of `TaskInfoAccessor` /// that uses GCS service as the backend. diff --git a/src/ray/gcs/gcs_client/service_based_gcs_client.cc b/src/ray/gcs/gcs_client/service_based_gcs_client.cc index 9bbbbde25..359f6cd81 100644 --- a/src/ray/gcs/gcs_client/service_based_gcs_client.cc +++ b/src/ray/gcs/gcs_client/service_based_gcs_client.cc @@ -57,6 +57,7 @@ Status ServiceBasedGcsClient::Connect(boost::asio::io_service &io_service) { job_accessor_->AsyncResubscribe(is_pubsub_server_restarted); actor_accessor_->AsyncResubscribe(is_pubsub_server_restarted); node_accessor_->AsyncResubscribe(is_pubsub_server_restarted); + node_resource_accessor_->AsyncResubscribe(is_pubsub_server_restarted); task_accessor_->AsyncResubscribe(is_pubsub_server_restarted); object_accessor_->AsyncResubscribe(is_pubsub_server_restarted); worker_accessor_->AsyncResubscribe(is_pubsub_server_restarted); @@ -70,6 +71,7 @@ Status ServiceBasedGcsClient::Connect(boost::asio::io_service &io_service) { job_accessor_.reset(new ServiceBasedJobInfoAccessor(this)); actor_accessor_.reset(new ServiceBasedActorInfoAccessor(this)); node_accessor_.reset(new ServiceBasedNodeInfoAccessor(this)); + node_resource_accessor_.reset(new ServiceBasedNodeResourceInfoAccessor(this)); task_accessor_.reset(new ServiceBasedTaskInfoAccessor(this)); object_accessor_.reset(new ServiceBasedObjectInfoAccessor(this)); stats_accessor_.reset(new ServiceBasedStatsInfoAccessor(this)); diff --git a/src/ray/gcs/gcs_client/test/global_state_accessor_test.cc b/src/ray/gcs/gcs_client/test/global_state_accessor_test.cc index 7bfc07545..0df0c7763 100644 --- a/src/ray/gcs/gcs_client/test/global_state_accessor_test.cc +++ b/src/ray/gcs/gcs_client/test/global_state_accessor_test.cc @@ -139,12 +139,12 @@ TEST_F(GlobalStateAccessorTest, TestNodeResourceTable) { RAY_CHECK_OK(gcs_client_->Nodes().AsyncRegister( *node_table_data, [&promise](Status status) { promise.set_value(status.ok()); })); WaitReady(promise.get_future(), timeout_ms_); - ray::gcs::NodeInfoAccessor::ResourceMap resources; + ray::gcs::NodeResourceInfoAccessor::ResourceMap resources; rpc::ResourceTableData resource_table_data; resource_table_data.set_resource_capacity(static_cast(index + 1) + 0.1); resources[std::to_string(index)] = std::make_shared(resource_table_data); - RAY_IGNORE_EXPR(gcs_client_->Nodes().AsyncUpdateResources( + RAY_IGNORE_EXPR(gcs_client_->NodeResources().AsyncUpdateResources( node_id, resources, [](Status status) { RAY_CHECK(status.ok()); })); } auto node_table = global_state_->GetAllNodeInfo(); diff --git a/src/ray/gcs/gcs_client/test/service_based_gcs_client_test.cc b/src/ray/gcs/gcs_client/test/service_based_gcs_client_test.cc index 32294b242..9df66dfff 100644 --- a/src/ray/gcs/gcs_client/test/service_based_gcs_client_test.cc +++ b/src/ray/gcs/gcs_client/test/service_based_gcs_client_test.cc @@ -286,18 +286,19 @@ class ServiceBasedGcsClientTest : public ::testing::Test { bool SubscribeToResources(const gcs::ItemCallback &subscribe) { std::promise promise; - RAY_CHECK_OK(gcs_client_->Nodes().AsyncSubscribeToResources( + RAY_CHECK_OK(gcs_client_->NodeResources().AsyncSubscribeToResources( subscribe, [&promise](Status status) { promise.set_value(status.ok()); })); return WaitReady(promise.get_future(), timeout_ms_); } - gcs::NodeInfoAccessor::ResourceMap GetResources(const NodeID &node_id) { - gcs::NodeInfoAccessor::ResourceMap resource_map; + gcs::NodeResourceInfoAccessor::ResourceMap GetResources(const NodeID &node_id) { + gcs::NodeResourceInfoAccessor::ResourceMap resource_map; std::promise promise; - RAY_CHECK_OK(gcs_client_->Nodes().AsyncGetResources( - node_id, [&resource_map, &promise]( - Status status, - const boost::optional &result) { + RAY_CHECK_OK(gcs_client_->NodeResources().AsyncGetResources( + node_id, + [&resource_map, &promise]( + Status status, + const boost::optional &result) { if (result) { resource_map.insert(result->begin(), result->end()); } @@ -309,11 +310,11 @@ class ServiceBasedGcsClientTest : public ::testing::Test { bool UpdateResources(const NodeID &node_id, const std::string &key) { std::promise promise; - gcs::NodeInfoAccessor::ResourceMap resource_map; + gcs::NodeResourceInfoAccessor::ResourceMap resource_map; auto resource = std::make_shared(); resource->set_resource_capacity(1.0); resource_map[key] = resource; - RAY_CHECK_OK(gcs_client_->Nodes().AsyncUpdateResources( + RAY_CHECK_OK(gcs_client_->NodeResources().AsyncUpdateResources( node_id, resource_map, [&promise](Status status) { promise.set_value(status.ok()); })); return WaitReady(promise.get_future(), timeout_ms_); @@ -322,7 +323,7 @@ class ServiceBasedGcsClientTest : public ::testing::Test { bool DeleteResources(const NodeID &node_id, const std::vector &resource_names) { std::promise promise; - RAY_CHECK_OK(gcs_client_->Nodes().AsyncDeleteResources( + RAY_CHECK_OK(gcs_client_->NodeResources().AsyncDeleteResources( node_id, resource_names, [&promise](Status status) { promise.set_value(status.ok()); })); return WaitReady(promise.get_future(), timeout_ms_); @@ -353,7 +354,7 @@ class ServiceBasedGcsClientTest : public ::testing::Test { std::vector GetAllAvailableResources() { std::promise promise; std::vector resources; - RAY_CHECK_OK(gcs_client_->Nodes().AsyncGetAllAvailableResources( + RAY_CHECK_OK(gcs_client_->NodeResources().AsyncGetAllAvailableResources( [&resources, &promise](Status status, const std::vector &result) { EXPECT_TRUE(!result.empty()); diff --git a/src/ray/gcs/gcs_server/gcs_node_manager.cc b/src/ray/gcs/gcs_server/gcs_node_manager.cc index f2d37f4ea..7549d098f 100644 --- a/src/ray/gcs/gcs_server/gcs_node_manager.cc +++ b/src/ray/gcs/gcs_server/gcs_node_manager.cc @@ -120,92 +120,6 @@ void GcsNodeManager::HandleReportResourceUsage( ++counts_[CountType::REPORT_RESOURCE_USAGE_REQUEST]; } -void GcsNodeManager::HandleGetResources(const rpc::GetResourcesRequest &request, - rpc::GetResourcesReply *reply, - rpc::SendReplyCallback send_reply_callback) { - NodeID node_id = NodeID::FromBinary(request.node_id()); - auto iter = cluster_resources_.find(node_id); - if (iter != cluster_resources_.end()) { - for (auto &resource : iter->second.items()) { - (*reply->mutable_resources())[resource.first] = resource.second; - } - } - GCS_RPC_SEND_REPLY(send_reply_callback, reply, Status::OK()); - ++counts_[CountType::GET_RESOURCES_REQUEST]; -} - -void GcsNodeManager::HandleUpdateResources(const rpc::UpdateResourcesRequest &request, - rpc::UpdateResourcesReply *reply, - rpc::SendReplyCallback send_reply_callback) { - NodeID node_id = NodeID::FromBinary(request.node_id()); - RAY_LOG(DEBUG) << "Updating resources, node id = " << node_id; - auto iter = cluster_resources_.find(node_id); - auto to_be_updated_resources = request.resources(); - if (iter != cluster_resources_.end()) { - for (auto &entry : to_be_updated_resources) { - (*iter->second.mutable_items())[entry.first] = entry.second; - } - auto on_done = [this, node_id, to_be_updated_resources, reply, - send_reply_callback](const Status &status) { - RAY_CHECK_OK(status); - rpc::NodeResourceChange node_resource_change; - node_resource_change.set_node_id(node_id.Binary()); - for (auto &it : to_be_updated_resources) { - (*node_resource_change.mutable_updated_resources())[it.first] = - it.second.resource_capacity(); - } - RAY_CHECK_OK(gcs_pub_sub_->Publish(NODE_RESOURCE_CHANNEL, node_id.Hex(), - node_resource_change.SerializeAsString(), - nullptr)); - - GCS_RPC_SEND_REPLY(send_reply_callback, reply, status); - RAY_LOG(DEBUG) << "Finished updating resources, node id = " << node_id; - }; - - RAY_CHECK_OK( - gcs_table_storage_->NodeResourceTable().Put(node_id, iter->second, on_done)); - } else { - GCS_RPC_SEND_REPLY(send_reply_callback, reply, Status::Invalid("Node is not exist.")); - RAY_LOG(ERROR) << "Failed to update resources as node " << node_id - << " is not registered."; - } - ++counts_[CountType::UPDATE_RESOURCES_REQUEST]; -} - -void GcsNodeManager::HandleDeleteResources(const rpc::DeleteResourcesRequest &request, - rpc::DeleteResourcesReply *reply, - rpc::SendReplyCallback send_reply_callback) { - NodeID node_id = NodeID::FromBinary(request.node_id()); - RAY_LOG(DEBUG) << "Deleting node resources, node id = " << node_id; - auto resource_names = VectorFromProtobuf(request.resource_name_list()); - auto iter = cluster_resources_.find(node_id); - if (iter != cluster_resources_.end()) { - for (auto &resource_name : resource_names) { - RAY_IGNORE_EXPR(iter->second.mutable_items()->erase(resource_name)); - } - auto on_done = [this, node_id, resource_names, reply, - send_reply_callback](const Status &status) { - RAY_CHECK_OK(status); - rpc::NodeResourceChange node_resource_change; - node_resource_change.set_node_id(node_id.Binary()); - for (const auto &resource_name : resource_names) { - node_resource_change.add_deleted_resources(resource_name); - } - RAY_CHECK_OK(gcs_pub_sub_->Publish(NODE_RESOURCE_CHANNEL, node_id.Hex(), - node_resource_change.SerializeAsString(), - nullptr)); - - GCS_RPC_SEND_REPLY(send_reply_callback, reply, status); - }; - RAY_CHECK_OK( - gcs_table_storage_->NodeResourceTable().Put(node_id, iter->second, on_done)); - } else { - GCS_RPC_SEND_REPLY(send_reply_callback, reply, Status::OK()); - RAY_LOG(DEBUG) << "Finished deleting node resources, node id = " << node_id; - } - ++counts_[CountType::DELETE_RESOURCES_REQUEST]; -} - void GcsNodeManager::HandleSetInternalConfig(const rpc::SetInternalConfigRequest &request, rpc::SetInternalConfigReply *reply, rpc::SendReplyCallback send_reply_callback) { @@ -234,22 +148,6 @@ void GcsNodeManager::HandleGetInternalConfig(const rpc::GetInternalConfigRequest ++counts_[CountType::GET_INTERNAL_CONFIG_REQUEST]; } -void GcsNodeManager::HandleGetAllAvailableResources( - const rpc::GetAllAvailableResourcesRequest &request, - rpc::GetAllAvailableResourcesReply *reply, - rpc::SendReplyCallback send_reply_callback) { - for (const auto &iter : gcs_resource_manager_->GetClusterResources()) { - rpc::AvailableResources resource; - resource.set_node_id(iter.first.Binary()); - for (const auto &res : iter.second.GetResourceAmountMap()) { - (*resource.mutable_resources_available())[res.first] = res.second.ToDouble(); - } - reply->add_resources_list()->CopyFrom(resource); - } - GCS_RPC_SEND_REPLY(send_reply_callback, reply, Status::OK()); - ++counts_[CountType::GET_ALL_AVAILABLE_RESOURCES_REQUEST]; -} - void GcsNodeManager::HandleGetAllResourceUsage( const rpc::GetAllResourceUsageRequest &request, rpc::GetAllResourceUsageReply *reply, rpc::SendReplyCallback send_reply_callback) { @@ -337,13 +235,12 @@ void GcsNodeManager::AddNode(std::shared_ptr node) { auto iter = alive_nodes_.find(node_id); if (iter == alive_nodes_.end()) { alive_nodes_.emplace(node_id, node); - // Add an empty resources for this node. - RAY_CHECK(cluster_resources_.emplace(node_id, rpc::ResourceMap()).second); // Notify all listeners. for (auto &listener : node_added_listeners_) { listener(node); } + gcs_resource_manager_->OnNodeAdd(node_id); } } @@ -359,7 +256,7 @@ std::shared_ptr GcsNodeManager::RemoveNode( // Remove from alive nodes. alive_nodes_.erase(iter); // Remove from cluster resources. - cluster_resources_.erase(node_id); + gcs_resource_manager_->OnNodeDead(node_id); resources_buffer_.erase(node_id); if (!is_intended) { // Broadcast a warning to all of the drivers indicating that the node @@ -413,12 +310,6 @@ void GcsNodeManager::Initialize(const GcsInitData &gcs_init_data) { sorted_dead_node_list_.sort( [](const std::pair &left, const std::pair &right) { return left.second < right.second; }); - - for (auto &entry : gcs_init_data.ClusterResources()) { - if (alive_nodes_.count(entry.first)) { - cluster_resources_[entry.first] = entry.second; - } - } } void GcsNodeManager::UpdateNodeRealtimeResources( @@ -486,14 +377,8 @@ std::string GcsNodeManager::DebugString() const { << counts_[CountType::GET_ALL_NODE_INFO_REQUEST] << ", ReportResourceUsage request count: " << counts_[CountType::REPORT_RESOURCE_USAGE_REQUEST] - << ", GetHeartbeat request count: " << counts_[CountType::GET_HEARTBEAT_REQUEST] << ", GetAllResourceUsage request count: " << counts_[CountType::GET_ALL_RESOURCE_USAGE_REQUEST] - << ", GetResources request count: " << counts_[CountType::GET_RESOURCES_REQUEST] - << ", UpdateResources request count: " - << counts_[CountType::UPDATE_RESOURCES_REQUEST] - << ", DeleteResources request count: " - << counts_[CountType::DELETE_RESOURCES_REQUEST] << ", SetInternalConfig request count: " << counts_[CountType::SET_INTERNAL_CONFIG_REQUEST] << ", GetInternalConfig request count: " diff --git a/src/ray/gcs/gcs_server/gcs_node_manager.h b/src/ray/gcs/gcs_server/gcs_node_manager.h index 69f3bdebd..3c62af3a7 100644 --- a/src/ray/gcs/gcs_server/gcs_node_manager.h +++ b/src/ray/gcs/gcs_server/gcs_node_manager.h @@ -70,21 +70,6 @@ class GcsNodeManager : public rpc::NodeInfoHandler { rpc::GetAllResourceUsageReply *reply, rpc::SendReplyCallback send_reply_callback) override; - /// Handle get resource rpc request. - void HandleGetResources(const rpc::GetResourcesRequest &request, - rpc::GetResourcesReply *reply, - rpc::SendReplyCallback send_reply_callback) override; - - /// Handle update resource rpc request. - void HandleUpdateResources(const rpc::UpdateResourcesRequest &request, - rpc::UpdateResourcesReply *reply, - rpc::SendReplyCallback send_reply_callback) override; - - /// Handle delete resource rpc request. - void HandleDeleteResources(const rpc::DeleteResourcesRequest &request, - rpc::DeleteResourcesReply *reply, - rpc::SendReplyCallback send_reply_callback) override; - /// Handle set internal config. void HandleSetInternalConfig(const rpc::SetInternalConfigRequest &request, rpc::SetInternalConfigReply *reply, @@ -95,12 +80,6 @@ class GcsNodeManager : public rpc::NodeInfoHandler { rpc::GetInternalConfigReply *reply, rpc::SendReplyCallback send_reply_callback) override; - /// Handle get available resources of all nodes. - void HandleGetAllAvailableResources( - const rpc::GetAllAvailableResourcesRequest &request, - rpc::GetAllAvailableResourcesReply *reply, - rpc::SendReplyCallback send_reply_callback) override; - /// Update resource usage of given node. /// /// \param node_id Node id. @@ -196,8 +175,6 @@ class GcsNodeManager : public rpc::NodeInfoHandler { /// The nodes are sorted according to the timestamp, and the oldest is at the head of /// the list. std::list> sorted_dead_node_list_; - /// Cluster resources. - absl::flat_hash_map cluster_resources_; /// Newest resource usage of all nodes. absl::flat_hash_map node_resource_usages_; /// A buffer containing resource usage received from node managers in the last tick. @@ -223,15 +200,10 @@ class GcsNodeManager : public rpc::NodeInfoHandler { UNREGISTER_NODE_REQUEST = 1, GET_ALL_NODE_INFO_REQUEST = 2, REPORT_RESOURCE_USAGE_REQUEST = 3, - GET_HEARTBEAT_REQUEST = 4, - GET_ALL_RESOURCE_USAGE_REQUEST = 5, - GET_RESOURCES_REQUEST = 6, - UPDATE_RESOURCES_REQUEST = 7, - DELETE_RESOURCES_REQUEST = 8, - SET_INTERNAL_CONFIG_REQUEST = 9, - GET_INTERNAL_CONFIG_REQUEST = 10, - GET_ALL_AVAILABLE_RESOURCES_REQUEST = 11, - CountType_MAX = 12, + GET_ALL_RESOURCE_USAGE_REQUEST = 4, + SET_INTERNAL_CONFIG_REQUEST = 5, + GET_INTERNAL_CONFIG_REQUEST = 6, + CountType_MAX = 7, }; uint64_t counts_[CountType::CountType_MAX] = {0}; }; diff --git a/src/ray/gcs/gcs_server/gcs_resource_manager.cc b/src/ray/gcs/gcs_server/gcs_resource_manager.cc index 0655570e4..a5349757b 100644 --- a/src/ray/gcs/gcs_server/gcs_resource_manager.cc +++ b/src/ray/gcs/gcs_server/gcs_resource_manager.cc @@ -17,35 +17,161 @@ namespace ray { namespace gcs { +GcsResourceManager::GcsResourceManager( + std::shared_ptr gcs_pub_sub, + std::shared_ptr gcs_table_storage) + : gcs_pub_sub_(gcs_pub_sub), gcs_table_storage_(gcs_table_storage) {} + +void GcsResourceManager::HandleGetResources(const rpc::GetResourcesRequest &request, + rpc::GetResourcesReply *reply, + rpc::SendReplyCallback send_reply_callback) { + NodeID node_id = NodeID::FromBinary(request.node_id()); + auto iter = cluster_resources_.find(node_id); + if (iter != cluster_resources_.end()) { + for (auto &resource : iter->second.items()) { + (*reply->mutable_resources())[resource.first] = resource.second; + } + } + GCS_RPC_SEND_REPLY(send_reply_callback, reply, Status::OK()); + ++counts_[CountType::GET_RESOURCES_REQUEST]; +} + +void GcsResourceManager::HandleUpdateResources( + const rpc::UpdateResourcesRequest &request, rpc::UpdateResourcesReply *reply, + rpc::SendReplyCallback send_reply_callback) { + NodeID node_id = NodeID::FromBinary(request.node_id()); + RAY_LOG(DEBUG) << "Updating resources, node id = " << node_id; + auto iter = cluster_resources_.find(node_id); + auto to_be_updated_resources = request.resources(); + if (iter != cluster_resources_.end()) { + for (auto &entry : to_be_updated_resources) { + (*iter->second.mutable_items())[entry.first] = entry.second; + } + auto on_done = [this, node_id, to_be_updated_resources, reply, + send_reply_callback](const Status &status) { + RAY_CHECK_OK(status); + rpc::NodeResourceChange node_resource_change; + node_resource_change.set_node_id(node_id.Binary()); + for (auto &it : to_be_updated_resources) { + (*node_resource_change.mutable_updated_resources())[it.first] = + it.second.resource_capacity(); + } + RAY_CHECK_OK(gcs_pub_sub_->Publish(NODE_RESOURCE_CHANNEL, node_id.Hex(), + node_resource_change.SerializeAsString(), + nullptr)); + + GCS_RPC_SEND_REPLY(send_reply_callback, reply, status); + RAY_LOG(DEBUG) << "Finished updating resources, node id = " << node_id; + }; + + RAY_CHECK_OK( + gcs_table_storage_->NodeResourceTable().Put(node_id, iter->second, on_done)); + } else { + GCS_RPC_SEND_REPLY(send_reply_callback, reply, Status::Invalid("Node is not exist.")); + RAY_LOG(ERROR) << "Failed to update resources as node " << node_id + << " is not registered."; + } + ++counts_[CountType::UPDATE_RESOURCES_REQUEST]; +} + +void GcsResourceManager::HandleDeleteResources( + const rpc::DeleteResourcesRequest &request, rpc::DeleteResourcesReply *reply, + rpc::SendReplyCallback send_reply_callback) { + NodeID node_id = NodeID::FromBinary(request.node_id()); + RAY_LOG(DEBUG) << "Deleting node resources, node id = " << node_id; + auto resource_names = VectorFromProtobuf(request.resource_name_list()); + auto iter = cluster_resources_.find(node_id); + if (iter != cluster_resources_.end()) { + for (auto &resource_name : resource_names) { + RAY_IGNORE_EXPR(iter->second.mutable_items()->erase(resource_name)); + } + auto on_done = [this, node_id, resource_names, reply, + send_reply_callback](const Status &status) { + RAY_CHECK_OK(status); + rpc::NodeResourceChange node_resource_change; + node_resource_change.set_node_id(node_id.Binary()); + for (const auto &resource_name : resource_names) { + node_resource_change.add_deleted_resources(resource_name); + } + RAY_CHECK_OK(gcs_pub_sub_->Publish(NODE_RESOURCE_CHANNEL, node_id.Hex(), + node_resource_change.SerializeAsString(), + nullptr)); + + GCS_RPC_SEND_REPLY(send_reply_callback, reply, status); + }; + RAY_CHECK_OK( + gcs_table_storage_->NodeResourceTable().Put(node_id, iter->second, on_done)); + } else { + GCS_RPC_SEND_REPLY(send_reply_callback, reply, Status::OK()); + RAY_LOG(DEBUG) << "Finished deleting node resources, node id = " << node_id; + } + ++counts_[CountType::DELETE_RESOURCES_REQUEST]; +} + +void GcsResourceManager::HandleGetAllAvailableResources( + const rpc::GetAllAvailableResourcesRequest &request, + rpc::GetAllAvailableResourcesReply *reply, + rpc::SendReplyCallback send_reply_callback) { + for (const auto &iter : cluster_scheduling_resources_) { + rpc::AvailableResources resource; + resource.set_node_id(iter.first.Binary()); + for (const auto &res : iter.second.GetResourceAmountMap()) { + (*resource.mutable_resources_available())[res.first] = res.second.ToDouble(); + } + reply->add_resources_list()->CopyFrom(resource); + } + GCS_RPC_SEND_REPLY(send_reply_callback, reply, Status::OK()); + ++counts_[CountType::GET_ALL_AVAILABLE_RESOURCES_REQUEST]; +} + +void GcsResourceManager::Initialize(const GcsInitData &gcs_init_data) { + const auto &nodes = gcs_init_data.Nodes(); + for (auto &entry : gcs_init_data.ClusterResources()) { + const auto &iter = nodes.find(entry.first); + if (iter->second.state() == rpc::GcsNodeInfo::ALIVE) { + cluster_resources_[entry.first] = entry.second; + } + } +} + const absl::flat_hash_map &GcsResourceManager::GetClusterResources() const { - return cluster_resources_; + return cluster_scheduling_resources_; } void GcsResourceManager::UpdateResources(const NodeID &node_id, const ResourceSet &resources) { - cluster_resources_[node_id] = resources; + cluster_scheduling_resources_[node_id] = resources; } -void GcsResourceManager::RemoveResources(const NodeID &node_id) { +void GcsResourceManager::OnNodeAdd(const NodeID &node_id) { + // Add an empty resources for this node. + cluster_resources_.emplace(node_id, rpc::ResourceMap()); +} + +void GcsResourceManager::OnNodeDead(const NodeID &node_id) { cluster_resources_.erase(node_id); + cluster_scheduling_resources_.erase(node_id); } bool GcsResourceManager::AcquireResources(const NodeID &node_id, const ResourceSet &required_resources) { - auto iter = cluster_resources_.find(node_id); - RAY_CHECK(iter != cluster_resources_.end()) << "Node " << node_id << " not exist."; - if (!required_resources.IsSubset(iter->second)) { - return false; + auto iter = cluster_scheduling_resources_.find(node_id); + if (iter != cluster_scheduling_resources_.end()) { + if (!required_resources.IsSubset(iter->second)) { + return false; + } + iter->second.SubtractResourcesStrict(required_resources); } - iter->second.SubtractResourcesStrict(required_resources); + // If node dead, we will not find the node. This is a normal scenario, so it returns + // true. return true; } bool GcsResourceManager::ReleaseResources(const NodeID &node_id, const ResourceSet &acquired_resources) { - auto iter = cluster_resources_.find(node_id); - if (iter != cluster_resources_.end()) { + auto iter = cluster_scheduling_resources_.find(node_id); + if (iter != cluster_scheduling_resources_.end()) { iter->second.AddResources(acquired_resources); } // If node dead, we will not find the node. This is a normal scenario, so it returns @@ -53,5 +179,18 @@ bool GcsResourceManager::ReleaseResources(const NodeID &node_id, return true; } +std::string GcsResourceManager::DebugString() const { + std::ostringstream stream; + stream << "GcsResourceManager: {GetResources request count: " + << counts_[CountType::GET_RESOURCES_REQUEST] + << ", GetAllAvailableResources request count" + << counts_[CountType::GET_ALL_AVAILABLE_RESOURCES_REQUEST] + << ", UpdateResources request count: " + << counts_[CountType::UPDATE_RESOURCES_REQUEST] + << ", DeleteResources request count: " + << counts_[CountType::DELETE_RESOURCES_REQUEST] << "}"; + return stream.str(); +} + } // namespace gcs } // namespace ray diff --git a/src/ray/gcs/gcs_server/gcs_resource_manager.h b/src/ray/gcs/gcs_server/gcs_resource_manager.h index 0a38d07ac..4a3e0dfca 100644 --- a/src/ray/gcs/gcs_server/gcs_resource_manager.h +++ b/src/ray/gcs/gcs_server/gcs_resource_manager.h @@ -14,33 +14,77 @@ #pragma once #include "absl/container/flat_hash_map.h" +#include "absl/container/flat_hash_set.h" +#include "ray/common/id.h" #include "ray/common/task/scheduling_resources.h" +#include "ray/gcs/accessor.h" +#include "ray/gcs/gcs_server/gcs_init_data.h" +#include "ray/gcs/gcs_server/gcs_resource_manager.h" +#include "ray/gcs/gcs_server/gcs_table_storage.h" +#include "ray/gcs/pubsub/gcs_pub_sub.h" +#include "ray/rpc/client_call.h" +#include "ray/rpc/gcs_server/gcs_rpc_server.h" +#include "src/ray/protobuf/gcs.pb.h" namespace ray { namespace gcs { /// Gcs resource manager interface. -/// It is used for actor and placement group scheduling. -/// Non-thread safe. -class GcsResourceManagerInterface { +/// It is responsible for handing node resource related rpc requests and it is used for +/// actor and placement group scheduling. It obtains the available resources of nodes +/// through heartbeat reporting. Non-thread safe. +class GcsResourceManager : public rpc::NodeResourceInfoHandler { public: - virtual ~GcsResourceManagerInterface() {} + /// Create a GcsResourceManager. + /// + /// \param gcs_pub_sub GCS message publisher. + /// \param gcs_table_storage GCS table external storage accessor. + explicit GcsResourceManager(std::shared_ptr gcs_pub_sub, + std::shared_ptr gcs_table_storage); + + virtual ~GcsResourceManager() {} + + /// Handle get resource rpc request. + void HandleGetResources(const rpc::GetResourcesRequest &request, + rpc::GetResourcesReply *reply, + rpc::SendReplyCallback send_reply_callback) override; + + /// Handle update resource rpc request. + void HandleUpdateResources(const rpc::UpdateResourcesRequest &request, + rpc::UpdateResourcesReply *reply, + rpc::SendReplyCallback send_reply_callback) override; + + /// Handle delete resource rpc request. + void HandleDeleteResources(const rpc::DeleteResourcesRequest &request, + rpc::DeleteResourcesReply *reply, + rpc::SendReplyCallback send_reply_callback) override; + + /// Handle get available resources of all nodes. + void HandleGetAllAvailableResources( + const rpc::GetAllAvailableResourcesRequest &request, + rpc::GetAllAvailableResourcesReply *reply, + rpc::SendReplyCallback send_reply_callback) override; /// Get the resources of all nodes in the cluster. /// /// \return The resources of all nodes in the cluster. - virtual const absl::flat_hash_map &GetClusterResources() const = 0; + const absl::flat_hash_map &GetClusterResources() const; + + /// Handle a node registration. + /// + /// \param node_id The specified node id. + void OnNodeAdd(const NodeID &node_id); + + /// Handle a node death. + /// + /// \param node_id The specified node id. + void OnNodeDead(const NodeID &node_id); /// Update the resources of the specified node. /// /// \param node_id Id of a node. /// \param resources Resources of a node. - virtual void UpdateResources(const NodeID &node_id, const ResourceSet &resources) = 0; - - /// Remove the resources of the specified node. - /// - /// \param node_id Id of a node. - virtual void RemoveResources(const NodeID &node_id) = 0; + void UpdateResources(const NodeID &node_id, const ResourceSet &resources); /// Acquire resources from the specified node. It will deduct directly from the node /// resources. @@ -48,8 +92,7 @@ class GcsResourceManagerInterface { /// \param node_id Id of a node. /// \param required_resources Resources to apply for. /// \return True if acquire resources successfully. False otherwise. - virtual bool AcquireResources(const NodeID &node_id, - const ResourceSet &required_resources) = 0; + bool AcquireResources(const NodeID &node_id, const ResourceSet &required_resources); /// Release the resources of the specified node. It will be added directly to the node /// resources. @@ -57,29 +100,35 @@ class GcsResourceManagerInterface { /// \param node_id Id of a node. /// \param acquired_resources Resources to release. /// \return True if release resources successfully. False otherwise. - virtual bool ReleaseResources(const NodeID &node_id, - const ResourceSet &acquired_resources) = 0; -}; - -/// Gcs resource manager implementation. It obtains the available resources of nodes -/// through heartbeat reporting. Non-thread safe. -class GcsResourceManager : public GcsResourceManagerInterface { - public: - virtual ~GcsResourceManager() = default; - - const absl::flat_hash_map &GetClusterResources() const; - - void UpdateResources(const NodeID &node_id, const ResourceSet &resources); - - void RemoveResources(const NodeID &node_id); - - bool AcquireResources(const NodeID &node_id, const ResourceSet &required_resources); - bool ReleaseResources(const NodeID &node_id, const ResourceSet &acquired_resources); + /// 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); + + std::string DebugString() const; + private: + /// A publisher for publishing gcs messages. + std::shared_ptr gcs_pub_sub_; + /// Storage for GCS tables. + std::shared_ptr gcs_table_storage_; + /// Cluster resources. + absl::flat_hash_map cluster_resources_; /// Map from node id to the resources of the node. - absl::flat_hash_map cluster_resources_; + absl::flat_hash_map cluster_scheduling_resources_; + + /// Debug info. + enum CountType { + GET_RESOURCES_REQUEST = 0, + UPDATE_RESOURCES_REQUEST = 1, + DELETE_RESOURCES_REQUEST = 2, + GET_ALL_AVAILABLE_RESOURCES_REQUEST = 3, + CountType_MAX = 4, + }; + uint64_t counts_[CountType::CountType_MAX] = {0}; }; } // namespace gcs diff --git a/src/ray/gcs/gcs_server/gcs_server.cc b/src/ray/gcs/gcs_server/gcs_server.cc index 489484a94..bf8ca289d 100644 --- a/src/ray/gcs/gcs_server/gcs_server.cc +++ b/src/ray/gcs/gcs_server/gcs_server.cc @@ -68,7 +68,7 @@ void GcsServer::Start() { void GcsServer::DoStart(const GcsInitData &gcs_init_data) { // Init gcs resource manager. - InitGcsResourceManager(); + InitGcsResourceManager(gcs_init_data); // Init gcs node manager. InitGcsNodeManager(gcs_init_data); @@ -160,8 +160,16 @@ void GcsServer::InitGcsHeartbeatManager(const GcsInitData &gcs_init_data) { rpc_server_.RegisterService(*heartbeat_info_service_); } -void GcsServer::InitGcsResourceManager() { - gcs_resource_manager_ = std::make_shared(); +void GcsServer::InitGcsResourceManager(const GcsInitData &gcs_init_data) { + RAY_CHECK(gcs_table_storage_ && gcs_pub_sub_); + gcs_resource_manager_ = + std::make_shared(gcs_pub_sub_, gcs_table_storage_); + // Initialize by gcs tables data. + gcs_resource_manager_->Initialize(gcs_init_data); + // Register service. + node_resource_info_service_.reset( + new rpc::NodeResourceInfoGrpcService(main_service_, *gcs_resource_manager_)); + rpc_server_.RegisterService(*node_resource_info_service_); } void GcsServer::InitGcsJobManager() { @@ -296,7 +304,6 @@ void GcsServer::InstallEventListeners() { // node is removed from the GCS. gcs_placement_group_manager_->OnNodeDead(node_id); gcs_actor_manager_->OnNodeDead(node_id); - gcs_resource_manager_->RemoveResources(node_id); raylet_client_pool_->Disconnect(NodeID::FromBinary(node->node_id())); }); diff --git a/src/ray/gcs/gcs_server/gcs_server.h b/src/ray/gcs/gcs_server/gcs_server.h index 3dde5d06b..a2082539f 100644 --- a/src/ray/gcs/gcs_server/gcs_server.h +++ b/src/ray/gcs/gcs_server/gcs_server.h @@ -84,7 +84,7 @@ class GcsServer { void InitGcsHeartbeatManager(const GcsInitData &gcs_init_data); /// Initialize gcs resource manager. - void InitGcsResourceManager(); + void InitGcsResourceManager(const GcsInitData &gcs_init_data); /// Initialize gcs job manager. void InitGcsJobManager(); @@ -156,6 +156,8 @@ class GcsServer { std::unique_ptr actor_info_service_; /// Node info handler and service std::unique_ptr node_info_service_; + /// Node resource info handler and service + std::unique_ptr node_resource_info_service_; /// Heartbeat info handler and service std::unique_ptr heartbeat_info_service_; /// Object info handler and service 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 67e857ab9..4ddba0627 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 @@ -26,7 +26,7 @@ class GcsActorSchedulerTest : public ::testing::Test { worker_client_ = std::make_shared(); gcs_pub_sub_ = std::make_shared(redis_client_); gcs_table_storage_ = std::make_shared(redis_client_); - gcs_resource_manager_ = std::make_shared(); + 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_); store_client_ = std::make_shared(io_service_); 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 8dab7280f..a904512ac 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 @@ -23,6 +23,7 @@ class GcsNodeManagerTest : public ::testing::Test { public: GcsNodeManagerTest() { gcs_pub_sub_ = std::make_shared(redis_client_); + gcs_resource_manager_ = std::make_shared(nullptr, nullptr); } protected: 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 9ec639c84..15f96a6a8 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 @@ -54,7 +54,7 @@ class GcsObjectManagerTest : public ::testing::Test { public: void SetUp() override { gcs_table_storage_ = std::make_shared(io_service_); - gcs_resource_manager_ = std::make_shared(); + 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_object_manager_ = std::make_shared( 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 dbb4860c5..e74b5fe1b 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 @@ -68,7 +68,7 @@ class GcsPlacementGroupManagerTest : public ::testing::Test { : mock_placement_group_scheduler_(new MockPlacementGroupScheduler()) { gcs_pub_sub_ = std::make_shared(redis_client_); gcs_table_storage_ = std::make_shared(io_service_); - gcs_resource_manager_ = std::make_shared(); + 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_placement_group_manager_.reset( 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 603f8d45c..3bf5923c9 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 @@ -39,7 +39,7 @@ 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(); + 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_table_storage_ = std::make_shared(io_service_); diff --git a/src/ray/gcs/gcs_server/test/gcs_resource_manager_test.cc b/src/ray/gcs/gcs_server/test/gcs_resource_manager_test.cc index 5e3343d3e..35d425826 100644 --- a/src/ray/gcs/gcs_server/test/gcs_resource_manager_test.cc +++ b/src/ray/gcs/gcs_server/test/gcs_resource_manager_test.cc @@ -25,10 +25,10 @@ using ::testing::_; class GcsResourceManagerTest : public ::testing::Test { public: GcsResourceManagerTest() { - gcs_resource_manager_ = std::make_shared(); + gcs_resource_manager_ = std::make_shared(nullptr, nullptr); } - std::shared_ptr gcs_resource_manager_; + std::shared_ptr gcs_resource_manager_; }; TEST_F(GcsResourceManagerTest, TestBasic) { 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 4511bc722..093b7c462 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 @@ -366,29 +366,6 @@ struct GcsServerMocker { bool IsRemoved(const NodeID &node_id) const override { return false; } - Status AsyncGetResources( - const NodeID &node_id, - const gcs::OptionalItemCallback &callback) override { - return Status::NotImplemented(""); - } - - Status AsyncUpdateResources(const NodeID &node_id, const ResourceMap &resources, - const gcs::StatusCallback &callback) override { - return Status::NotImplemented(""); - } - - Status AsyncDeleteResources(const NodeID &node_id, - const std::vector &resource_names, - const gcs::StatusCallback &callback) override { - return Status::NotImplemented(""); - } - - Status AsyncSubscribeToResources( - const gcs::ItemCallback &subscribe, - const gcs::StatusCallback &done) override { - return Status::NotImplemented(""); - } - Status AsyncReportHeartbeat(const std::shared_ptr &data_ptr, const gcs::StatusCallback &callback) override { return Status::NotImplemented(""); diff --git a/src/ray/gcs/redis_accessor.cc b/src/ray/gcs/redis_accessor.cc index 24ba0ac06..bd3fe0604 100644 --- a/src/ray/gcs/redis_accessor.cc +++ b/src/ray/gcs/redis_accessor.cc @@ -421,7 +421,6 @@ Status RedisObjectInfoAccessor::AsyncUnsubscribeToLocations(const ObjectID &obje RedisNodeInfoAccessor::RedisNodeInfoAccessor(RedisGcsClient *client_impl) : client_impl_(client_impl), - resource_sub_executor_(client_impl_->resource_table()), resource_usage_batch_sub_executor_(client_impl->resource_usage_batch_table()) {} Status RedisNodeInfoAccessor::RegisterSelf(const GcsNodeInfo &local_node_info, @@ -549,7 +548,10 @@ Status RedisNodeInfoAccessor::AsyncSubscribeBatchedResourceUsage( done); } -Status RedisNodeInfoAccessor::AsyncGetResources( +RedisNodeResourceInfoAccessor::RedisNodeResourceInfoAccessor(RedisGcsClient *client_impl) + : client_impl_(client_impl), resource_sub_executor_(client_impl_->resource_table()) {} + +Status RedisNodeResourceInfoAccessor::AsyncGetResources( const NodeID &node_id, const OptionalItemCallback &callback) { RAY_CHECK(callback != nullptr); auto on_done = [callback](RedisGcsClient *client, const NodeID &id, @@ -565,9 +567,8 @@ Status RedisNodeInfoAccessor::AsyncGetResources( return resource_table.Lookup(JobID::Nil(), node_id, on_done); } -Status RedisNodeInfoAccessor::AsyncUpdateResources(const NodeID &node_id, - const ResourceMap &resources, - const StatusCallback &callback) { +Status RedisNodeResourceInfoAccessor::AsyncUpdateResources( + const NodeID &node_id, const ResourceMap &resources, const StatusCallback &callback) { Hash::HashCallback on_done = nullptr; if (callback != nullptr) { on_done = [callback](RedisGcsClient *client, const NodeID &node_id, @@ -578,7 +579,7 @@ Status RedisNodeInfoAccessor::AsyncUpdateResources(const NodeID &node_id, return resource_table.Update(JobID::Nil(), node_id, resources, on_done); } -Status RedisNodeInfoAccessor::AsyncDeleteResources( +Status RedisNodeResourceInfoAccessor::AsyncDeleteResources( const NodeID &node_id, const std::vector &resource_names, const StatusCallback &callback) { Hash::HashRemoveCallback on_done = nullptr; @@ -593,7 +594,7 @@ Status RedisNodeInfoAccessor::AsyncDeleteResources( return resource_table.RemoveEntries(JobID::Nil(), node_id, resource_names, on_done); } -Status RedisNodeInfoAccessor::AsyncSubscribeToResources( +Status RedisNodeResourceInfoAccessor::AsyncSubscribeToResources( const ItemCallback &subscribe, const StatusCallback &done) { RAY_CHECK(subscribe != nullptr); auto on_subscribe = [subscribe](const NodeID &id, diff --git a/src/ray/gcs/redis_accessor.h b/src/ray/gcs/redis_accessor.h index 682f06e5c..c8263d0c8 100644 --- a/src/ray/gcs/redis_accessor.h +++ b/src/ray/gcs/redis_accessor.h @@ -326,24 +326,6 @@ class RedisNodeInfoAccessor : public NodeInfoAccessor { bool IsRemoved(const NodeID &node_id) const override; - Status AsyncGetResources(const NodeID &node_id, - const OptionalItemCallback &callback) override; - - Status AsyncGetAllAvailableResources( - const MultiItemCallback &callback) override { - return Status::NotImplemented("AsyncGetAllAvailableResources not implemented"); - } - - Status AsyncUpdateResources(const NodeID &node_id, const ResourceMap &resources, - const StatusCallback &callback) override; - - Status AsyncDeleteResources(const NodeID &node_id, - const std::vector &resource_names, - const StatusCallback &callback) override; - - Status AsyncSubscribeToResources(const ItemCallback &subscribe, - const StatusCallback &done) override; - Status AsyncReportHeartbeat(const std::shared_ptr &data_ptr, const StatusCallback &callback) override; @@ -377,15 +359,48 @@ class RedisNodeInfoAccessor : public NodeInfoAccessor { private: RedisGcsClient *client_impl_{nullptr}; - typedef SubscriptionExecutor - DynamicResourceSubscriptionExecutor; - DynamicResourceSubscriptionExecutor resource_sub_executor_; - typedef SubscriptionExecutor HeartbeatBatchSubscriptionExecutor; HeartbeatBatchSubscriptionExecutor resource_usage_batch_sub_executor_; }; +/// \class RedisNodeResourceInfoAccessor +/// RedisNodeResourceInfoAccessor is an implementation of `NodeResourceInfoAccessor` +/// that uses Redis as the backend storage. +class RedisNodeResourceInfoAccessor : public NodeResourceInfoAccessor { + public: + explicit RedisNodeResourceInfoAccessor(RedisGcsClient *client_impl); + + virtual ~RedisNodeResourceInfoAccessor() {} + + Status AsyncGetResources(const NodeID &node_id, + const OptionalItemCallback &callback) override; + + Status AsyncGetAllAvailableResources( + const MultiItemCallback &callback) override { + return Status::NotImplemented("AsyncGetAllAvailableResources not implemented"); + } + + Status AsyncUpdateResources(const NodeID &node_id, const ResourceMap &resources, + const StatusCallback &callback) override; + + Status AsyncDeleteResources(const NodeID &node_id, + const std::vector &resource_names, + const StatusCallback &callback) override; + + Status AsyncSubscribeToResources(const ItemCallback &subscribe, + const StatusCallback &done) override; + + void AsyncResubscribe(bool is_pubsub_server_restarted) override {} + + private: + RedisGcsClient *client_impl_{nullptr}; + + typedef SubscriptionExecutor + DynamicResourceSubscriptionExecutor; + DynamicResourceSubscriptionExecutor resource_sub_executor_; +}; + /// \class RedisErrorInfoAccessor /// RedisErrorInfoAccessor is an implementation of `ErrorInfoAccessor` /// that uses Redis as the backend storage. diff --git a/src/ray/gcs/redis_gcs_client.cc b/src/ray/gcs/redis_gcs_client.cc index 20ef35b5e..1b2359346 100644 --- a/src/ray/gcs/redis_gcs_client.cc +++ b/src/ray/gcs/redis_gcs_client.cc @@ -71,6 +71,7 @@ Status RedisGcsClient::Connect(boost::asio::io_service &io_service) { job_accessor_.reset(new RedisJobInfoAccessor(this)); object_accessor_.reset(new RedisObjectInfoAccessor(this)); node_accessor_.reset(new RedisNodeInfoAccessor(this)); + node_resource_accessor_.reset(new RedisNodeResourceInfoAccessor(this)); task_accessor_.reset(new RedisTaskInfoAccessor(this)); error_accessor_.reset(new RedisErrorInfoAccessor(this)); stats_accessor_.reset(new RedisStatsInfoAccessor(this)); diff --git a/src/ray/gcs/test/redis_node_info_accessor_test.cc b/src/ray/gcs/test/redis_node_info_accessor_test.cc index 49b31b09f..e4435184e 100644 --- a/src/ray/gcs/test/redis_node_info_accessor_test.cc +++ b/src/ray/gcs/test/redis_node_info_accessor_test.cc @@ -25,7 +25,7 @@ namespace gcs { class NodeDynamicResourceTest : public AccessorTestBase { protected: - typedef NodeInfoAccessor::ResourceMap ResourceMap; + typedef NodeResourceInfoAccessor::ResourceMap ResourceMap; virtual void GenTestData() { for (size_t node_index = 0; node_index < node_number_; ++node_index) { NodeID id = NodeID::FromRandom(); @@ -56,13 +56,14 @@ class NodeDynamicResourceTest : public AccessorTestBaseNodes(); + NodeResourceInfoAccessor &node_resource_accessor = gcs_client_->NodeResources(); for (const auto &node_rs : id_to_resource_map_) { ++pending_count_; const NodeID &id = node_rs.first; // Update - Status status = node_accessor.AsyncUpdateResources( - node_rs.first, node_rs.second, [this, &node_accessor, id](Status status) { + Status status = node_resource_accessor.AsyncUpdateResources( + node_rs.first, node_rs.second, + [this, &node_resource_accessor, id](Status status) { RAY_CHECK_OK(status); auto get_callback = [this, id](Status status, const boost::optional &result) { @@ -73,7 +74,7 @@ TEST_F(NodeDynamicResourceTest, UpdateAndGet) { ASSERT_EQ(it->second.size(), result->size()); }; // Get - status = node_accessor.AsyncGetResources(id, get_callback); + status = node_resource_accessor.AsyncGetResources(id, get_callback); RAY_CHECK_OK(status); }); } @@ -81,15 +82,15 @@ TEST_F(NodeDynamicResourceTest, UpdateAndGet) { } TEST_F(NodeDynamicResourceTest, Delete) { - NodeInfoAccessor &node_accessor = gcs_client_->Nodes(); + NodeResourceInfoAccessor &node_resource_accessor = gcs_client_->NodeResources(); for (const auto &node_rs : id_to_resource_map_) { ++pending_count_; // Update - Status status = node_accessor.AsyncUpdateResources(node_rs.first, node_rs.second, - [this](Status status) { - RAY_CHECK_OK(status); - --pending_count_; - }); + Status status = node_resource_accessor.AsyncUpdateResources( + node_rs.first, node_rs.second, [this](Status status) { + RAY_CHECK_OK(status); + --pending_count_; + }); } WaitPendingDone(wait_pending_timeout_); @@ -97,11 +98,11 @@ TEST_F(NodeDynamicResourceTest, Delete) { ++pending_count_; const NodeID &id = node_rs.first; // Delete - Status status = node_accessor.AsyncDeleteResources( - id, resource_to_delete_, [this, &node_accessor, id](Status status) { + Status status = node_resource_accessor.AsyncDeleteResources( + id, resource_to_delete_, [this, &node_resource_accessor, id](Status status) { RAY_CHECK_OK(status); // Get - status = node_accessor.AsyncGetResources( + status = node_resource_accessor.AsyncGetResources( id, [this, id](Status status, const boost::optional &result) { --pending_count_; RAY_CHECK_OK(status); @@ -115,15 +116,15 @@ TEST_F(NodeDynamicResourceTest, Delete) { } TEST_F(NodeDynamicResourceTest, Subscribe) { - NodeInfoAccessor &node_accessor = gcs_client_->Nodes(); + NodeResourceInfoAccessor &node_resource_accessor = gcs_client_->NodeResources(); for (const auto &node_rs : id_to_resource_map_) { ++pending_count_; // Update - Status status = node_accessor.AsyncUpdateResources(node_rs.first, node_rs.second, - [this](Status status) { - RAY_CHECK_OK(status); - --pending_count_; - }); + Status status = node_resource_accessor.AsyncUpdateResources( + node_rs.first, node_rs.second, [this](Status status) { + RAY_CHECK_OK(status); + --pending_count_; + }); } WaitPendingDone(wait_pending_timeout_); @@ -147,18 +148,18 @@ TEST_F(NodeDynamicResourceTest, Subscribe) { // Subscribe ++pending_count_; - Status status = node_accessor.AsyncSubscribeToResources(subscribe, done); + Status status = node_resource_accessor.AsyncSubscribeToResources(subscribe, done); RAY_CHECK_OK(status); for (const auto &node_rs : id_to_resource_map_) { // Delete ++pending_count_; ++sub_pending_count_; - Status status = node_accessor.AsyncDeleteResources(node_rs.first, resource_to_delete_, - [this](Status status) { - RAY_CHECK_OK(status); - --pending_count_; - }); + Status status = node_resource_accessor.AsyncDeleteResources( + node_rs.first, resource_to_delete_, [this](Status status) { + RAY_CHECK_OK(status); + --pending_count_; + }); RAY_CHECK_OK(status); } diff --git a/src/ray/protobuf/gcs_service.proto b/src/ray/protobuf/gcs_service.proto index 11d847a61..eb730a7cf 100644 --- a/src/ray/protobuf/gcs_service.proto +++ b/src/ray/protobuf/gcs_service.proto @@ -180,6 +180,40 @@ message GetAllResourceUsageReply { ResourceUsageBatchData resource_usage_data = 2; } +message SetInternalConfigRequest { + StoredConfig config = 1; +} + +message SetInternalConfigReply { + GcsStatus status = 1; +} + +message GetInternalConfigRequest { +} + +message GetInternalConfigReply { + GcsStatus status = 1; + StoredConfig config = 2; +} + +// Service for node info access. +service NodeInfoGcsService { + // Register a node to GCS Service. + rpc RegisterNode(RegisterNodeRequest) returns (RegisterNodeReply); + // Unregister a node from GCS Service. + rpc UnregisterNode(UnregisterNodeRequest) returns (UnregisterNodeReply); + // Get information of all nodes from GCS Service. + rpc GetAllNodeInfo(GetAllNodeInfoRequest) returns (GetAllNodeInfoReply); + // Report resource usage of a node to GCS Service. + rpc ReportResourceUsage(ReportResourceUsageRequest) returns (ReportResourceUsageReply); + // Get resource usage of all nodes from GCS Service. + rpc GetAllResourceUsage(GetAllResourceUsageRequest) returns (GetAllResourceUsageReply); + // Set cluster internal config. + rpc SetInternalConfig(SetInternalConfigRequest) returns (SetInternalConfigReply); + // Get cluster internal config. + rpc GetInternalConfig(GetInternalConfigRequest) returns (GetInternalConfigReply); +} + message GetResourcesRequest { bytes node_id = 1; } @@ -207,22 +241,6 @@ message DeleteResourcesReply { GcsStatus status = 1; } -message SetInternalConfigRequest { - StoredConfig config = 1; -} - -message SetInternalConfigReply { - GcsStatus status = 1; -} - -message GetInternalConfigRequest { -} - -message GetInternalConfigReply { - GcsStatus status = 1; - StoredConfig config = 2; -} - message GetAllAvailableResourcesRequest { } @@ -231,28 +249,14 @@ message GetAllAvailableResourcesReply { repeated AvailableResources resources_list = 2; } -// Service for node info access. -service NodeInfoGcsService { - // Register a node to GCS Service. - rpc RegisterNode(RegisterNodeRequest) returns (RegisterNodeReply); - // Unregister a node from GCS Service. - rpc UnregisterNode(UnregisterNodeRequest) returns (UnregisterNodeReply); - // Get information of all nodes from GCS Service. - rpc GetAllNodeInfo(GetAllNodeInfoRequest) returns (GetAllNodeInfoReply); - // Get newest heartbeat of all nodes from GCS Service. - rpc ReportResourceUsage(ReportResourceUsageRequest) returns (ReportResourceUsageReply); - // Get resource usage of all nodes from GCS Service. - rpc GetAllResourceUsage(GetAllResourceUsageRequest) returns (GetAllResourceUsageReply); +// Service for node resource info access. +service NodeResourceInfoGcsService { // Get node's resources from GCS Service. rpc GetResources(GetResourcesRequest) returns (GetResourcesReply); // Update resources of a node in GCS Service. rpc UpdateResources(UpdateResourcesRequest) returns (UpdateResourcesReply); // Delete resources of a node in GCS Service. rpc DeleteResources(DeleteResourcesRequest) returns (DeleteResourcesReply); - // Set cluster internal config. - rpc SetInternalConfig(SetInternalConfigRequest) returns (SetInternalConfigReply); - // Get cluster internal config. - rpc GetInternalConfig(GetInternalConfigRequest) returns (GetInternalConfigReply); // Get available resources of all nodes. rpc GetAllAvailableResources(GetAllAvailableResourcesRequest) returns (GetAllAvailableResourcesReply); diff --git a/src/ray/raylet/node_manager.cc b/src/ray/raylet/node_manager.cc index 1af9bf8d2..277e3c9df 100644 --- a/src/ray/raylet/node_manager.cc +++ b/src/ray/raylet/node_manager.cc @@ -268,7 +268,7 @@ ray::Status NodeManager::RegisterGcs() { id, VectorFromProtobuf(resource_notification.deleted_resources())); } }; - RAY_CHECK_OK(gcs_client_->Nodes().AsyncSubscribeToResources( + RAY_CHECK_OK(gcs_client_->NodeResources().AsyncSubscribeToResources( /*subscribe_callback=*/resources_changed, /*done_callback=*/nullptr)); }; @@ -749,10 +749,11 @@ void NodeManager::NodeAdded(const GcsNodeInfo &node_info) { std::make_pair(node_info.node_manager_address(), node_info.node_manager_port()); // Fetch resource info for the remote node and update cluster resource map. - RAY_CHECK_OK(gcs_client_->Nodes().AsyncGetResources( + RAY_CHECK_OK(gcs_client_->NodeResources().AsyncGetResources( node_id, - [this, node_id](Status status, - const boost::optional &data) { + [this, node_id]( + Status status, + const boost::optional &data) { if (data) { ResourceSet resource_set; for (auto &resource_entry : *data) { @@ -1930,14 +1931,15 @@ void NodeManager::ProcessSetResourceRequest( // Submit to the resource table. This calls the ResourceCreateUpdated or ResourceDeleted // callback, which updates cluster_resource_map_. if (is_deletion) { - RAY_CHECK_OK( - gcs_client_->Nodes().AsyncDeleteResources(node_id, {resource_name}, nullptr)); + RAY_CHECK_OK(gcs_client_->NodeResources().AsyncDeleteResources( + node_id, {resource_name}, nullptr)); } else { std::unordered_map> data_map; auto resource_table_data = std::make_shared(); resource_table_data->set_resource_capacity(capacity); data_map.emplace(resource_name, resource_table_data); - RAY_CHECK_OK(gcs_client_->Nodes().AsyncUpdateResources(node_id, data_map, nullptr)); + RAY_CHECK_OK( + gcs_client_->NodeResources().AsyncUpdateResources(node_id, data_map, nullptr)); } } diff --git a/src/ray/raylet/raylet.cc b/src/ray/raylet/raylet.cc index 26e03b12e..d2ddead62 100644 --- a/src/ray/raylet/raylet.cc +++ b/src/ray/raylet/raylet.cc @@ -136,8 +136,8 @@ ray::Status Raylet::RegisterGcs() { resource->set_resource_capacity(resource_pair.second); resources.emplace(resource_pair.first, resource); } - RAY_CHECK_OK( - gcs_client_->Nodes().AsyncUpdateResources(self_node_id_, resources, nullptr)); + RAY_CHECK_OK(gcs_client_->NodeResources().AsyncUpdateResources(self_node_id_, + resources, nullptr)); RAY_CHECK_OK(node_manager_.RegisterGcs()); }; diff --git a/src/ray/rpc/gcs_server/gcs_rpc_client.h b/src/ray/rpc/gcs_server/gcs_rpc_client.h index 67a2abb15..82857123e 100644 --- a/src/ray/rpc/gcs_server/gcs_rpc_client.h +++ b/src/ray/rpc/gcs_server/gcs_rpc_client.h @@ -97,6 +97,10 @@ class GcsRpcClient { new GrpcClient(address, port, client_call_manager)); node_info_grpc_client_ = std::unique_ptr>( new GrpcClient(address, port, client_call_manager)); + node_resource_info_grpc_client_ = + std::unique_ptr>( + new GrpcClient(address, port, + client_call_manager)); heartbeat_info_grpc_client_ = std::unique_ptr>( new GrpcClient(address, port, client_call_manager)); object_info_grpc_client_ = std::unique_ptr>( @@ -165,17 +169,6 @@ class GcsRpcClient { VOID_GCS_RPC_CLIENT_METHOD(NodeInfoGcsService, GetAllResourceUsage, node_info_grpc_client_, ) - /// Get node's resources from GCS Service. - VOID_GCS_RPC_CLIENT_METHOD(NodeInfoGcsService, GetResources, node_info_grpc_client_, ) - - /// Update resources of a node in GCS Service. - VOID_GCS_RPC_CLIENT_METHOD(NodeInfoGcsService, UpdateResources, - node_info_grpc_client_, ) - - /// Delete resources of a node in GCS Service. - VOID_GCS_RPC_CLIENT_METHOD(NodeInfoGcsService, DeleteResources, - node_info_grpc_client_, ) - /// Set internal config of the cluster in the GCS Service. VOID_GCS_RPC_CLIENT_METHOD(NodeInfoGcsService, SetInternalConfig, node_info_grpc_client_, ) @@ -184,9 +177,21 @@ class GcsRpcClient { VOID_GCS_RPC_CLIENT_METHOD(NodeInfoGcsService, GetInternalConfig, node_info_grpc_client_, ) + /// Get node's resources from GCS Service. + VOID_GCS_RPC_CLIENT_METHOD(NodeResourceInfoGcsService, GetResources, + node_resource_info_grpc_client_, ) + + /// Update resources of a node in GCS Service. + VOID_GCS_RPC_CLIENT_METHOD(NodeResourceInfoGcsService, UpdateResources, + node_resource_info_grpc_client_, ) + + /// Delete resources of a node in GCS Service. + VOID_GCS_RPC_CLIENT_METHOD(NodeResourceInfoGcsService, DeleteResources, + node_resource_info_grpc_client_, ) + /// Get available resources of all nodes from the GCS Service. - VOID_GCS_RPC_CLIENT_METHOD(NodeInfoGcsService, GetAllAvailableResources, - node_info_grpc_client_, ) + VOID_GCS_RPC_CLIENT_METHOD(NodeResourceInfoGcsService, GetAllAvailableResources, + node_resource_info_grpc_client_, ) /// Report heartbeat of a node to GCS Service. VOID_GCS_RPC_CLIENT_METHOD(HeartbeatInfoGcsService, ReportHeartbeat, @@ -275,6 +280,7 @@ class GcsRpcClient { std::unique_ptr> job_info_grpc_client_; std::unique_ptr> actor_info_grpc_client_; std::unique_ptr> node_info_grpc_client_; + std::unique_ptr> node_resource_info_grpc_client_; std::unique_ptr> heartbeat_info_grpc_client_; std::unique_ptr> object_info_grpc_client_; std::unique_ptr> task_info_grpc_client_; diff --git a/src/ray/rpc/gcs_server/gcs_rpc_server.h b/src/ray/rpc/gcs_server/gcs_rpc_server.h index c5c24c56f..248ec9837 100644 --- a/src/ray/rpc/gcs_server/gcs_rpc_server.h +++ b/src/ray/rpc/gcs_server/gcs_rpc_server.h @@ -33,6 +33,9 @@ namespace rpc { #define HEARTBEAT_INFO_SERVICE_RPC_HANDLER(HANDLER) \ RPC_SERVICE_HANDLER(HeartbeatInfoGcsService, HANDLER) +#define NODE_RESOURCE_INFO_SERVICE_RPC_HANDLER(HANDLER) \ + RPC_SERVICE_HANDLER(NodeResourceInfoGcsService, HANDLER) + #define OBJECT_INFO_SERVICE_RPC_HANDLER(HANDLER) \ RPC_SERVICE_HANDLER(ObjectInfoGcsService, HANDLER) @@ -188,18 +191,6 @@ class NodeInfoGcsServiceHandler { GetAllResourceUsageReply *reply, SendReplyCallback send_reply_callback) = 0; - virtual void HandleGetResources(const GetResourcesRequest &request, - GetResourcesReply *reply, - SendReplyCallback send_reply_callback) = 0; - - virtual void HandleUpdateResources(const UpdateResourcesRequest &request, - UpdateResourcesReply *reply, - SendReplyCallback send_reply_callback) = 0; - - virtual void HandleDeleteResources(const DeleteResourcesRequest &request, - DeleteResourcesReply *reply, - SendReplyCallback send_reply_callback) = 0; - virtual void HandleSetInternalConfig(const SetInternalConfigRequest &request, SetInternalConfigReply *reply, SendReplyCallback send_reply_callback) = 0; @@ -207,11 +198,6 @@ class NodeInfoGcsServiceHandler { virtual void HandleGetInternalConfig(const GetInternalConfigRequest &request, GetInternalConfigReply *reply, SendReplyCallback send_reply_callback) = 0; - - virtual void HandleGetAllAvailableResources( - const rpc::GetAllAvailableResourcesRequest &request, - rpc::GetAllAvailableResourcesReply *reply, - rpc::SendReplyCallback send_reply_callback) = 0; }; /// The `GrpcService` for `NodeInfoGcsService`. @@ -235,12 +221,8 @@ class NodeInfoGrpcService : public GrpcService { NODE_INFO_SERVICE_RPC_HANDLER(GetAllNodeInfo); NODE_INFO_SERVICE_RPC_HANDLER(ReportResourceUsage); NODE_INFO_SERVICE_RPC_HANDLER(GetAllResourceUsage); - NODE_INFO_SERVICE_RPC_HANDLER(GetResources); - NODE_INFO_SERVICE_RPC_HANDLER(UpdateResources); - NODE_INFO_SERVICE_RPC_HANDLER(DeleteResources); NODE_INFO_SERVICE_RPC_HANDLER(SetInternalConfig); NODE_INFO_SERVICE_RPC_HANDLER(GetInternalConfig); - NODE_INFO_SERVICE_RPC_HANDLER(GetAllAvailableResources); } private: @@ -250,6 +232,57 @@ class NodeInfoGrpcService : public GrpcService { NodeInfoGcsServiceHandler &service_handler_; }; +class NodeResourceInfoGcsServiceHandler { + public: + virtual ~NodeResourceInfoGcsServiceHandler() = default; + + virtual void HandleGetResources(const GetResourcesRequest &request, + GetResourcesReply *reply, + SendReplyCallback send_reply_callback) = 0; + + virtual void HandleUpdateResources(const UpdateResourcesRequest &request, + UpdateResourcesReply *reply, + SendReplyCallback send_reply_callback) = 0; + + virtual void HandleDeleteResources(const DeleteResourcesRequest &request, + DeleteResourcesReply *reply, + SendReplyCallback send_reply_callback) = 0; + + virtual void HandleGetAllAvailableResources( + const rpc::GetAllAvailableResourcesRequest &request, + rpc::GetAllAvailableResourcesReply *reply, + rpc::SendReplyCallback send_reply_callback) = 0; +}; + +/// The `GrpcService` for `NodeResourceInfoGcsService`. +class NodeResourceInfoGrpcService : public GrpcService { + public: + /// Constructor. + /// + /// \param[in] handler The service handler that actually handle the requests. + explicit NodeResourceInfoGrpcService(boost::asio::io_service &io_service, + NodeResourceInfoGcsServiceHandler &handler) + : GrpcService(io_service), service_handler_(handler){}; + + protected: + grpc::Service &GetGrpcService() override { return service_; } + + void InitServerCallFactories( + const std::unique_ptr &cq, + std::vector> *server_call_factories) override { + NODE_RESOURCE_INFO_SERVICE_RPC_HANDLER(GetResources); + NODE_RESOURCE_INFO_SERVICE_RPC_HANDLER(UpdateResources); + NODE_RESOURCE_INFO_SERVICE_RPC_HANDLER(DeleteResources); + NODE_RESOURCE_INFO_SERVICE_RPC_HANDLER(GetAllAvailableResources); + } + + private: + /// The grpc async service object. + NodeResourceInfoGcsService::AsyncService service_; + /// The service handler that actually handle the requests. + NodeResourceInfoGcsServiceHandler &service_handler_; +}; + class HeartbeatInfoGcsServiceHandler { public: virtual ~HeartbeatInfoGcsServiceHandler() = default; @@ -539,6 +572,7 @@ class PlacementGroupInfoGrpcService : public GrpcService { using JobInfoHandler = JobInfoGcsServiceHandler; using ActorInfoHandler = ActorInfoGcsServiceHandler; using NodeInfoHandler = NodeInfoGcsServiceHandler; +using NodeResourceInfoHandler = NodeResourceInfoGcsServiceHandler; using HeartbeatInfoHandler = HeartbeatInfoGcsServiceHandler; using ObjectInfoHandler = ObjectInfoGcsServiceHandler; using TaskInfoHandler = TaskInfoGcsServiceHandler;