- refactored internals

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