Merge pull request #457 from niphlod/enhancement/scheduler
scheduler's refactoring
This commit is contained in:
+238
-155
@@ -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
|
||||||
@@ -143,9 +144,10 @@ class Task(object):
|
|||||||
|
|
||||||
|
|
||||||
class TaskReport(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
|
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
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user