From a76d91f50d4a975c511821f1c39b56036e072574 Mon Sep 17 00:00:00 2001 From: Philipp Moritz Date: Tue, 17 May 2016 15:08:49 -0700 Subject: [PATCH] Locality aware scheduler --- README.md | 1 + doc/scheduler.md | 25 +++++++++++++ src/scheduler.cc | 94 +++++++++++++++++++++++++++++++++++++++++++----- src/scheduler.h | 14 +++++++- 4 files changed, 124 insertions(+), 10 deletions(-) create mode 100644 doc/scheduler.md diff --git a/README.md b/README.md index ad5d738a4..0aea9c26e 100644 --- a/README.md +++ b/README.md @@ -8,6 +8,7 @@ For a description of our design decisions, see - [Reference Counting](doc/reference-counting.md) - [Aliasing](doc/aliasing.md) +- [Scheduler](doc/scheduler.md) ## Setup diff --git a/doc/scheduler.md b/doc/scheduler.md new file mode 100644 index 000000000..57f506826 --- /dev/null +++ b/doc/scheduler.md @@ -0,0 +1,25 @@ +# Scheduler + +The scheduling strategies currently implemented in Photon are fairly basic and +all use a central scheduler. + +* The naive scheduler assigns tasks to workers just taking into account +dependencies between tasks (no other information like data locality). It is +supposed to be an example for how to write a scheduler. We do not recommend +its use and it only works well for single node setups. Tasks are assigned in the +following way: For each idle worker, we iterate over the tasks in the task +queue. The first task that has all its requirements satisfied will be scheduled +on the worker. + +* The locality aware scheduler is more suited for multi node setups, but still +inappropriate for very large clusters. This is because the computational +overhead for each scheduling decision is O(mn) where m is the number of idle +workers and n is the number of tasks in the task queue. For each idle worker, +all tasks in the task queue are considered and the one that requires the +smallest number of objects to be shipped will be executed. + +We expect to implement more refined scheduling strategies in the future, +including more computationally efficient location aware scheduling, +scheduling that takes into account sizes of the shipped objects, and strategies +that do not require a central scheduler (which is a bottleneck for large +clusters). diff --git a/src/scheduler.cc b/src/scheduler.cc index 49c1224c7..a61aa71ca 100644 --- a/src/scheduler.cc +++ b/src/scheduler.cc @@ -6,6 +6,8 @@ #include "utils.h" +SchedulerService::SchedulerService(SchedulingAlgorithmType scheduling_algorithm) : scheduling_algorithm_(scheduling_algorithm) {} + Status SchedulerService::RemoteCall(ServerContext* context, const RemoteCallRequest* request, RemoteCallReply* reply) { std::unique_ptr task(new Call(request->call())); // need to copy, because request is const fntable_lock_.lock(); @@ -220,7 +222,13 @@ void SchedulerService::deliver_object(ObjRef objref, ObjStoreId from, ObjStoreId void SchedulerService::schedule() { // TODO(rkn): Do this more intelligently. perform_pulls(); // See what we can do in pull_queue_ - schedule_tasks(); // See what we can do in task_queue_ + if (scheduling_algorithm_ == SCHEDULING_ALGORITHM_NAIVE) { + schedule_tasks_naively(); // See what we can do in task_queue_ + } else if (scheduling_algorithm_ == SCHEDULING_ALGORITHM_LOCALITY_AWARE) { + schedule_tasks_location_aware(); // See what we can do in task_queue_ + } else { + ORCH_LOG(ORCH_FATAL, "scheduling algorithm not known"); + } perform_notify_aliases(); // See what we can do in alias_notification_queue_ } @@ -472,7 +480,7 @@ void SchedulerService::perform_pulls() { } } -void SchedulerService::schedule_tasks() { +void SchedulerService::schedule_tasks_naively() { std::lock_guard fntable_lock(fntable_lock_); std::lock_guard avail_workers_lock(avail_workers_lock_); std::lock_guard task_queue_lock(task_queue_lock_); @@ -497,6 +505,52 @@ void SchedulerService::schedule_tasks() { } } +void SchedulerService::schedule_tasks_location_aware() { + std::lock_guard fntable_lock(fntable_lock_); + std::lock_guard avail_workers_lock(avail_workers_lock_); + std::lock_guard task_queue_lock(task_queue_lock_); + for (int i = 0; i < avail_workers_.size(); ++i) { + // Submit all tasks whose arguments are ready. + WorkerId workerid = avail_workers_[i]; + ObjStoreId objstoreid = workers_[workerid].objstoreid; + auto bestit = task_queue_.end(); // keep track of the task that fits the worker best so far + size_t min_num_shipped_objects = std::numeric_limits::max(); // number of objects that need to be transfered for this worker + for (auto it = task_queue_.begin(); it != task_queue_.end(); ++it) { + const Call& task = *(*it); + auto& workers = fntable_[task.name()].workers(); + if (std::binary_search(workers.begin(), workers.end(), workerid) && can_run(task)) { + // determine how many objects would need to be shipped + size_t num_shipped_objects = 0; + for (int j = 0; j < task.arg_size(); ++j) { + if (!task.arg(j).has_obj()) { + ObjRef objref = task.arg(j).ref(); + if (!has_canonical_objref(objref)) { + ORCH_LOG(ORCH_FATAL, "no canonical object ref found even though task is ready; that should not be possible!"); + } + ObjRef canonical_objref = get_canonical_objref(objref); + // check if the object is already in the local object store + if (!std::binary_search(objtable_[canonical_objref].begin(), objtable_[canonical_objref].end(), objstoreid)) { + num_shipped_objects += 1; + } + } + } + if (num_shipped_objects < min_num_shipped_objects) { + min_num_shipped_objects = num_shipped_objects; + bestit = it; + } + } + } + // if we found a suitable task + if (bestit != task_queue_.end()) { + submit_task(std::move(*bestit), workerid); + task_queue_.erase(bestit); + std::swap(avail_workers_[i], avail_workers_[avail_workers_.size() - 1]); + avail_workers_.pop_back(); + i -= 1; + } + } +} + void SchedulerService::perform_notify_aliases() { std::lock_guard alias_notification_queue_lock(alias_notification_queue_lock_); for (int i = 0; i < alias_notification_queue_.size(); ++i) { @@ -663,12 +717,12 @@ void SchedulerService::get_equivalent_objrefs(ObjRef objref, std::vector upstream_objrefs(downstream_objref, equivalent_objrefs); } -void start_scheduler_service(const char* service_addr) { +void start_scheduler_service(const char* service_addr, SchedulingAlgorithmType scheduling_algorithm) { std::string service_address(service_addr); std::string::iterator split_point = split_ip_address(service_address); std::string port; port.assign(split_point, service_address.end()); - SchedulerService service; + SchedulerService service(scheduling_algorithm); ServerBuilder builder; builder.AddListeningPort(std::string("0.0.0.0:") + port, grpc::InsecureServerCredentials()); builder.RegisterService(&service); @@ -676,11 +730,33 @@ void start_scheduler_service(const char* service_addr) { server->Wait(); } -int main(int argc, char** argv) { - if (argc != 2) { - ORCH_LOG(ORCH_FATAL, "scheduler: expected one argument (scheduler ip address)"); - return 1; +char* get_cmd_option(char** begin, char** end, const std::string& option) { + char** it = std::find(begin, end, option); + if (it != end && ++it != end) { + return *it; } - start_scheduler_service(argv[1]); + return 0; +} + +int main(int argc, char** argv) { + SchedulingAlgorithmType scheduling_algorithm = SCHEDULING_ALGORITHM_LOCALITY_AWARE; + if (argc < 2) { + ORCH_LOG(ORCH_FATAL, "scheduler: expected at least one argument (scheduler ip address)"); + return 1; + } + if (argc > 2) { + char* scheduling_algorithm_name = get_cmd_option(argv, argv + argc, "--scheduler-algorithm"); + if (scheduling_algorithm_name) { + if(std::string(scheduling_algorithm_name) == "naive") { + std::cout << "using 'naive' scheduler" << std::endl; + scheduling_algorithm = SCHEDULING_ALGORITHM_NAIVE; + } + if(std::string(scheduling_algorithm_name) == "locality_aware") { + std::cout << "using 'locality aware' scheduler" << std::endl; + scheduling_algorithm = SCHEDULING_ALGORITHM_LOCALITY_AWARE; + } + } + } + start_scheduler_service(argv[1], scheduling_algorithm); return 0; } diff --git a/src/scheduler.h b/src/scheduler.h index c00643a41..63d6d241c 100644 --- a/src/scheduler.h +++ b/src/scheduler.h @@ -41,8 +41,15 @@ struct ObjStoreHandle { std::string address; }; +enum SchedulingAlgorithmType { + SCHEDULING_ALGORITHM_NAIVE = 0, + SCHEDULING_ALGORITHM_LOCALITY_AWARE = 1 +}; + class SchedulerService : public Scheduler::Service { public: + SchedulerService(SchedulingAlgorithmType scheduling_algorithm); + Status RemoteCall(ServerContext* context, const RemoteCallRequest* request, RemoteCallReply* reply) override; Status PushObj(ServerContext* context, const PushObjRequest* request, PushObjReply* reply) override; Status RequestObj(ServerContext* context, const RequestObjRequest* request, AckReply* reply) override; @@ -86,7 +93,10 @@ private: bool is_canonical(ObjRef objref); void perform_pulls(); - void schedule_tasks(); + // schedule tasks using the naive algorithm + void schedule_tasks_naively(); + // schedule tasks using a scheduling algorithm that takes into account data locality + void schedule_tasks_location_aware(); void perform_notify_aliases(); // checks if aliasing for objref has been completed @@ -149,6 +159,8 @@ private: // contained_objrefs_[objref] is a vector of all of the objrefs contained inside the object referred to by objref std::vector > contained_objrefs_; std::mutex contained_objrefs_lock_; + // the scheduling algorithm that will be used + SchedulingAlgorithmType scheduling_algorithm_; }; #endif