Fix flaky placement group test bug (#10915)

Co-authored-by: 灵洵 <fengbin.ffb@antfin.com>
This commit is contained in:
fangfengbin
2020-09-20 19:50:55 -07:00
committed by GitHub
co-authored by 灵洵
parent d7c42d6d92
commit 3e94c690c7
2 changed files with 65 additions and 71 deletions
@@ -29,9 +29,10 @@ class GcsPlacementGroupSchedulerTest : public ::testing::Test {
io_service_.run();
}));
raylet_client_ = std::make_shared<GcsServerMocker::MockRayletResourceClient>();
raylet_client1_ = std::make_shared<GcsServerMocker::MockRayletResourceClient>();
raylet_client2_ = std::make_shared<GcsServerMocker::MockRayletResourceClient>();
for (int index = 0; index < 3; ++index) {
raylet_clients_.push_back(
std::make_shared<GcsServerMocker::MockRayletResourceClient>());
}
gcs_table_storage_ = std::make_shared<gcs::InMemoryGcsTableStorage>(io_service_);
gcs_pub_sub_ = std::make_shared<GcsServerMocker::MockGcsPubSub>(redis_client_);
gcs_node_manager_ = std::make_shared<gcs::GcsNodeManager>(
@@ -41,15 +42,7 @@ class GcsPlacementGroupSchedulerTest : public ::testing::Test {
scheduler_ = std::make_shared<GcsServerMocker::MockedGcsPlacementGroupScheduler>(
io_service_, gcs_table_storage_, *gcs_node_manager_,
/*lease_client_fplacement_groupy=*/
[this](const rpc::Address &address) {
if (0 == address.port()) {
return raylet_client_;
} else if (1 == address.port()) {
return raylet_client1_;
} else {
return raylet_client2_;
}
});
[this](const rpc::Address &address) { return raylet_clients_[address.port()]; });
}
void TearDown() override {
@@ -91,7 +84,7 @@ class GcsPlacementGroupSchedulerTest : public ::testing::Test {
// The lease request should not be send and the scheduling of placement_group should
// fail as there are no available nodes.
ASSERT_EQ(raylet_client_->num_lease_requested, 0);
ASSERT_EQ(raylet_clients_[0]->num_lease_requested, 0);
ASSERT_EQ(0, success_placement_groups_.size());
ASSERT_EQ(1, failure_placement_groups_.size());
ASSERT_EQ(placement_group, failure_placement_groups_.front());
@@ -118,10 +111,10 @@ class GcsPlacementGroupSchedulerTest : public ::testing::Test {
success_placement_groups_.emplace_back(std::move(placement_group));
});
ASSERT_EQ(2, raylet_client_->num_lease_requested);
ASSERT_EQ(2, raylet_client_->lease_callbacks.size());
ASSERT_TRUE(raylet_client_->GrantPrepareBundleResources());
ASSERT_TRUE(raylet_client_->GrantPrepareBundleResources());
ASSERT_EQ(2, raylet_clients_[0]->num_lease_requested);
ASSERT_EQ(2, raylet_clients_[0]->lease_callbacks.size());
ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources());
ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources());
WaitPendingDone(failure_placement_groups_, 0);
WaitPendingDone(success_placement_groups_, 1);
ASSERT_EQ(placement_group, success_placement_groups_.front());
@@ -152,8 +145,8 @@ class GcsPlacementGroupSchedulerTest : public ::testing::Test {
AddNode(Mocker::GenNodeInfo(0), 2);
scheduler_->ScheduleUnplacedBundles(placement_group, failure_handler,
success_handler);
ASSERT_TRUE(raylet_client_->GrantPrepareBundleResources());
ASSERT_TRUE(raylet_client_->GrantPrepareBundleResources());
ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources());
ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources());
WaitPendingDone(success_placement_groups_, 1);
}
@@ -164,9 +157,7 @@ class GcsPlacementGroupSchedulerTest : public ::testing::Test {
boost::asio::io_service io_service_;
std::shared_ptr<gcs::StoreClient> store_client_;
std::shared_ptr<GcsServerMocker::MockRayletResourceClient> raylet_client_;
std::shared_ptr<GcsServerMocker::MockRayletResourceClient> raylet_client1_;
std::shared_ptr<GcsServerMocker::MockRayletResourceClient> raylet_client2_;
std::vector<std::shared_ptr<GcsServerMocker::MockRayletResourceClient>> raylet_clients_;
std::shared_ptr<gcs::GcsNodeManager> gcs_node_manager_;
std::shared_ptr<GcsServerMocker::MockedGcsPlacementGroupScheduler> scheduler_;
std::vector<std::shared_ptr<gcs::GcsPlacementGroup>> success_placement_groups_;
@@ -225,12 +216,12 @@ TEST_F(GcsPlacementGroupSchedulerTest, TestSchedulePlacementGroupReplyFailure) {
success_placement_groups_.emplace_back(std::move(placement_group));
});
ASSERT_EQ(2, raylet_client_->num_lease_requested);
ASSERT_EQ(2, raylet_client_->lease_callbacks.size());
ASSERT_EQ(2, raylet_clients_[0]->num_lease_requested);
ASSERT_EQ(2, raylet_clients_[0]->lease_callbacks.size());
// Reply failure, so the placement group scheduling failed.
ASSERT_TRUE(raylet_client_->GrantPrepareBundleResources(false));
ASSERT_TRUE(raylet_client_->GrantPrepareBundleResources(false));
ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources(false));
ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources(false));
WaitPendingDone(failure_placement_groups_, 1);
WaitPendingDone(success_placement_groups_, 0);
ASSERT_EQ(placement_group, failure_placement_groups_.front());
@@ -282,12 +273,12 @@ TEST_F(GcsPlacementGroupSchedulerTest, TestSchedulePlacementGroupReturnResource)
success_placement_groups_.emplace_back(std::move(placement_group));
});
ASSERT_EQ(2, raylet_client_->num_lease_requested);
ASSERT_EQ(2, raylet_client_->lease_callbacks.size());
ASSERT_EQ(2, raylet_clients_[0]->num_lease_requested);
ASSERT_EQ(2, raylet_clients_[0]->lease_callbacks.size());
// One bundle success and the other failed.
ASSERT_TRUE(raylet_client_->GrantPrepareBundleResources());
ASSERT_TRUE(raylet_client_->GrantPrepareBundleResources(false));
ASSERT_EQ(1, raylet_client_->num_return_requested);
ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources());
ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources(false));
ASSERT_EQ(1, raylet_clients_[0]->num_return_requested);
// Reply the placement_group creation request, then the placement_group should be
// scheduled successfully.
WaitPendingDone(failure_placement_groups_, 1);
@@ -308,8 +299,9 @@ TEST_F(GcsPlacementGroupSchedulerTest, TestStrictPackStrategyBalancedScheduling)
};
// Schedule placement group, it will be evenly distributed over the two nodes.
int select_node0_count = 0;
int select_node1_count = 0;
int node_select_count[2] = {0, 0};
int node_commit_count[2] = {0, 0};
int node_index = 0;
for (int index = 0; index < 10; ++index) {
auto request =
Mocker::GenCreatePlacementGroupRequest("", rpc::PlacementStrategy::STRICT_PACK);
@@ -317,19 +309,20 @@ TEST_F(GcsPlacementGroupSchedulerTest, TestStrictPackStrategyBalancedScheduling)
scheduler_->ScheduleUnplacedBundles(placement_group, failure_handler,
success_handler);
if (!raylet_client_->lease_callbacks.empty()) {
ASSERT_TRUE(raylet_client_->GrantPrepareBundleResources());
ASSERT_TRUE(raylet_client_->GrantPrepareBundleResources());
++select_node0_count;
} else {
ASSERT_TRUE(raylet_client1_->GrantPrepareBundleResources());
ASSERT_TRUE(raylet_client1_->GrantPrepareBundleResources());
++select_node1_count;
}
node_index = !raylet_clients_[0]->lease_callbacks.empty() ? 0 : 1;
++node_select_count[node_index];
node_commit_count[node_index] += 2;
ASSERT_TRUE(raylet_clients_[node_index]->GrantPrepareBundleResources());
ASSERT_TRUE(raylet_clients_[node_index]->GrantPrepareBundleResources());
auto condition = [this, node_index, node_commit_count]() {
return raylet_clients_[node_index]->num_commit_requested ==
node_commit_count[node_index];
};
EXPECT_TRUE(WaitForCondition(condition, timeout_ms_.count()));
}
WaitPendingDone(success_placement_groups_, 10);
ASSERT_EQ(select_node0_count, 5);
ASSERT_EQ(select_node1_count, 5);
ASSERT_EQ(node_select_count[0], 5);
ASSERT_EQ(node_select_count[1], 5);
}
TEST_F(GcsPlacementGroupSchedulerTest, TestStrictPackStrategyReschedulingWhenNodeAdd) {
@@ -351,8 +344,8 @@ TEST_F(GcsPlacementGroupSchedulerTest, TestStrictPackStrategyResourceCheck) {
Mocker::GenCreatePlacementGroupRequest("", rpc::PlacementStrategy::STRICT_PACK);
auto placement_group = std::make_shared<gcs::GcsPlacementGroup>(request);
scheduler_->ScheduleUnplacedBundles(placement_group, failure_handler, success_handler);
ASSERT_TRUE(raylet_client_->GrantPrepareBundleResources());
ASSERT_TRUE(raylet_client_->GrantPrepareBundleResources());
ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources());
ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources());
WaitPendingDone(success_placement_groups_, 1);
// Node1 has less number of bundles, but it doesn't satisfy the resource
@@ -364,8 +357,8 @@ TEST_F(GcsPlacementGroupSchedulerTest, TestStrictPackStrategyResourceCheck) {
auto placement_group2 =
std::make_shared<gcs::GcsPlacementGroup>(create_placement_group_request2);
scheduler_->ScheduleUnplacedBundles(placement_group2, failure_handler, success_handler);
ASSERT_TRUE(raylet_client_->GrantPrepareBundleResources());
ASSERT_TRUE(raylet_client_->GrantPrepareBundleResources());
ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources());
ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources());
WaitPendingDone(success_placement_groups_, 2);
}
@@ -390,19 +383,18 @@ TEST_F(GcsPlacementGroupSchedulerTest, DestroyPlacementGroup) {
absl::MutexLock lock(&vector_mutex_);
success_placement_groups_.emplace_back(std::move(placement_group));
});
ASSERT_TRUE(raylet_client_->GrantPrepareBundleResources());
ASSERT_TRUE(raylet_client_->GrantPrepareBundleResources());
ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources());
ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources());
WaitPendingDone(failure_placement_groups_, 0);
WaitPendingDone(success_placement_groups_, 1);
RAY_LOG(ERROR) << "sanngbin1";
const auto &placement_group_id = placement_group->GetPlacementGroupID();
scheduler_->DestroyPlacementGroupBundleResourcesIfExists(placement_group_id);
ASSERT_TRUE(raylet_client_->GrantCancelResourceReserve());
ASSERT_TRUE(raylet_client_->GrantCancelResourceReserve());
ASSERT_TRUE(raylet_clients_[0]->GrantCancelResourceReserve());
ASSERT_TRUE(raylet_clients_[0]->GrantCancelResourceReserve());
// Subsequent destroy request should not do anything.
scheduler_->DestroyPlacementGroupBundleResourcesIfExists(placement_group_id);
ASSERT_FALSE(raylet_client_->GrantCancelResourceReserve());
ASSERT_FALSE(raylet_client_->GrantCancelResourceReserve());
ASSERT_FALSE(raylet_clients_[0]->GrantCancelResourceReserve());
ASSERT_FALSE(raylet_clients_[0]->GrantCancelResourceReserve());
}
TEST_F(GcsPlacementGroupSchedulerTest, DestroyCancelledPlacementGroup) {
@@ -429,11 +421,11 @@ TEST_F(GcsPlacementGroupSchedulerTest, DestroyCancelledPlacementGroup) {
});
// Now, cancel the schedule request.
ASSERT_TRUE(raylet_client_->GrantPrepareBundleResources());
ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources());
scheduler_->MarkScheduleCancelled(placement_group_id);
ASSERT_TRUE(raylet_client_->GrantPrepareBundleResources());
ASSERT_TRUE(raylet_client_->GrantCancelResourceReserve());
ASSERT_TRUE(raylet_client_->GrantCancelResourceReserve());
ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources());
ASSERT_TRUE(raylet_clients_[0]->GrantCancelResourceReserve());
ASSERT_TRUE(raylet_clients_[0]->GrantCancelResourceReserve());
WaitPendingDone(failure_placement_groups_, 1);
}
@@ -459,13 +451,13 @@ TEST_F(GcsPlacementGroupSchedulerTest, TestPackStrategyLargeBundlesScheduling) {
Mocker::GenCreatePlacementGroupRequest("", rpc::PlacementStrategy::PACK, 15);
auto placement_group = std::make_shared<gcs::GcsPlacementGroup>(request);
scheduler_->ScheduleUnplacedBundles(placement_group, failure_handler, success_handler);
RAY_CHECK(raylet_client_->num_lease_requested > 0);
RAY_CHECK(raylet_client1_->num_lease_requested > 0);
for (int index = 0; index < raylet_client_->num_lease_requested; ++index) {
ASSERT_TRUE(raylet_client_->GrantPrepareBundleResources());
RAY_CHECK(raylet_clients_[0]->num_lease_requested > 0);
RAY_CHECK(raylet_clients_[1]->num_lease_requested > 0);
for (int index = 0; index < raylet_clients_[0]->num_lease_requested; ++index) {
ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources());
}
for (int index = 0; index < raylet_client1_->num_lease_requested; ++index) {
ASSERT_TRUE(raylet_client1_->GrantPrepareBundleResources());
for (int index = 0; index < raylet_clients_[1]->num_lease_requested; ++index) {
ASSERT_TRUE(raylet_clients_[1]->GrantPrepareBundleResources());
}
WaitPendingDone(success_placement_groups_, 1);
}
@@ -492,8 +484,8 @@ TEST_F(GcsPlacementGroupSchedulerTest, TestRescheduleWhenNodeDead) {
};
scheduler_->ScheduleUnplacedBundles(placement_group, failure_handler, success_handler);
ASSERT_TRUE(raylet_client_->GrantPrepareBundleResources());
ASSERT_TRUE(raylet_client1_->GrantPrepareBundleResources());
ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources());
ASSERT_TRUE(raylet_clients_[1]->GrantPrepareBundleResources());
WaitPendingDone(success_placement_groups_, 1);
auto bundles_on_node0 =
@@ -509,8 +501,8 @@ TEST_F(GcsPlacementGroupSchedulerTest, TestRescheduleWhenNodeDead) {
// TODO(ffbin): We need to see which node the other bundles that have been placed are
// deployed on, and spread them as far as possible. It will be implemented in the next
// pr.
raylet_client_->GrantPrepareBundleResources();
raylet_client1_->GrantPrepareBundleResources();
raylet_clients_[0]->GrantPrepareBundleResources();
raylet_clients_[1]->GrantPrepareBundleResources();
WaitPendingDone(success_placement_groups_, 2);
}
@@ -543,8 +535,8 @@ TEST_F(GcsPlacementGroupSchedulerTest, TestStrictSpreadStrategyResourceCheck) {
auto node2 = Mocker::GenNodeInfo(2);
AddNode(node2);
scheduler_->ScheduleUnplacedBundles(placement_group, failure_handler, success_handler);
ASSERT_TRUE(raylet_client_->GrantPrepareBundleResources());
ASSERT_TRUE(raylet_client2_->GrantPrepareBundleResources());
ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources());
ASSERT_TRUE(raylet_clients_[2]->GrantPrepareBundleResources());
WaitPendingDone(success_placement_groups_, 1);
}
@@ -170,6 +170,7 @@ struct GcsServerMocker {
const BundleSpecification &bundle_spec,
const ray::rpc::ClientCallback<ray::rpc::CommitBundleResourcesReply> &callback)
override {
num_commit_requested += 1;
rpc::CommitBundleResourcesReply reply;
callback(Status::OK(), reply);
}
@@ -215,6 +216,7 @@ struct GcsServerMocker {
int num_lease_requested = 0;
int num_return_requested = 0;
int num_commit_requested = 0;
ClientID node_id = ClientID::FromRandom();
std::list<rpc::ClientCallback<rpc::PrepareBundleResourcesReply>> lease_callbacks = {};
std::list<rpc::ClientCallback<rpc::CancelResourceReserveReply>> return_callbacks = {};