diff --git a/gluon/scheduler.py b/gluon/scheduler.py index 37f50012..d18f955f 100644 --- a/gluon/scheduler.py +++ b/gluon/scheduler.py @@ -91,7 +91,7 @@ try: except: from simplejson import loads, dumps -IDENTIFIER = "%s#%s" % (socket.gethostname(),os.getpid()) +IDENTIFIER = "%s#%s" % (socket.gethostname(), os.getpid()) logger = logging.getLogger('web2py.scheduler.%s' % IDENTIFIER) @@ -117,7 +117,7 @@ STOP_TASK = 'STOP_TASK' EXPIRED = 'EXPIRED' SECONDS = 1 HEARTBEAT = 3 * SECONDS -MAXHIBERNATION = 10 +MAXHIBERNATION = 10 * HEARTBEAT CLEAROUT = '!clear!' CALLABLETYPES = (types.LambdaType, types.FunctionType, @@ -129,6 +129,7 @@ class Task(object): """Defines a "task" object that gets passed from the main thread to the executor's one """ + def __init__(self, app, function, timeout, args='[]', vars='{}', **kwargs): logger.debug(' new task allocated: %s.%s', app, function) self.app = app @@ -143,9 +144,10 @@ class Task(object): class TaskReport(object): - """Defines a "task report" object that gets passed from the executor's + """Defines a "task report" object that gets passed from the executor's thread to the main one """ + def __init__(self, status, result=None, output=None, tb=None): logger.debug(' new task report: %s', status) if tb: @@ -168,9 +170,10 @@ def demo_function(*argv, **kwargs): time.sleep(1) return 'done' -#the two functions below deal with simplejson decoding as unicode, esp for the dict decode -#and subsequent usage as function Keyword arguments unicode variable names won't work! -#borrowed from http://stackoverflow.com/questions/956867/how-to-get-string-objects-instead-unicode-ones-from-json-in-python +#the two functions below deal with simplejson decoding as unicode, +#esp for the dict decode and subsequent usage as function Keyword arguments +#unicode variable names won't work! +#borrowed from http://stackoverflow.com/questions/956867/ def _decode_list(lst): @@ -217,7 +220,7 @@ def executor(queue, task, out): def write(self, data): self.out_queue.put(data) - W2P_TASK = Storage({'id' : task.task_id, 'uuid' : task.uuid}) + W2P_TASK = Storage({'id': task.task_id, 'uuid': task.uuid}) stdout = LogOutput(out) try: if task.app: @@ -242,7 +245,7 @@ def executor(queue, task, out): raise NameError( "name '%s' not found in scheduler's environment" % f) #Inject W2P_TASK into environment - _env.update({'W2P_TASK' : W2P_TASK}) + _env.update({'W2P_TASK': W2P_TASK}) #Inject W2P_TASK into current from gluon import current current.W2P_TASK = W2P_TASK @@ -472,9 +475,10 @@ class Scheduler(MetaScheduler): worker_name(str): force worker_name to identify each process. Leave it to None to autoassign a name (hostname#pid) group_names(list): process tasks belonging to this group - heartbeat(int): how many seconds the worker sleeps between one execution - and the following one. Indirectly sets how many seconds will pass - between checks for new tasks + defaults to ['main'] if nothing gets passed + heartbeat(int): how many seconds the worker sleeps between one + execution and the following one. Indirectly sets how many seconds + will pass between checks for new tasks max_empty_runs(int): how many loops are allowed to pass without processing any tasks before exiting the process. 0 to keep always the process alive @@ -489,7 +493,7 @@ class Scheduler(MetaScheduler): """ def __init__(self, db, tasks=None, migrate=True, - worker_name=None, group_names=['main'], heartbeat=HEARTBEAT, + worker_name=None, group_names=None, heartbeat=HEARTBEAT, max_empty_runs=0, discard_results=False, utc_time=False): MetaScheduler.__init__(self) @@ -497,18 +501,26 @@ class Scheduler(MetaScheduler): self.db = db self.db_thread = None self.tasks = tasks - self.group_names = group_names + self.group_names = group_names or ['main'] self.heartbeat = heartbeat self.worker_name = worker_name or IDENTIFIER - #list containing status as recorded in the table plus a boost parameter - #for hibernation (i.e. when someone stop the worker acting on the worker table) - self.worker_status = [RUNNING, 1] self.max_empty_runs = max_empty_runs self.discard_results = discard_results self.is_a_ticker = False self.do_assign_tasks = False self.greedy = False self.utc_time = utc_time + self.w_stats = Storage( + dict( + status=RUNNING, + sleep=heartbeat, + total=0, + errors=0, + empty_runs=0, + queue=0, + distribution=None, + workers=0) + ) #dict holding statistics from gluon import current current._scheduler = self @@ -521,7 +533,7 @@ class Scheduler(MetaScheduler): elif migrate is True: return True elif isinstance(migrate, str): - return "%s%s.table" % (migrate , tablename) + return "%s%s.table" % (migrate, tablename) return True def now(self): @@ -604,6 +616,7 @@ class Scheduler(MetaScheduler): Field('status', requires=IS_IN_SET(WORKER_STATUS)), Field('is_ticker', 'boolean', default=False, writable=False), Field('group_names', 'list:string', default=self.group_names), + Field('worker_stats', 'json'), migrate=self.__get_migrate('scheduler_worker', migrate) ) @@ -632,24 +645,26 @@ class Scheduler(MetaScheduler): try: self.start_heartbeats() while True and self.have_heartbeat: - if self.worker_status[0] == DISABLED: - logger.debug('Someone stopped me, sleeping until better times come (%s)', self.worker_status[1]) + if self.w_stats.status == DISABLED: + logger.debug('Someone stopped me, sleeping until better' + ' times come (%s)', self.w_stats.sleep) self.sleep() continue logger.debug('looping...') task = self.wrapped_pop_task() if task: - self.empty_runs = 0 - self.worker_status[0] = RUNNING - self.report_task(task, self.async(task)) - self.worker_status[0] = ACTIVE + self.w_stats.empty_runs = 0 + self.w_stats.status = RUNNING + self.w_stats.total += 1 + self.wrapped_report_task(task, self.async(task)) + self.w_stats.status = ACTIVE else: - self.empty_runs += 1 + self.w_stats.empty_runs += 1 logger.debug('sleeping...') if self.max_empty_runs != 0: logger.debug('empty runs %s/%s', - self.empty_runs, self.max_empty_runs) - if self.empty_runs >= self.max_empty_runs: + self.w_stats.empty_runs, self.max_empty_runs) + if self.w_stats.empty_runs >= self.max_empty_runs: logger.info( 'empty runs limit reached, killing myself') self.die() @@ -673,6 +688,7 @@ class Scheduler(MetaScheduler): logger.debug('Tasks assigned...') break except: + self.w_stats.errors += 1 db.rollback() logger.error('TICKER: error assigning tasks (%s)', x) x += 1 @@ -692,6 +708,7 @@ class Scheduler(MetaScheduler): return rtn break except: + self.w_stats.errors += 1 db.rollback() logger.error(' error popping tasks') x += 1 @@ -702,13 +719,15 @@ class Scheduler(MetaScheduler): now = self.now() st = self.db.scheduler_task if self.is_a_ticker and self.do_assign_tasks: - #I'm a ticker, and 5 loops passed without reassigning tasks, let's do - #that and loop again + #I'm a ticker, and 5 loops passed without reassigning tasks, + #let's do that and loop again self.wrapped_assign_tasks(db) return None #ready to process something - grabbed = db(st.assigned_worker_name == self.worker_name)( - st.status == ASSIGNED) + grabbed = db( + (st.assigned_worker_name == self.worker_name) & + (st.status == ASSIGNED) + ) task = grabbed.select(limitby=(0, 1), orderby=st.next_run_time).first() if task: @@ -717,7 +736,7 @@ class Scheduler(MetaScheduler): db.commit() logger.debug(' work to do %s', task.id) else: - if self.greedy and self.is_a_ticker: + if self.is_a_ticker and self.greedy: #there are other tasks ready to be assigned logger.info('TICKER: greedy loop') self.wrapped_assign_tasks(db) @@ -753,7 +772,8 @@ class Scheduler(MetaScheduler): except: time.sleep(0.5) db.rollback() - logger.info('new task %(id)s "%(task_name)s" %(application_name)s.%(function_name)s' % task) + logger.info('new task %(id)s "%(task_name)s"' + ' %(application_name)s.%(function_name)s' % task) return Task( app=task.application_name, function=task.function_name, @@ -771,73 +791,81 @@ class Scheduler(MetaScheduler): sync_output=task.sync_output, uuid=task.uuid) + def wrapped_report_task(self, task, task_report): + """Commodity function to call `report_task` and trap exceptions + If an exception is raised, assume it happened because of database + contention and retries `pop_task` after 0.5 seconds + """ + db = self.db + while True: + try: + self.report_task(task, task_report) + db.commit() + break + except: + self.w_stats.errors += 1 + db.rollback() + logger.error(' error storing result') + time.sleep(0.5) + def report_task(self, task, task_report): """Takes care of storing the result according to preferences and deals with logic for repeating tasks""" db = self.db now = self.now() - while True: - try: - if not self.discard_results: - if task_report.result != 'null' or task_report.tb: - #result is 'null' as a string if task completed - #if it's stopped it's None as NoneType, so we record - #the STOPPED "run" anyway - logger.debug(' recording task report in db (%s)', - task_report.status) - db(db.scheduler_run.id == task.run_id).update( - status=task_report.status, - stop_time=now, - run_result=task_report.result, - run_output=task_report.output, - traceback=task_report.tb) - else: - logger.debug(' deleting task report in db because of no result') - db(db.scheduler_run.id == task.run_id).delete() - #if there is a stop_time and the following run would exceed it - is_expired = (task.stop_time - and task.next_run_time > task.stop_time - and True or False) - status = (task.run_again and is_expired and EXPIRED - or task.run_again and not is_expired - and QUEUED or COMPLETED) - if task_report.status == COMPLETED: - d = dict(status=status, - next_run_time=task.next_run_time, - times_run=task.times_run, - times_failed=0 - ) - db(db.scheduler_task.id == task.task_id)( - db.scheduler_task.status == RUNNING).update(**d) - else: - st_mapping = {'FAILED': 'FAILED', - 'TIMEOUT': 'TIMEOUT', - 'STOPPED': 'QUEUED'}[task_report.status] - status = (task.retry_failed - and task.times_failed < task.retry_failed - and QUEUED or task.retry_failed == -1 - and QUEUED or st_mapping) - db( - (db.scheduler_task.id == task.task_id) & - (db.scheduler_task.status == RUNNING) - ).update( - times_failed=db.scheduler_task.times_failed + 1, - next_run_time=task.next_run_time, - status=status - ) - db.commit() - logger.info('task completed (%s)', task_report.status) - break - except: - db.rollback() - time.sleep(0.5) + st = db.scheduler_task + sr = db.scheduler_run + if not self.discard_results: + if task_report.result != 'null' or task_report.tb: + #result is 'null' as a string if task completed + #if it's stopped it's None as NoneType, so we record + #the STOPPED "run" anyway + logger.debug(' recording task report in db (%s)', + task_report.status) + db(sr.id == task.run_id).update( + status=task_report.status, + stop_time=now, + run_result=task_report.result, + run_output=task_report.output, + traceback=task_report.tb) + else: + logger.debug(' deleting task report in db because of no result') + db(sr.id == task.run_id).delete() + #if there is a stop_time and the following run would exceed it + is_expired = (task.stop_time + and task.next_run_time > task.stop_time + and True or False) + status = (task.run_again and is_expired and EXPIRED + or task.run_again and not is_expired + and QUEUED or COMPLETED) + if task_report.status == COMPLETED: + d = dict(status=status, + next_run_time=task.next_run_time, + times_run=task.times_run, + times_failed=0 + ) + db(st.id == task.task_id).update(**d) + else: + st_mapping = {'FAILED': 'FAILED', + 'TIMEOUT': 'TIMEOUT', + 'STOPPED': 'QUEUED'}[task_report.status] + status = (task.retry_failed + and task.times_failed < task.retry_failed + and QUEUED or task.retry_failed == -1 + and QUEUED or st_mapping) + db(st.id == task.task_id).update( + times_failed=db.scheduler_task.times_failed + 1, + next_run_time=task.next_run_time, + status=status + ) + logger.info('task completed (%s)', task_report.status) def adj_hibernation(self): """Used to increase the "sleep" interval for DISABLED workers""" - if self.worker_status[0] == DISABLED: - wk_st = self.worker_status[1] - hibernation = wk_st + 1 if wk_st < MAXHIBERNATION else MAXHIBERNATION - self.worker_status[1] = hibernation + if self.w_stats.status == DISABLED: + wk_st = self.w_stats.sleep + hibernation = wk_st + HEARTBEAT if wk_st < MAXHIBERNATION else MAXHIBERNATION + self.w_stats.sleep = hibernation def send_heartbeat(self, counter): """This function is vital for proper coordination among available @@ -866,59 +894,68 @@ class Scheduler(MetaScheduler): if not mybackedstatus: sw.insert(status=ACTIVE, worker_name=self.worker_name, first_heartbeat=now, last_heartbeat=now, - group_names=self.group_names) - self.worker_status = [ACTIVE, 1] # activating the process + group_names=self.group_names, + worker_stats=self.w_stats) + self.w_stats.status = ACTIVE + self.w_stats.sleep = self.heartbeat mybackedstatus = ACTIVE else: mybackedstatus = mybackedstatus.status if mybackedstatus == DISABLED: # keep sleeping - self.worker_status[0] = DISABLED - if self.worker_status[1] == MAXHIBERNATION: - logger.debug('........recording heartbeat (%s)', self.worker_status[0]) + self.w_stats.status = DISABLED + if self.w_stats.sleep >= MAXHIBERNATION: + logger.debug('........recording heartbeat (%s)', + self.w_stats.status) db(sw.worker_name == self.worker_name).update( - last_heartbeat=now) + last_heartbeat=now, + worker_stats=self.w_stats) elif mybackedstatus == TERMINATE: - self.worker_status[0] = TERMINATE + self.w_stats.status = TERMINATE logger.debug("Waiting to terminate the current task") self.give_up() - return elif mybackedstatus == KILL: - self.worker_status[0] = KILL + self.w_stats.status = KILL self.die() + return else: if mybackedstatus == STOP_TASK: logger.info('Asked to kill the current task') self.terminate_process() - logger.debug('........recording heartbeat (%s)', self.worker_status[0]) + logger.debug('........recording heartbeat (%s)', + self.w_stats.status) db(sw.worker_name == self.worker_name).update( - last_heartbeat=now, status=ACTIVE) - self.worker_status[1] = 1 # re-activating the process - if self.worker_status[0] != RUNNING: - self.worker_status[0] = ACTIVE + last_heartbeat=now, status=ACTIVE, + worker_stats=self.w_stats) + self.w_stats.sleep = self.heartbeat # re-activating the process + if self.w_stats.status != RUNNING: + self.w_stats.status = ACTIVE self.do_assign_tasks = False if counter % 5 == 0 or mybackedstatus == PICK: try: - # delete inactive workers - expiration = now - datetime.timedelta(seconds=self.heartbeat * 3) + # delete dead workers + expiration = now - datetime.timedelta( + seconds=self.heartbeat * 3) departure = now - datetime.timedelta( - seconds=self.heartbeat * 3 * MAXHIBERNATION) + seconds=self.heartbeat * 3 * 15) logger.debug( ' freeing workers that have not sent heartbeat') - inactive_workers = db( + dead_workers = db( ((sw.last_heartbeat < expiration) & (sw.status == ACTIVE)) | ((sw.last_heartbeat < departure) & (sw.status != ACTIVE)) ) - db(st.assigned_worker_name.belongs( - inactive_workers._select(sw.worker_name)))(st.status == RUNNING)\ - .update(assigned_worker_name='', status=QUEUED) - inactive_workers.delete() + dead_workers_name = dead_workers._select(sw.worker_name) + db( + (st.assigned_worker_name.belongs(dead_workers_name)) & + (st.status == RUNNING) + ).update(assigned_worker_name='', status=QUEUED) + dead_workers.delete() try: self.is_a_ticker = self.being_a_ticker() except: logger.error('Error coordinating TICKER') - if self.worker_status[0] == ACTIVE: + if self.w_stats.status == ACTIVE: self.do_assign_tasks = True except: logger.error('Error cleaning up') @@ -936,23 +973,24 @@ class Scheduler(MetaScheduler): """ db = self.db_thread sw = db.scheduler_worker + my_name = self.worker_name all_active = db( - (sw.worker_name != self.worker_name) & (sw.status == ACTIVE) - ).select() + (sw.worker_name != my_name) & (sw.status == ACTIVE) + ).select(sw.is_ticker, sw.worker_name) ticker = all_active.find(lambda row: row.is_ticker is True).first() - not_busy = self.worker_status[0] == ACTIVE + not_busy = self.w_stats.status == ACTIVE if not ticker: #if no other tickers are around if not_busy: #only if I'm not busy - db(sw.worker_name == self.worker_name).update(is_ticker=True) - db(sw.worker_name != self.worker_name).update(is_ticker=False) + db(sw.worker_name == my_name).update(is_ticker=True) + db(sw.worker_name != my_name).update(is_ticker=False) logger.info("TICKER: I'm a ticker") else: #I'm busy if len(all_active) >= 1: #so I'll "downgrade" myself to a "poor worker" - db(sw.worker_name == self.worker_name).update(is_ticker=False) + db(sw.worker_name == my_name).update(is_ticker=False) else: not_busy = True db.commit() @@ -984,8 +1022,10 @@ class Scheduler(MetaScheduler): {'name': w.worker_name, 'c': 0}) #set queued tasks that expired between "runs" (i.e., you turned off #the scheduler): then it wasn't expired, but now it is - db(st.status.belongs( - (QUEUED, ASSIGNED)))(st.stop_time < now).update(status=EXPIRED) + db( + (st.status.belongs((QUEUED, ASSIGNED))) & + (st.stop_time < now) + ).update(status=EXPIRED) all_available = db( (st.status.belongs((QUEUED, ASSIGNED))) & @@ -996,16 +1036,19 @@ class Scheduler(MetaScheduler): (st.enabled == True) ) limit = len(all_workers) * (50 / (len(wkgroups) or 1)) - #if there are a moltitude of tasks, let's figure out a maximum of tasks per worker. - #this can be adjusted with some added intelligence (like esteeming how many tasks will - #a worker complete before the ticker reassign them around, but the gain is quite small - #50 is quite a sweet spot also for fast tasks, with sane heartbeat values - #NB: ticker reassign tasks every 5 cycles, so if a worker completes his 50 tasks in less - #than heartbeat*5 seconds, it won't pick new tasks until heartbeat*5 seconds pass. + #if there are a moltitude of tasks, let's figure out a maximum of + #tasks per worker. This can be further tuned with some added + #intelligence (like esteeming how many tasks will a worker complete + #before the ticker reassign them around, but the gain is quite small + #50 is a sweet spot also for fast tasks, with sane heartbeat values + #NB: ticker reassign tasks every 5 cycles, so if a worker completes its + #50 tasks in less than heartbeat*5 seconds, + #it won't pick new tasks until heartbeat*5 seconds pass. - #If a worker is currently elaborating a long task, all other tasks assigned - #to him needs to be reassigned "freely" to other workers, that may be free. - #this shuffles up things a bit, in order to maintain the idea of a semi-linear scalability + #If a worker is currently elaborating a long task, its tasks needs to + #be reassigned to other workers + #this shuffles up things a bit, in order to give a task equal chances + #to be executed #let's freeze it up db.commit() @@ -1025,63 +1068,96 @@ class Scheduler(MetaScheduler): if w['c'] < counter: myw = i counter = w['c'] + assigned_wn = wkgroups[gname]['workers'][myw]['name'] d = dict( status=ASSIGNED, - assigned_worker_name=wkgroups[gname]['workers'][myw]['name'] + assigned_worker_name=assigned_wn ) if not task.task_name: d['task_name'] = task.function_name - db((st.id==task.id) & (st.status.belongs((QUEUED, ASSIGNED)))).update(**d) - db.commit() + db( + (st.id==task.id) & + (st.status.belongs((QUEUED, ASSIGNED))) + ).update(**d) wkgroups[gname]['workers'][myw]['c'] += 1 - + db.commit() #I didn't report tasks but I'm working nonetheless!!!! if x > 0: - self.empty_runs = 0 + self.w_stats.empty_runs = 0 + self.w_stats.queue = x + self.w_stats.distribution = wkgroups + self.w_stats.workers = len(all_workers) #I'll be greedy only if tasks assigned are equal to the limit # (meaning there could be others ready to be assigned) - self.greedy = x >= limit and True or False + self.greedy = x >= limit logger.info('TICKER: workers are %s', len(all_workers)) logger.info('TICKER: tasks are %s', x) def sleep(self): """Calculates the number of seconds to sleep according to worker's status and `heartbeat` parameter""" - time.sleep(self.heartbeat * self.worker_status[1]) + time.sleep(self.w_stats.sleep) # should only sleep until next available task - def set_worker_status(self, group_names=None, action=ACTIVE): + def set_worker_status(self, group_names=None, action=ACTIVE, + exclude=None, limit=None): """Internal function to set worker's status""" + ws = self.db.scheduler_worker if not group_names: group_names = self.group_names elif isinstance(group_names, str): group_names = [group_names] - for group in group_names: - self.db( - self.db.scheduler_worker.group_names.contains(group) - ).update(status=action) + exclusion = exclude and exclude.append(action) or [action] + if not limit: + for group in group_names: + self.db( + (ws.group_names.contains(group)) & + (~ws.status.belongs(exclusion)) + ).update(status=action) + else: + for group in group_names: + workers = self.db( + (ws.group_names.contains(group)) & + (~ws.status.belongs(exclusion)) + )._select(ws.id, limitby=(0,limit)) + print self.db(ws.id.belongs(workers)).update(status=action) - def disable(self, group_names=None): + def disable(self, group_names=None, limit=None): """Sets DISABLED on the workers processing `group_names` tasks. A DISABLED worker will be kept alive but it won't be able to process any waiting tasks, essentially putting it to sleep. By default, all group_names of Scheduler's instantation are selected""" - self.set_worker_status(group_names=group_names,action=DISABLED) + self.set_worker_status( + group_names=group_names, + action=DISABLED, + exclude=[DISABLED, KILL, TERMINATE], + limit=limit) - def resume(self, group_names=None): + def resume(self, group_names=None, limit=None): """Wakes a worker up (it will be able to process queued tasks)""" - self.set_worker_status(group_names=group_names,action=ACTIVE) + self.set_worker_status( + group_names=group_names, + action=ACTIVE, + exclude=[KILL, TERMINATE], + limit=limit) - def terminate(self, group_names=None): + def terminate(self, group_names=None, limit=None): """Sets TERMINATE as worker status. The worker will wait for any currently running tasks to be executed and then it will exit gracefully """ - self.set_worker_status(group_names=group_names,action=TERMINATE) + self.set_worker_status( + group_names=group_names, + action=TERMINATE, + exclude=[KILL], + limit=limit) - def kill(self, group_names=None): + def kill(self, group_names=None, limit=None): """Sets KILL as worker status. The worker will be killed even if it's processing a task.""" - self.set_worker_status(group_names=group_names,action=KILL) + self.set_worker_status( + group_names=group_names, + action=KILL, + limit=limit) def queue_task(self, function, pargs=[], pvars={}, **kwargs): """ @@ -1121,7 +1197,9 @@ class Scheduler(MetaScheduler): if not rtn.errors: rtn.uuid = tuuid if immediate: - self.db(self.db.scheduler_worker.is_ticker == True).update(status=PICK) + self.db( + (self.db.scheduler_worker.is_ticker == True) + ).update(status=PICK) else: rtn.uuid = None return rtn @@ -1208,14 +1286,19 @@ class Scheduler(MetaScheduler): else: raise SyntaxError( "You can retrieve results only by id or uuid") - task = self.db(q).select(st.id, st.status, st.assigned_worker_name).first() + task = self.db(q).select(st.id, st.status, st.assigned_worker_name) + task = task.first() rtn = None if not task: return rtn if task.status == 'RUNNING': - rtn = self.db(sw.worker_name == task.assigned_worker_name).update(status=STOP_TASK) + q = sw.worker_name == task.assigned_worker_name + rtn = self.db(q).update(status=STOP_TASK) elif task.status == 'QUEUED': - rtn = self.db(q).update(stop_time=self.now(), enabled=False, status=STOPPED) + rtn = self.db(q).update( + stop_time=self.now(), + enabled=False, + status=STOPPED) return rtn