Publish job/worker info with Hex format instead of Binary (#13235)

This commit is contained in:
Tao Wang
2021-01-07 20:31:58 +08:00
committed by GitHub
parent 3669c02821
commit 82c54c67ee
7 changed files with 12 additions and 15 deletions
+1 -1
View File
@@ -692,7 +692,7 @@ class WorkerInfoAccessor {
/// \param done Callback that will be called when subscription is complete.
/// \return Status
virtual Status AsyncSubscribeToWorkerFailures(
const SubscribeCallback<WorkerID, rpc::WorkerTableData> &subscribe,
const ItemCallback<rpc::WorkerTableData> &subscribe,
const StatusCallback &done) = 0;
/// Report a worker failure to GCS asynchronously.
@@ -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<WorkerID, rpc::WorkerTableData> &subscribe,
const StatusCallback &done) {
const ItemCallback<rpc::WorkerTableData> &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);
};
@@ -428,7 +428,7 @@ class ServiceBasedWorkerInfoAccessor : public WorkerInfoAccessor {
virtual ~ServiceBasedWorkerInfoAccessor() = default;
Status AsyncSubscribeToWorkerFailures(
const SubscribeCallback<WorkerID, rpc::WorkerTableData> &subscribe,
const ItemCallback<rpc::WorkerTableData> &subscribe,
const StatusCallback &done) override;
Status AsyncReportWorkerFailure(const std::shared_ptr<rpc::WorkerTableData> &data_ptr,
@@ -516,7 +516,7 @@ class ServiceBasedGcsClientTest : public ::testing::Test {
}
bool SubscribeToWorkerFailures(
const gcs::SubscribeCallback<WorkerID, rpc::WorkerTableData> &subscribe) {
const gcs::ItemCallback<rpc::WorkerTableData> &subscribe) {
std::promise<bool> 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<int> 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<int> 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));
+2 -2
View File
@@ -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;
+1 -1
View File
@@ -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));
}
+1 -1
View File
@@ -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(