From 3f90ec5963ad0df6d00bdc15a6e228947130aa0b Mon Sep 17 00:00:00 2001 From: fangfengbin <869218239a@zju.edu.cn> Date: Sun, 20 Sep 2020 12:35:45 +0800 Subject: [PATCH] [GCS]Fix actor idempotent bug (#10856) --- .../test/direct_actor_transport_test.cc | 14 ++++++++++++++ .../transport/direct_actor_transport.cc | 7 +++++++ 2 files changed, 21 insertions(+) diff --git a/src/ray/core_worker/test/direct_actor_transport_test.cc b/src/ray/core_worker/test/direct_actor_transport_test.cc index d92fc95ee..35101dbbf 100644 --- a/src/ray/core_worker/test/direct_actor_transport_test.cc +++ b/src/ray/core_worker/test/direct_actor_transport_test.cc @@ -140,6 +140,12 @@ TEST_F(DirectActorSubmitterTest, TestSubmitTask) { ASSERT_TRUE(worker_client_->ReplyPushTask()); } ASSERT_THAT(worker_client_->received_seq_nos, ElementsAre(0, 1)); + + // Connect to the actor again. + // Because the IP and port of address are not modified, it will skip directly and will + // not reset `received_seq_nos`. + submitter_.ConnectActor(actor_id, addr, 0); + ASSERT_THAT(worker_client_->received_seq_nos, ElementsAre(0, 1)); } TEST_F(DirectActorSubmitterTest, TestDependencies) { @@ -251,6 +257,7 @@ TEST_F(DirectActorSubmitterTest, TestActorRestartNoRetry) { ActorID actor_id = ActorID::Of(JobID::FromInt(0), TaskID::Nil(), 0); submitter_.AddActorQueueIfNotExists(actor_id); gcs::ActorTableData actor_data; + addr.set_port(0); submitter_.ConnectActor(actor_id, addr, 0); ASSERT_EQ(worker_client_->callbacks.size(), 0); @@ -277,6 +284,7 @@ TEST_F(DirectActorSubmitterTest, TestActorRestartNoRetry) { ASSERT_TRUE(worker_client_->ReplyPushTask(Status::IOError(""))); // Actor gets restarted. + addr.set_port(1); submitter_.ConnectActor(actor_id, addr, 1); ASSERT_TRUE(submitter_.SubmitTask(task4).ok()); ASSERT_TRUE(worker_client_->ReplyPushTask(Status::OK())); @@ -292,6 +300,7 @@ TEST_F(DirectActorSubmitterTest, TestActorRestartRetry) { ActorID actor_id = ActorID::Of(JobID::FromInt(0), TaskID::Nil(), 0); submitter_.AddActorQueueIfNotExists(actor_id); gcs::ActorTableData actor_data; + addr.set_port(0); submitter_.ConnectActor(actor_id, addr, 0); ASSERT_EQ(worker_client_->callbacks.size(), 0); @@ -321,6 +330,7 @@ TEST_F(DirectActorSubmitterTest, TestActorRestartRetry) { ASSERT_TRUE(worker_client_->ReplyPushTask(Status::IOError(""))); // Actor gets restarted. + addr.set_port(1); submitter_.ConnectActor(actor_id, addr, 1); // A new task is submitted. ASSERT_TRUE(submitter_.SubmitTask(task4).ok()); @@ -342,6 +352,7 @@ TEST_F(DirectActorSubmitterTest, TestActorRestartOutOfOrderGcs) { ActorID actor_id = ActorID::Of(JobID::FromInt(0), TaskID::Nil(), 0); submitter_.AddActorQueueIfNotExists(actor_id); gcs::ActorTableData actor_data; + addr.set_port(0); submitter_.ConnectActor(actor_id, addr, 0); ASSERT_EQ(worker_client_->callbacks.size(), 0); ASSERT_EQ(num_clients_connected_, 1); @@ -354,6 +365,7 @@ TEST_F(DirectActorSubmitterTest, TestActorRestartOutOfOrderGcs) { ASSERT_TRUE(worker_client_->ReplyPushTask(Status::OK())); // Actor restarts, but we don't receive the disconnect message until later. + addr.set_port(1); submitter_.ConnectActor(actor_id, addr, 1); ASSERT_EQ(num_clients_connected_, 2); // Submit a task. @@ -381,6 +393,7 @@ TEST_F(DirectActorSubmitterTest, TestActorRestartOutOfOrderGcs) { ASSERT_FALSE(worker_client_->ReplyPushTask(Status::OK())); // We receive the late messages. Nothing happens. + addr.set_port(2); submitter_.ConnectActor(actor_id, addr, 2); submitter_.DisconnectActor(actor_id, 2, /*dead=*/false); ASSERT_EQ(num_clients_connected_, 2); @@ -392,6 +405,7 @@ TEST_F(DirectActorSubmitterTest, TestActorRestartOutOfOrderGcs) { // We receive more late messages. Nothing happens because the actor is dead. submitter_.DisconnectActor(actor_id, 4, /*dead=*/false); + addr.set_port(3); submitter_.ConnectActor(actor_id, addr, 4); ASSERT_EQ(num_clients_connected_, 2); // Submit a task. diff --git a/src/ray/core_worker/transport/direct_actor_transport.cc b/src/ray/core_worker/transport/direct_actor_transport.cc index c921301a1..ac256febc 100644 --- a/src/ray/core_worker/transport/direct_actor_transport.cc +++ b/src/ray/core_worker/transport/direct_actor_transport.cc @@ -135,6 +135,13 @@ void CoreWorkerDirectActorTaskSubmitter::ConnectActor(const ActorID &actor_id, return; } + if (queue->second.rpc_client && + queue->second.rpc_client->Addr().ip_address() == address.ip_address() && + queue->second.rpc_client->Addr().port() == address.port()) { + RAY_LOG(DEBUG) << "Skip actor that has already been connected, actor_id=" << actor_id; + return; + } + if (queue->second.state == rpc::ActorTableData::DEAD) { // This message is about an old version of the actor and the actor has // already died since then. Skip the connection.