[GCS]Move node resource info to gcs resource manager (#12775)

* 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: 灵洵 <fengbin.ffb@antgroup.com>
This commit is contained in:
fangfengbin
2020-12-13 20:37:34 +08:00
committed by GitHub
co-authored by 灵洵
parent ac24d1db30
commit 1e02b28abe
30 changed files with 701 additions and 543 deletions
+66 -45
View File
@@ -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<std::string, std::shared_ptr<rpc::ResourceTableData>>
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<ResourceMap> &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<rpc::AvailableResources> &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<std::string> &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<rpc::NodeResourceChange> &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<SchedulingResources>();
};
/// \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<std::string, std::shared_ptr<rpc::ResourceTableData>>
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<ResourceMap> &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<rpc::AvailableResources> &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<std::string> &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<rpc::NodeResourceChange> &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
+8
View File
@@ -107,6 +107,13 @@ class GcsClient : public std::enable_shared_from_this<GcsClient> {
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<GcsClient> {
std::unique_ptr<JobInfoAccessor> job_accessor_;
std::unique_ptr<ObjectInfoAccessor> object_accessor_;
std::unique_ptr<NodeInfoAccessor> node_accessor_;
std::unique_ptr<NodeResourceInfoAccessor> node_resource_accessor_;
std::unique_ptr<TaskInfoAccessor> task_accessor_;
std::unique_ptr<ErrorInfoAccessor> error_accessor_;
std::unique_ptr<StatsInfoAccessor> stats_accessor_;
@@ -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<ray::gcs::NodeInfoAccessor::ResourceMap> &result) {
const boost::optional<ray::gcs::NodeResourceInfoAccessor::ResourceMap>
&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<std::string> GlobalStateAccessor::GetAllAvailableResources() {
std::vector<std::string> available_resources;
std::promise<bool> promise;
RAY_CHECK_OK(gcs_client_->Nodes().AsyncGetAllAvailableResources(
RAY_CHECK_OK(gcs_client_->NodeResources().AsyncGetAllAvailableResources(
TransformForMultiItemCallback<rpc::AvailableResources>(available_resources,
promise)));
promise.get_future().get();
+119 -108
View File
@@ -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<ResourceMap> &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<rpc::ResourceTableData>(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<rpc::AvailableResources> &callback) {
rpc::GetAllAvailableResourcesRequest request;
client_impl_->GetGcsRpcClient().GetAllAvailableResources(
request,
[callback](const Status &status, const rpc::GetAllAvailableResourcesReply &reply) {
std::vector<rpc::AvailableResources> 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<std::string> &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<rpc::NodeResourceChange> &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<rpc::HeartbeatTableData> &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<ResourceMap> &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<rpc::ResourceTableData>(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<rpc::AvailableResources> &callback) {
rpc::GetAllAvailableResourcesRequest request;
client_impl_->GetGcsRpcClient().GetAllAvailableResources(
request,
[callback](const Status &status, const rpc::GetAllAvailableResourcesReply &reply) {
std::vector<rpc::AvailableResources> 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<std::string> &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<rpc::NodeResourceChange> &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) {}
+37 -19
View File
@@ -163,22 +163,6 @@ class ServiceBasedNodeInfoAccessor : public NodeInfoAccessor {
bool IsRemoved(const NodeID &node_id) const override;
Status AsyncGetResources(const NodeID &node_id,
const OptionalItemCallback<ResourceMap> &callback) override;
Status AsyncGetAllAvailableResources(
const MultiItemCallback<rpc::AvailableResources> &callback) override;
Status AsyncUpdateResources(const NodeID &node_id, const ResourceMap &resources,
const StatusCallback &callback) override;
Status AsyncDeleteResources(const NodeID &node_id,
const std::vector<std::string> &resource_names,
const StatusCallback &callback) override;
Status AsyncSubscribeToResources(const ItemCallback<rpc::NodeResourceChange> &subscribe,
const StatusCallback &done) override;
Status AsyncReportHeartbeat(const std::shared_ptr<rpc::HeartbeatTableData> &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<NodeID> 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<NodeID> 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<ResourceMap> &callback) override;
Status AsyncGetAllAvailableResources(
const MultiItemCallback<rpc::AvailableResources> &callback) override;
Status AsyncUpdateResources(const NodeID &node_id, const ResourceMap &resources,
const StatusCallback &callback) override;
Status AsyncDeleteResources(const NodeID &node_id,
const std::vector<std::string> &resource_names,
const StatusCallback &callback) override;
Status AsyncSubscribeToResources(const ItemCallback<rpc::NodeResourceChange> &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<NodeID> sequencer_;
};
/// \class ServiceBasedTaskInfoAccessor
/// ServiceBasedTaskInfoAccessor is an implementation of `TaskInfoAccessor`
/// that uses GCS service as the backend.
@@ -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));
@@ -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<double>(index + 1) + 0.1);
resources[std::to_string(index)] =
std::make_shared<rpc::ResourceTableData>(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();
@@ -286,18 +286,19 @@ class ServiceBasedGcsClientTest : public ::testing::Test {
bool SubscribeToResources(const gcs::ItemCallback<rpc::NodeResourceChange> &subscribe) {
std::promise<bool> 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<bool> promise;
RAY_CHECK_OK(gcs_client_->Nodes().AsyncGetResources(
node_id, [&resource_map, &promise](
Status status,
const boost::optional<gcs::NodeInfoAccessor::ResourceMap> &result) {
RAY_CHECK_OK(gcs_client_->NodeResources().AsyncGetResources(
node_id,
[&resource_map, &promise](
Status status,
const boost::optional<gcs::NodeResourceInfoAccessor::ResourceMap> &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<bool> promise;
gcs::NodeInfoAccessor::ResourceMap resource_map;
gcs::NodeResourceInfoAccessor::ResourceMap resource_map;
auto resource = std::make_shared<rpc::ResourceTableData>();
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<std::string> &resource_names) {
std::promise<bool> 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<rpc::AvailableResources> GetAllAvailableResources() {
std::promise<bool> promise;
std::vector<rpc::AvailableResources> resources;
RAY_CHECK_OK(gcs_client_->Nodes().AsyncGetAllAvailableResources(
RAY_CHECK_OK(gcs_client_->NodeResources().AsyncGetAllAvailableResources(
[&resources, &promise](Status status,
const std::vector<rpc::AvailableResources> &result) {
EXPECT_TRUE(!result.empty());
+2 -117
View File
@@ -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<rpc::GcsNodeInfo> 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<rpc::GcsNodeInfo> 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<NodeID, int64_t> &left,
const std::pair<NodeID, int64_t> &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: "
+4 -32
View File
@@ -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<std::pair<NodeID, int64_t>> sorted_dead_node_list_;
/// Cluster resources.
absl::flat_hash_map<NodeID, rpc::ResourceMap> cluster_resources_;
/// Newest resource usage of all nodes.
absl::flat_hash_map<NodeID, rpc::ResourcesData> 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};
};
+149 -10
View File
@@ -17,35 +17,161 @@
namespace ray {
namespace gcs {
GcsResourceManager::GcsResourceManager(
std::shared_ptr<gcs::GcsPubSub> gcs_pub_sub,
std::shared_ptr<gcs::GcsTableStorage> 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<NodeID, ResourceSet> &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
+81 -32
View File
@@ -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::GcsPubSub> gcs_pub_sub,
std::shared_ptr<gcs::GcsTableStorage> 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<NodeID, ResourceSet> &GetClusterResources() const = 0;
const absl::flat_hash_map<NodeID, ResourceSet> &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<NodeID, ResourceSet> &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::GcsPubSub> gcs_pub_sub_;
/// Storage for GCS tables.
std::shared_ptr<gcs::GcsTableStorage> gcs_table_storage_;
/// Cluster resources.
absl::flat_hash_map<NodeID, rpc::ResourceMap> cluster_resources_;
/// Map from node id to the resources of the node.
absl::flat_hash_map<NodeID, ResourceSet> cluster_resources_;
absl::flat_hash_map<NodeID, ResourceSet> 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
+11 -4
View File
@@ -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<GcsResourceManager>();
void GcsServer::InitGcsResourceManager(const GcsInitData &gcs_init_data) {
RAY_CHECK(gcs_table_storage_ && gcs_pub_sub_);
gcs_resource_manager_ =
std::make_shared<GcsResourceManager>(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()));
});
+3 -1
View File
@@ -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<rpc::ActorInfoGrpcService> actor_info_service_;
/// Node info handler and service
std::unique_ptr<rpc::NodeInfoGrpcService> node_info_service_;
/// Node resource info handler and service
std::unique_ptr<rpc::NodeResourceInfoGrpcService> node_resource_info_service_;
/// Heartbeat info handler and service
std::unique_ptr<rpc::HeartbeatInfoGrpcService> heartbeat_info_service_;
/// Object info handler and service
@@ -26,7 +26,7 @@ class GcsActorSchedulerTest : public ::testing::Test {
worker_client_ = std::make_shared<GcsServerMocker::MockWorkerClient>();
gcs_pub_sub_ = std::make_shared<GcsServerMocker::MockGcsPubSub>(redis_client_);
gcs_table_storage_ = std::make_shared<gcs::RedisGcsTableStorage>(redis_client_);
gcs_resource_manager_ = std::make_shared<gcs::GcsResourceManager>();
gcs_resource_manager_ = std::make_shared<gcs::GcsResourceManager>(nullptr, nullptr);
gcs_node_manager_ = std::make_shared<gcs::GcsNodeManager>(
io_service_, gcs_pub_sub_, gcs_table_storage_, gcs_resource_manager_);
store_client_ = std::make_shared<gcs::InMemoryStoreClient>(io_service_);
@@ -23,6 +23,7 @@ class GcsNodeManagerTest : public ::testing::Test {
public:
GcsNodeManagerTest() {
gcs_pub_sub_ = std::make_shared<GcsServerMocker::MockGcsPubSub>(redis_client_);
gcs_resource_manager_ = std::make_shared<gcs::GcsResourceManager>(nullptr, nullptr);
}
protected:
@@ -54,7 +54,7 @@ class GcsObjectManagerTest : public ::testing::Test {
public:
void SetUp() override {
gcs_table_storage_ = std::make_shared<gcs::InMemoryGcsTableStorage>(io_service_);
gcs_resource_manager_ = std::make_shared<gcs::GcsResourceManager>();
gcs_resource_manager_ = std::make_shared<gcs::GcsResourceManager>(nullptr, nullptr);
gcs_node_manager_ = std::make_shared<gcs::GcsNodeManager>(
io_service_, gcs_pub_sub_, gcs_table_storage_, gcs_resource_manager_);
gcs_object_manager_ = std::make_shared<MockedGcsObjectManager>(
@@ -68,7 +68,7 @@ class GcsPlacementGroupManagerTest : public ::testing::Test {
: mock_placement_group_scheduler_(new MockPlacementGroupScheduler()) {
gcs_pub_sub_ = std::make_shared<GcsServerMocker::MockGcsPubSub>(redis_client_);
gcs_table_storage_ = std::make_shared<gcs::InMemoryGcsTableStorage>(io_service_);
gcs_resource_manager_ = std::make_shared<gcs::GcsResourceManager>();
gcs_resource_manager_ = std::make_shared<gcs::GcsResourceManager>(nullptr, nullptr);
gcs_node_manager_ = std::make_shared<gcs::GcsNodeManager>(
io_service_, gcs_pub_sub_, gcs_table_storage_, gcs_resource_manager_);
gcs_placement_group_manager_.reset(
@@ -39,7 +39,7 @@ class GcsPlacementGroupSchedulerTest : public ::testing::Test {
}
gcs_table_storage_ = std::make_shared<gcs::InMemoryGcsTableStorage>(io_service_);
gcs_pub_sub_ = std::make_shared<GcsServerMocker::MockGcsPubSub>(redis_client_);
gcs_resource_manager_ = std::make_shared<gcs::GcsResourceManager>();
gcs_resource_manager_ = std::make_shared<gcs::GcsResourceManager>(nullptr, nullptr);
gcs_node_manager_ = std::make_shared<gcs::GcsNodeManager>(
io_service_, gcs_pub_sub_, gcs_table_storage_, gcs_resource_manager_);
gcs_table_storage_ = std::make_shared<gcs::InMemoryGcsTableStorage>(io_service_);
@@ -25,10 +25,10 @@ using ::testing::_;
class GcsResourceManagerTest : public ::testing::Test {
public:
GcsResourceManagerTest() {
gcs_resource_manager_ = std::make_shared<gcs::GcsResourceManager>();
gcs_resource_manager_ = std::make_shared<gcs::GcsResourceManager>(nullptr, nullptr);
}
std::shared_ptr<gcs::GcsResourceManagerInterface> gcs_resource_manager_;
std::shared_ptr<gcs::GcsResourceManager> gcs_resource_manager_;
};
TEST_F(GcsResourceManagerTest, TestBasic) {
@@ -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<ResourceMap> &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<std::string> &resource_names,
const gcs::StatusCallback &callback) override {
return Status::NotImplemented("");
}
Status AsyncSubscribeToResources(
const gcs::ItemCallback<rpc::NodeResourceChange> &subscribe,
const gcs::StatusCallback &done) override {
return Status::NotImplemented("");
}
Status AsyncReportHeartbeat(const std::shared_ptr<rpc::HeartbeatTableData> &data_ptr,
const gcs::StatusCallback &callback) override {
return Status::NotImplemented("");
+8 -7
View File
@@ -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<ResourceMap> &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<NodeID, ResourceTableData>::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<std::string> &resource_names,
const StatusCallback &callback) {
Hash<NodeID, ResourceTableData>::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<rpc::NodeResourceChange> &subscribe, const StatusCallback &done) {
RAY_CHECK(subscribe != nullptr);
auto on_subscribe = [subscribe](const NodeID &id,
+37 -22
View File
@@ -326,24 +326,6 @@ class RedisNodeInfoAccessor : public NodeInfoAccessor {
bool IsRemoved(const NodeID &node_id) const override;
Status AsyncGetResources(const NodeID &node_id,
const OptionalItemCallback<ResourceMap> &callback) override;
Status AsyncGetAllAvailableResources(
const MultiItemCallback<rpc::AvailableResources> &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<std::string> &resource_names,
const StatusCallback &callback) override;
Status AsyncSubscribeToResources(const ItemCallback<rpc::NodeResourceChange> &subscribe,
const StatusCallback &done) override;
Status AsyncReportHeartbeat(const std::shared_ptr<HeartbeatTableData> &data_ptr,
const StatusCallback &callback) override;
@@ -377,15 +359,48 @@ class RedisNodeInfoAccessor : public NodeInfoAccessor {
private:
RedisGcsClient *client_impl_{nullptr};
typedef SubscriptionExecutor<NodeID, ResourceChangeNotification, DynamicResourceTable>
DynamicResourceSubscriptionExecutor;
DynamicResourceSubscriptionExecutor resource_sub_executor_;
typedef SubscriptionExecutor<NodeID, ResourceUsageBatchData, ResourceUsageBatchTable>
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<ResourceMap> &callback) override;
Status AsyncGetAllAvailableResources(
const MultiItemCallback<rpc::AvailableResources> &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<std::string> &resource_names,
const StatusCallback &callback) override;
Status AsyncSubscribeToResources(const ItemCallback<rpc::NodeResourceChange> &subscribe,
const StatusCallback &done) override;
void AsyncResubscribe(bool is_pubsub_server_restarted) override {}
private:
RedisGcsClient *client_impl_{nullptr};
typedef SubscriptionExecutor<NodeID, ResourceChangeNotification, DynamicResourceTable>
DynamicResourceSubscriptionExecutor;
DynamicResourceSubscriptionExecutor resource_sub_executor_;
};
/// \class RedisErrorInfoAccessor
/// RedisErrorInfoAccessor is an implementation of `ErrorInfoAccessor`
/// that uses Redis as the backend storage.
+1
View File
@@ -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));
@@ -25,7 +25,7 @@ namespace gcs {
class NodeDynamicResourceTest : public AccessorTestBase<NodeID, ResourceTableData> {
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 AccessorTestBase<NodeID, ResourceTableDat
};
TEST_F(NodeDynamicResourceTest, UpdateAndGet) {
NodeInfoAccessor &node_accessor = gcs_client_->Nodes();
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<ResourceMap> &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<ResourceMap> &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);
}
+36 -32
View File
@@ -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);
+9 -7
View File
@@ -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<gcs::NodeInfoAccessor::ResourceMap> &data) {
[this, node_id](
Status status,
const boost::optional<gcs::NodeResourceInfoAccessor::ResourceMap> &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<std::string, std::shared_ptr<gcs::ResourceTableData>> data_map;
auto resource_table_data = std::make_shared<gcs::ResourceTableData>();
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));
}
}
+2 -2
View File
@@ -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());
};
+19 -13
View File
@@ -97,6 +97,10 @@ class GcsRpcClient {
new GrpcClient<ActorInfoGcsService>(address, port, client_call_manager));
node_info_grpc_client_ = std::unique_ptr<GrpcClient<NodeInfoGcsService>>(
new GrpcClient<NodeInfoGcsService>(address, port, client_call_manager));
node_resource_info_grpc_client_ =
std::unique_ptr<GrpcClient<NodeResourceInfoGcsService>>(
new GrpcClient<NodeResourceInfoGcsService>(address, port,
client_call_manager));
heartbeat_info_grpc_client_ = std::unique_ptr<GrpcClient<HeartbeatInfoGcsService>>(
new GrpcClient<HeartbeatInfoGcsService>(address, port, client_call_manager));
object_info_grpc_client_ = std::unique_ptr<GrpcClient<ObjectInfoGcsService>>(
@@ -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<GrpcClient<JobInfoGcsService>> job_info_grpc_client_;
std::unique_ptr<GrpcClient<ActorInfoGcsService>> actor_info_grpc_client_;
std::unique_ptr<GrpcClient<NodeInfoGcsService>> node_info_grpc_client_;
std::unique_ptr<GrpcClient<NodeResourceInfoGcsService>> node_resource_info_grpc_client_;
std::unique_ptr<GrpcClient<HeartbeatInfoGcsService>> heartbeat_info_grpc_client_;
std::unique_ptr<GrpcClient<ObjectInfoGcsService>> object_info_grpc_client_;
std::unique_ptr<GrpcClient<TaskInfoGcsService>> task_info_grpc_client_;
+55 -21
View File
@@ -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<grpc::ServerCompletionQueue> &cq,
std::vector<std::unique_ptr<ServerCallFactory>> *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;