[New Scheduler] Additional unit tests (#9990)

This commit is contained in:
Alex Wu
2020-08-10 11:44:06 -07:00
committed by GitHub
parent 4b10bdf8fc
commit 2ebf76c7a3
2 changed files with 100 additions and 2 deletions
@@ -31,6 +31,8 @@ bool ClusterTaskManager::SchedulePendingTasks() {
auto request_resources =
task.GetTaskSpecification().GetRequiredResources().GetResourceMap();
int64_t _unused;
// TODO (Alex): We should distinguish between infeasible tasks and a fully
// utilized cluster.
std::string node_id_string =
cluster_resource_scheduler_->GetBestSchedulableNode(request_resources, &_unused);
if (node_id_string.empty()) {
@@ -250,7 +250,8 @@ std::shared_ptr<ClusterResourceScheduler> CreateSingleNodeScheduler(
return scheduler;
}
Task CreateTask(const std::unordered_map<std::string, double> &required_resources) {
Task CreateTask(const std::unordered_map<std::string, double> &required_resources,
int num_args = 0) {
TaskSpecBuilder spec_builder;
TaskID id = RandomTaskId();
JobID job_id = RandomJobId();
@@ -258,6 +259,12 @@ Task CreateTask(const std::unordered_map<std::string, double> &required_resource
spec_builder.SetCommonTaskSpec(
id, Language::PYTHON, FunctionDescriptorBuilder::BuildPython("", "", "", ""),
job_id, TaskID::Nil(), 0, TaskID::Nil(), address, 0, required_resources, {});
for (int i = 0; i < num_args; i++) {
ObjectID put_id = ObjectID::ForPut(TaskID::Nil(), /*index=*/i + 1);
spec_builder.AddArg(TaskArgByReference(put_id, rpc::Address()));
}
rpc::TaskExecutionSpec execution_spec_message;
execution_spec_message.set_num_forwards(1);
return Task(spec_builder.Build(), TaskExecutionSpecification(execution_spec_message));
@@ -326,13 +333,102 @@ TEST_F(ClusterTaskManagerTest, BasicTest) {
task_manager_.DispatchScheduledTasksToWorkers(pool_, leased_workers_);
ASSERT_EQ(callback_occurred, true);
ASSERT_TRUE(callback_occurred);
ASSERT_EQ(leased_workers_.size(), 1);
ASSERT_EQ(pool_.workers.size(), 0);
ASSERT_EQ(fulfills_dependencies_calls_, 0);
ASSERT_EQ(node_info_calls_, 0);
}
TEST_F(ClusterTaskManagerTest, NoFeasibleNodeTest) {
std::shared_ptr<MockWorker> worker =
std::make_shared<MockWorker>(WorkerID::FromRandom(), 1234);
pool_.PushWorker(std::dynamic_pointer_cast<WorkerInterface>(worker));
Task task = CreateTask({{ray::kCPU_ResourceLabel, 999}});
rpc::RequestWorkerLeaseReply reply;
bool callback_called = false;
bool *callback_called_ptr = &callback_called;
auto callback = [callback_called_ptr]() { *callback_called_ptr = true; };
task_manager_.QueueTask(task, &reply, callback);
task_manager_.SchedulePendingTasks();
task_manager_.DispatchScheduledTasksToWorkers(pool_, leased_workers_);
ASSERT_FALSE(callback_called);
ASSERT_EQ(leased_workers_.size(), 0);
// Worker is unused.
ASSERT_EQ(pool_.workers.size(), 1);
ASSERT_EQ(fulfills_dependencies_calls_, 0);
ASSERT_EQ(node_info_calls_, 0);
}
TEST_F(ClusterTaskManagerTest, ResourceTakenWhileResolving) {
/*
Test the race condition in which a task is assigned to a node, but cannot
run because its dependencies are unresolved. Once its dependencies are
resolved, the node no longer has available resources.
*/
std::shared_ptr<MockWorker> worker =
std::make_shared<MockWorker>(WorkerID::FromRandom(), 1234);
std::shared_ptr<MockWorker> worker2 =
std::make_shared<MockWorker>(WorkerID::FromRandom(), 12345);
pool_.PushWorker(std::dynamic_pointer_cast<WorkerInterface>(worker2));
pool_.PushWorker(std::dynamic_pointer_cast<WorkerInterface>(worker));
rpc::RequestWorkerLeaseReply reply;
int num_callbacks = 0;
int *num_callbacks_ptr = &num_callbacks;
auto callback = [num_callbacks_ptr]() {
(*num_callbacks_ptr) = *num_callbacks_ptr + 1;
};
/* Blocked on dependencies */
auto task = CreateTask({{ray::kCPU_ResourceLabel, 5}}, 1);
dependencies_fulfilled_ = false;
task_manager_.QueueTask(task, &reply, callback);
task_manager_.SchedulePendingTasks();
task_manager_.DispatchScheduledTasksToWorkers(pool_, leased_workers_);
ASSERT_EQ(num_callbacks, 0);
ASSERT_EQ(leased_workers_.size(), 0);
ASSERT_EQ(pool_.workers.size(), 2);
/* This task can run */
auto task2 = CreateTask({{ray::kCPU_ResourceLabel, 5}});
task_manager_.QueueTask(task2, &reply, callback);
task_manager_.SchedulePendingTasks();
task_manager_.DispatchScheduledTasksToWorkers(pool_, leased_workers_);
ASSERT_EQ(num_callbacks, 1);
ASSERT_EQ(leased_workers_.size(), 1);
ASSERT_EQ(pool_.workers.size(), 1);
/* First task is unblocked now, but resources are no longer available */
auto id = task.GetTaskSpecification().TaskId();
std::vector<TaskID> unblocked = {id};
task_manager_.TasksUnblocked(unblocked);
task_manager_.DispatchScheduledTasksToWorkers(pool_, leased_workers_);
ASSERT_EQ(num_callbacks, 1);
ASSERT_EQ(leased_workers_.size(), 1);
ASSERT_EQ(pool_.workers.size(), 1);
/* Second task finishes, making space for the original task */
single_node_resource_scheduler_->FreeLocalTaskResources(
worker->GetAllocatedInstances());
// single_node_resource_scheduler_->UpdateLocalAvailableResourcesFromResourceInstances();
leased_workers_.clear();
task_manager_.DispatchScheduledTasksToWorkers(pool_, leased_workers_);
// Task2 is now done so task can run.
ASSERT_EQ(num_callbacks, 2);
ASSERT_EQ(leased_workers_.size(), 1);
ASSERT_EQ(pool_.workers.size(), 0);
}
} // namespace raylet
} // namespace ray