From acf1703afdef0581502489a45e26600f00eea26a Mon Sep 17 00:00:00 2001 From: Robert Nishihara Date: Wed, 28 Dec 2016 22:33:20 -0800 Subject: [PATCH] Implement naive scheduling algorithm using local scheduler load. (#164) * Implement naive scheduling algorithm using local scheduler load. * Have the global scheduler estimate load on local schedulers better. * Fixes. --- src/common/state/local_scheduler_table.h | 7 ++ src/global_scheduler/global_scheduler.c | 45 ++++++++++- src/global_scheduler/global_scheduler.h | 23 +++++- .../global_scheduler_algorithm.c | 79 +++++++++++++------ src/photon/photon_algorithm.c | 1 + src/photon/photon_scheduler.h | 2 +- 6 files changed, 128 insertions(+), 29 deletions(-) diff --git a/src/common/state/local_scheduler_table.h b/src/common/state/local_scheduler_table.h index 48d7d556d..633b1a607 100644 --- a/src/common/state/local_scheduler_table.h +++ b/src/common/state/local_scheduler_table.h @@ -4,8 +4,15 @@ #include "db.h" #include "table.h" +/** This struct is sent with heartbeat messages from the local scheduler to the + * global scheduler, and it contains information about the load on the local + * scheduler. */ typedef struct { + /** The total number of workers that are connected to this local scheduler. */ + int total_num_workers; + /** The number of tasks queued in this local scheduler. */ int task_queue_length; + /** The number of workers that are available and waiting for tasks. */ int available_workers; } local_scheduler_info; diff --git a/src/global_scheduler/global_scheduler.c b/src/global_scheduler/global_scheduler.c index 66ac936d4..d3fdb5984 100644 --- a/src/global_scheduler/global_scheduler.c +++ b/src/global_scheduler/global_scheduler.c @@ -33,6 +33,14 @@ void assign_task_to_local_scheduler(global_scheduler_state *state, object_id_to_string(task_task_id(task), id_string, ID_STRING_SIZE)); UNUSED(id_string); task_table_update(state->db, copy_task(task), &retry, NULL, NULL); + + /* TODO(rkn): We should probably pass around local_scheduler struct pointers + * instead of db_client_id objects. */ + /* Update the local scheduler info. */ + local_scheduler *local_scheduler = + get_local_scheduler(state, local_scheduler_id); + local_scheduler->num_tasks_sent += 1; + local_scheduler->num_recent_tasks_sent += 1; } global_scheduler_state *init_global_scheduler(event_loop *loop, @@ -86,6 +94,20 @@ void signal_handler(int signal) { /* End of the cleanup code. */ +local_scheduler *get_local_scheduler(global_scheduler_state *state, + db_client_id photon_id) { + local_scheduler *local_scheduler_ptr; + for (int i = 0; i < utarray_len(state->local_schedulers); ++i) { + local_scheduler_ptr = + (local_scheduler *) utarray_eltptr(state->local_schedulers, i); + if (db_client_ids_equal(local_scheduler_ptr->id, photon_id)) { + LOG_DEBUG("photon_id matched cached local scheduler entry."); + return local_scheduler_ptr; + } + } + return NULL; +} + void process_task_waiting(task *task, void *user_context) { global_scheduler_state *state = (global_scheduler_state *) user_context; LOG_DEBUG("Task waiting callback is called."); @@ -129,6 +151,15 @@ void process_new_db_client(db_client_id db_client_id, } /* Add new local scheduler to the state. */ + local_scheduler local_scheduler; + local_scheduler.id = db_client_id; + local_scheduler.num_tasks_sent = 0; + local_scheduler.num_recent_tasks_sent = 0; + local_scheduler.info.task_queue_length = 0; + local_scheduler.info.available_workers = 0; + utarray_push_back(state->local_schedulers, &local_scheduler); + + /* Allow the scheduling algorithm to process this event. */ handle_new_local_scheduler(state, state->policy_state, db_client_id); } } @@ -204,8 +235,18 @@ void local_scheduler_table_handler(db_client_id client_id, "Local scheduler heartbeat from db_client_id %s", object_id_to_string((object_id) client_id, id_string, ID_STRING_SIZE)); UNUSED(id_string); - LOG_DEBUG("Task queue length is %d", info.task_queue_length); - LOG_DEBUG("Num available workers is %d", info.available_workers); + LOG_DEBUG( + "total workers = %d, task queue length = %d, available workers = %d", + info.num_total_workers, info.task_queue_length, info.available_workers); + /* Update the local scheduler info struct. */ + local_scheduler *local_scheduler_ptr = get_local_scheduler(state, client_id); + if (local_scheduler_ptr != NULL) { + /* Reset the number of tasks sent since the last heartbeat. */ + local_scheduler_ptr->num_recent_tasks_sent = 0; + local_scheduler_ptr->info = info; + } else { + LOG_WARN("client_id didn't match any cached local scheduler entries"); + } } void start_server(const char *redis_addr, int redis_port) { diff --git a/src/global_scheduler/global_scheduler.h b/src/global_scheduler/global_scheduler.h index 6239896c2..b052babe1 100644 --- a/src/global_scheduler/global_scheduler.h +++ b/src/global_scheduler/global_scheduler.h @@ -4,6 +4,7 @@ #include "task.h" #include "state/db.h" +#include "state/local_scheduler_table.h" #include "utarray.h" #include "uthash.h" @@ -14,6 +15,12 @@ typedef struct { /** The number of tasks sent from the global scheduler to this local * scheduler. */ int64_t num_tasks_sent; + /** The number of tasks sent from the global scheduler to this local scheduler + * since the last heartbeat arrived. */ + int64_t num_recent_tasks_sent; + /** The latest information about the local scheduler capacity. This is updated + * every time a new local scheduler heartbeat arrives. */ + local_scheduler_info info; } local_scheduler; typedef struct global_scheduler_policy_state global_scheduler_policy_state; @@ -43,7 +50,9 @@ typedef struct { event_loop *loop; /** The global state store database. */ db_handle *db; - /** The local schedulers that are connected to Redis. */ + /** The local schedulers that are connected to Redis. TODO(rkn): This probably + * needs to be a hashtable since we often look up the local_scheduler struct + * based on its db_client_id. */ UT_array *local_schedulers; /** The state managed by the scheduling policy. */ global_scheduler_policy_state *policy_state; @@ -52,6 +61,18 @@ typedef struct { scheduler_object_info *scheduler_object_info_table; } global_scheduler_state; +/** + * This is a helper method to look up the local scheduler struct that + * corresponds to a particular photon_id. + * + * @param state The state of the global scheduler. + * @param The photon_id of the local scheduler. + * @return The corresponding local scheduler struct. If the global scheduler is + * not aware of the local scheduler, then this will be NULL. + */ +local_scheduler *get_local_scheduler(global_scheduler_state *state, + db_client_id photon_id); + void assign_task_to_local_scheduler(global_scheduler_state *state, task *task, db_client_id local_scheduler_id); diff --git a/src/global_scheduler/global_scheduler_algorithm.c b/src/global_scheduler/global_scheduler_algorithm.c index 997bfbd77..5c6f09c83 100644 --- a/src/global_scheduler/global_scheduler_algorithm.c +++ b/src/global_scheduler/global_scheduler_algorithm.c @@ -1,3 +1,5 @@ +#include + #include "object_info.h" #include "task.h" #include "task_table.h" @@ -16,19 +18,45 @@ void destroy_global_scheduler_policy( free(policy_state); } +/** + * This is a helper method that assigns a task to the next local scheduler in a + * round robin fashion. + */ void handle_task_round_robin(global_scheduler_state *state, global_scheduler_policy_state *policy_state, task *task) { - if (utarray_len(state->local_schedulers) > 0) { - local_scheduler *scheduler = (local_scheduler *) utarray_eltptr( - state->local_schedulers, policy_state->round_robin_index); - scheduler->num_tasks_sent++; - policy_state->round_robin_index += 1; - policy_state->round_robin_index %= utarray_len(state->local_schedulers); - assign_task_to_local_scheduler(state, task, scheduler->id); - } else { - CHECKM(0, "No local schedulers. We currently don't handle this case."); + CHECKM(utarray_len(state->local_schedulers) > 0, + "No local schedulers. We currently don't handle this case.") + local_scheduler *scheduler = (local_scheduler *) utarray_eltptr( + state->local_schedulers, policy_state->round_robin_index); + policy_state->round_robin_index += 1; + policy_state->round_robin_index %= utarray_len(state->local_schedulers); + assign_task_to_local_scheduler(state, task, scheduler->id); +} + +/** + * This is a helper method that assigns a task to the local scheduler with the + * minimal load. + */ +void handle_task_minimum_load(global_scheduler_state *state, + global_scheduler_policy_state *policy_state, + task *task) { + CHECKM(utarray_len(state->local_schedulers) > 0, + "No local schedulers. We currently don't handle this case.") + int current_minimal_load_estimate = INT_MAX; + local_scheduler *current_local_scheduler_ptr = NULL; + for (int i = 0; i < utarray_len(state->local_schedulers); ++i) { + local_scheduler *local_scheduler_ptr = + (local_scheduler *) utarray_eltptr(state->local_schedulers, i); + int load_estimate = local_scheduler_ptr->info.task_queue_length + + local_scheduler_ptr->num_recent_tasks_sent; + if (load_estimate <= current_minimal_load_estimate) { + current_minimal_load_estimate = load_estimate; + current_local_scheduler_ptr = local_scheduler_ptr; + } } + DCHECK(current_local_scheduler_ptr != NULL); + assign_task_to_local_scheduler(state, task, current_local_scheduler_ptr->id); } object_size_entry *create_object_size_hashmap(global_scheduler_state *state, @@ -129,16 +157,8 @@ db_client_id get_photon_id(global_scheduler_state *state, /* Check to make sure this photon_db_client_id matches one of the * schedulers. */ - int i; - for (i = 0; i < utarray_len(state->local_schedulers); ++i) { - local_scheduler *local_scheduler_ptr = - (local_scheduler *) utarray_eltptr(state->local_schedulers, i); - if (memcmp(&local_scheduler_ptr->id, &photon_id, sizeof(photon_id)) == 0) { - LOG_DEBUG("photon_id matched cached local scheduler entry."); - break; - } - } - if (i == utarray_len(state->local_schedulers)) { + local_scheduler *local_scheduler_ptr = get_local_scheduler(state, photon_id); + if (local_scheduler_ptr == NULL) { LOG_WARN("photon_id didn't match any cached local scheduler entries"); } return photon_id; @@ -189,7 +209,20 @@ void handle_task_waiting(global_scheduler_state *state, " num_returns = %" PRId64 "\n", task_num_args(task_spec), task_num_returns(task_spec)); - assign_task_to_local_scheduler(state, task, photon_id); + /* Get the local scheduler for this photon ID. */ + local_scheduler *local_scheduler_ptr = get_local_scheduler(state, photon_id); + CHECK(local_scheduler_ptr != NULL); + /* If this local scheduler has enough capacity, assign the task to this local + * scheduler. Otherwise assign the task to the global scheduler with the + * minimal load. */ + int64_t load_estimate = local_scheduler_ptr->info.task_queue_length + + local_scheduler_ptr->num_recent_tasks_sent; + if (local_scheduler_ptr->info.available_workers > 0 && + load_estimate < local_scheduler_ptr->info.total_num_workers) { + assign_task_to_local_scheduler(state, task, photon_id); + } else { + handle_task_minimum_load(state, policy_state, task); + } free_object_size_hashmap(object_size_table); } @@ -208,9 +241,5 @@ void handle_local_scheduler_heartbeat( void handle_new_local_scheduler(global_scheduler_state *state, global_scheduler_policy_state *policy_state, db_client_id db_client_id) { - local_scheduler local_scheduler; - memset(&local_scheduler, 0, sizeof(local_scheduler)); - local_scheduler.id = db_client_id; - local_scheduler.num_tasks_sent = 0; - utarray_push_back(state->local_schedulers, &local_scheduler); + /* Do nothing for now. */ } diff --git a/src/photon/photon_algorithm.c b/src/photon/photon_algorithm.c index cc861bdf1..7066d3dee 100644 --- a/src/photon/photon_algorithm.c +++ b/src/photon/photon_algorithm.c @@ -95,6 +95,7 @@ void provide_scheduler_info(local_scheduler_state *state, scheduling_algorithm_state *algorithm_state, local_scheduler_info *info) { task_queue_entry *elt; + info->total_num_workers = utarray_len(state->workers); DL_COUNT(algorithm_state->task_queue, elt, info->task_queue_length); info->available_workers = utarray_len(algorithm_state->available_workers); } diff --git a/src/photon/photon_scheduler.h b/src/photon/photon_scheduler.h index f178d2a38..b62041887 100644 --- a/src/photon/photon_scheduler.h +++ b/src/photon/photon_scheduler.h @@ -5,7 +5,7 @@ #include "event_loop.h" /* The duration between local scheduler heartbeats. */ -#define LOCAL_SCHEDULER_HEARTBEAT_TIMEOUT_MILLISECONDS 1000 +#define LOCAL_SCHEDULER_HEARTBEAT_TIMEOUT_MILLISECONDS 100 /** * Establish a connection to a new client.