From 82c54c67eeb6b7cf063fa222d6a9f54726e90e18 Mon Sep 17 00:00:00 2001 From: Tao Wang Date: Thu, 7 Jan 2021 20:31:58 +0800 Subject: [PATCH] Publish job/worker info with Hex format instead of Binary (#13235) --- src/ray/gcs/accessor.h | 2 +- src/ray/gcs/gcs_client/service_based_accessor.cc | 7 +++---- src/ray/gcs/gcs_client/service_based_accessor.h | 2 +- .../gcs/gcs_client/test/service_based_gcs_client_test.cc | 8 +++----- src/ray/gcs/gcs_server/gcs_job_manager.cc | 4 ++-- src/ray/gcs/gcs_server/gcs_worker_manager.cc | 2 +- src/ray/raylet/node_manager.cc | 2 +- 7 files changed, 12 insertions(+), 15 deletions(-) diff --git a/src/ray/gcs/accessor.h b/src/ray/gcs/accessor.h index 4adbc0ddd..eac5ae5e7 100644 --- a/src/ray/gcs/accessor.h +++ b/src/ray/gcs/accessor.h @@ -692,7 +692,7 @@ class WorkerInfoAccessor { /// \param done Callback that will be called when subscription is complete. /// \return Status virtual Status AsyncSubscribeToWorkerFailures( - const SubscribeCallback &subscribe, + const ItemCallback &subscribe, const StatusCallback &done) = 0; /// Report a worker failure to GCS asynchronously. diff --git a/src/ray/gcs/gcs_client/service_based_accessor.cc b/src/ray/gcs/gcs_client/service_based_accessor.cc index 8a3fd8424..c276a7c5f 100644 --- a/src/ray/gcs/gcs_client/service_based_accessor.cc +++ b/src/ray/gcs/gcs_client/service_based_accessor.cc @@ -82,7 +82,7 @@ Status ServiceBasedJobInfoAccessor::AsyncSubscribeAll( auto on_subscribe = [subscribe](const std::string &id, const std::string &data) { JobTableData job_data; job_data.ParseFromString(data); - subscribe(JobID::FromBinary(id), job_data); + subscribe(JobID::FromHex(id), job_data); }; return client_impl_->GetGcsPubSub().SubscribeAll(JOB_CHANNEL, on_subscribe, done); }; @@ -1377,14 +1377,13 @@ ServiceBasedWorkerInfoAccessor::ServiceBasedWorkerInfoAccessor( : client_impl_(client_impl) {} Status ServiceBasedWorkerInfoAccessor::AsyncSubscribeToWorkerFailures( - const SubscribeCallback &subscribe, - const StatusCallback &done) { + const ItemCallback &subscribe, const StatusCallback &done) { RAY_CHECK(subscribe != nullptr); subscribe_operation_ = [this, subscribe](const StatusCallback &done) { auto on_subscribe = [subscribe](const std::string &id, const std::string &data) { rpc::WorkerTableData worker_failure_data; worker_failure_data.ParseFromString(data); - subscribe(WorkerID::FromBinary(id), worker_failure_data); + subscribe(worker_failure_data); }; return client_impl_->GetGcsPubSub().SubscribeAll(WORKER_CHANNEL, on_subscribe, done); }; diff --git a/src/ray/gcs/gcs_client/service_based_accessor.h b/src/ray/gcs/gcs_client/service_based_accessor.h index 47c763c66..5133a9512 100644 --- a/src/ray/gcs/gcs_client/service_based_accessor.h +++ b/src/ray/gcs/gcs_client/service_based_accessor.h @@ -428,7 +428,7 @@ class ServiceBasedWorkerInfoAccessor : public WorkerInfoAccessor { virtual ~ServiceBasedWorkerInfoAccessor() = default; Status AsyncSubscribeToWorkerFailures( - const SubscribeCallback &subscribe, + const ItemCallback &subscribe, const StatusCallback &done) override; Status AsyncReportWorkerFailure(const std::shared_ptr &data_ptr, 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 521b280d8..4a4868c5d 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 @@ -516,7 +516,7 @@ class ServiceBasedGcsClientTest : public ::testing::Test { } bool SubscribeToWorkerFailures( - const gcs::SubscribeCallback &subscribe) { + const gcs::ItemCallback &subscribe) { std::promise promise; RAY_CHECK_OK(gcs_client_->Workers().AsyncSubscribeToWorkerFailures( subscribe, [&promise](Status status) { promise.set_value(status.ok()); })); @@ -960,8 +960,7 @@ TEST_F(ServiceBasedGcsClientTest, TestStats) { TEST_F(ServiceBasedGcsClientTest, TestWorkerInfo) { // Subscribe to all unexpected failure of workers from GCS. std::atomic worker_failure_count(0); - auto on_subscribe = [&worker_failure_count](const WorkerID &worker_id, - const rpc::WorkerTableData &result) { + auto on_subscribe = [&worker_failure_count](const rpc::WorkerTableData &result) { ++worker_failure_count; }; ASSERT_TRUE(SubscribeToWorkerFailures(on_subscribe)); @@ -1217,8 +1216,7 @@ TEST_F(ServiceBasedGcsClientTest, TestTaskTableResubscribe) { TEST_F(ServiceBasedGcsClientTest, TestWorkerTableResubscribe) { // Subscribe to all unexpected failure of workers from GCS. std::atomic worker_failure_count(0); - auto on_subscribe = [&worker_failure_count](const WorkerID &worker_id, - const rpc::WorkerTableData &result) { + auto on_subscribe = [&worker_failure_count](const rpc::WorkerTableData &result) { ++worker_failure_count; }; ASSERT_TRUE(SubscribeToWorkerFailures(on_subscribe)); diff --git a/src/ray/gcs/gcs_server/gcs_job_manager.cc b/src/ray/gcs/gcs_server/gcs_job_manager.cc index 9b6cbf074..a991c15c7 100644 --- a/src/ray/gcs/gcs_server/gcs_job_manager.cc +++ b/src/ray/gcs/gcs_server/gcs_job_manager.cc @@ -31,7 +31,7 @@ void GcsJobManager::HandleAddJob(const rpc::AddJobRequest &request, RAY_LOG(ERROR) << "Failed to add job, job id = " << job_id << ", driver pid = " << request.data().driver_pid(); } else { - RAY_CHECK_OK(gcs_pub_sub_->Publish(JOB_CHANNEL, job_id.Binary(), + RAY_CHECK_OK(gcs_pub_sub_->Publish(JOB_CHANNEL, job_id.Hex(), request.data().SerializeAsString(), nullptr)); RAY_LOG(INFO) << "Finished adding job, job id = " << job_id << ", driver pid = " << request.data().driver_pid(); @@ -57,7 +57,7 @@ void GcsJobManager::HandleMarkJobFinished(const rpc::MarkJobFinishedRequest &req if (!status.ok()) { RAY_LOG(ERROR) << "Failed to mark job state, job id = " << job_id; } else { - RAY_CHECK_OK(gcs_pub_sub_->Publish(JOB_CHANNEL, job_id.Binary(), + RAY_CHECK_OK(gcs_pub_sub_->Publish(JOB_CHANNEL, job_id.Hex(), job_table_data->SerializeAsString(), nullptr)); ClearJobInfos(job_id); RAY_LOG(INFO) << "Finished marking job state, job id = " << job_id; diff --git a/src/ray/gcs/gcs_server/gcs_worker_manager.cc b/src/ray/gcs/gcs_server/gcs_worker_manager.cc index 55768bdfe..7f626d636 100644 --- a/src/ray/gcs/gcs_server/gcs_worker_manager.cc +++ b/src/ray/gcs/gcs_server/gcs_worker_manager.cc @@ -52,7 +52,7 @@ void GcsWorkerManager::HandleReportWorkerFailure( << ", address = " << worker_address.ip_address(); } else { stats::UnintentionalWorkerFailures.Record(1); - RAY_CHECK_OK(gcs_pub_sub_->Publish(WORKER_CHANNEL, worker_id.Binary(), + RAY_CHECK_OK(gcs_pub_sub_->Publish(WORKER_CHANNEL, worker_id.Hex(), worker_failure_data->SerializeAsString(), nullptr)); } diff --git a/src/ray/raylet/node_manager.cc b/src/ray/raylet/node_manager.cc index b1657de49..356b5fa34 100644 --- a/src/ray/raylet/node_manager.cc +++ b/src/ray/raylet/node_manager.cc @@ -296,7 +296,7 @@ ray::Status NodeManager::RegisterGcs() { // node failure. These workers can be identified by comparing the raylet_id // in their rpc::Address to the ID of a failed raylet. const auto &worker_failure_handler = - [this](const WorkerID &id, const rpc::WorkerTableData &worker_failure_data) { + [this](const rpc::WorkerTableData &worker_failure_data) { HandleUnexpectedWorkerFailure(worker_failure_data.worker_address()); }; RAY_CHECK_OK(gcs_client_->Workers().AsyncSubscribeToWorkerFailures(