Improve PEP8 gluon/scheduler.py
This commit is contained in:
+43
-44
@@ -231,9 +231,9 @@ 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, esp for the dict decode
|
||||||
#and subsequent usage as function Keyword arguments unicode variable names won't work!
|
# 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
|
# borrowed from http://stackoverflow.com/questions/956867/how-to-get-string-objects-instead-unicode-ones-from-json-in-python
|
||||||
|
|
||||||
|
|
||||||
def _decode_list(lst):
|
def _decode_list(lst):
|
||||||
@@ -304,9 +304,9 @@ def executor(queue, task, out):
|
|||||||
if not isinstance(_function, CALLABLETYPES):
|
if not isinstance(_function, CALLABLETYPES):
|
||||||
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
|
||||||
globals().update(_env)
|
globals().update(_env)
|
||||||
@@ -795,7 +795,7 @@ class Scheduler(MetaScheduler):
|
|||||||
#let's do 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(
|
grabbed = db(
|
||||||
(st.assigned_worker_name == self.worker_name) &
|
(st.assigned_worker_name == self.worker_name) &
|
||||||
(st.status == ASSIGNED)
|
(st.status == ASSIGNED)
|
||||||
@@ -804,12 +804,12 @@ class Scheduler(MetaScheduler):
|
|||||||
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:
|
||||||
task.update_record(status=RUNNING, last_run_time=now)
|
task.update_record(status=RUNNING, last_run_time=now)
|
||||||
#noone will touch my task!
|
# noone will touch my task!
|
||||||
db.commit()
|
db.commit()
|
||||||
logger.debug(' work to do %s', task.id)
|
logger.debug(' work to do %s', task.id)
|
||||||
else:
|
else:
|
||||||
if self.is_a_ticker and self.greedy:
|
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)
|
||||||
else:
|
else:
|
||||||
@@ -825,10 +825,10 @@ class Scheduler(MetaScheduler):
|
|||||||
seconds=task.period * times_run
|
seconds=task.period * times_run
|
||||||
)
|
)
|
||||||
if times_run < task.repeats or task.repeats == 0:
|
if times_run < task.repeats or task.repeats == 0:
|
||||||
#need to run (repeating task)
|
# need to run (repeating task)
|
||||||
run_again = True
|
run_again = True
|
||||||
else:
|
else:
|
||||||
#no need to run again
|
# no need to run again
|
||||||
run_again = False
|
run_again = False
|
||||||
run_id = 0
|
run_id = 0
|
||||||
while True and not self.discard_results:
|
while True and not self.discard_results:
|
||||||
@@ -889,9 +889,9 @@ class Scheduler(MetaScheduler):
|
|||||||
sr = db.scheduler_run
|
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(sr.id == task.run_id).update(
|
db(sr.id == task.run_id).update(
|
||||||
@@ -903,7 +903,7 @@ class Scheduler(MetaScheduler):
|
|||||||
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(sr.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)
|
||||||
@@ -1056,16 +1056,16 @@ class Scheduler(MetaScheduler):
|
|||||||
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.w_stats.status == 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 == my_name).update(is_ticker=True)
|
db(sw.worker_name == my_name).update(is_ticker=True)
|
||||||
db(sw.worker_name != my_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 == my_name).update(is_ticker=False)
|
db(sw.worker_name == my_name).update(is_ticker=False)
|
||||||
else:
|
else:
|
||||||
not_busy = True
|
not_busy = True
|
||||||
@@ -1085,7 +1085,7 @@ class Scheduler(MetaScheduler):
|
|||||||
sw, st, sd = db.scheduler_worker, db.scheduler_task, db.scheduler_task_deps
|
sw, st, sd = db.scheduler_worker, db.scheduler_task, db.scheduler_task_deps
|
||||||
now = self.now()
|
now = self.now()
|
||||||
all_workers = db(sw.status == ACTIVE).select()
|
all_workers = db(sw.status == ACTIVE).select()
|
||||||
#build workers as dict of groups
|
# build workers as dict of groups
|
||||||
wkgroups = {}
|
wkgroups = {}
|
||||||
for w in all_workers:
|
for w in all_workers:
|
||||||
if w.worker_stats['status'] == 'RUNNING':
|
if w.worker_stats['status'] == 'RUNNING':
|
||||||
@@ -1098,14 +1098,14 @@ class Scheduler(MetaScheduler):
|
|||||||
else:
|
else:
|
||||||
wkgroups[gname]['workers'].append(
|
wkgroups[gname]['workers'].append(
|
||||||
{'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(
|
db(
|
||||||
(st.status.belongs((QUEUED, ASSIGNED))) &
|
(st.status.belongs((QUEUED, ASSIGNED))) &
|
||||||
(st.stop_time < now)
|
(st.stop_time < now)
|
||||||
).update(status=EXPIRED)
|
).update(status=EXPIRED)
|
||||||
|
|
||||||
#calculate dependencies
|
# calculate dependencies
|
||||||
deps_with_no_deps = db(
|
deps_with_no_deps = db(
|
||||||
(sd.can_visit == False) &
|
(sd.can_visit == False) &
|
||||||
(~sd.task_child.belongs(
|
(~sd.task_child.belongs(
|
||||||
@@ -1114,7 +1114,7 @@ class Scheduler(MetaScheduler):
|
|||||||
)
|
)
|
||||||
)._select(sd.task_child)
|
)._select(sd.task_child)
|
||||||
no_deps = db(
|
no_deps = db(
|
||||||
(st.status.belongs((QUEUED,ASSIGNED))) &
|
(st.status.belongs((QUEUED, ASSIGNED))) &
|
||||||
(
|
(
|
||||||
(sd.id == None) | (st.id.belongs(deps_with_no_deps))
|
(sd.id == None) | (st.id.belongs(deps_with_no_deps))
|
||||||
|
|
||||||
@@ -1137,27 +1137,27 @@ class Scheduler(MetaScheduler):
|
|||||||
|
|
||||||
|
|
||||||
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
|
# 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
|
# tasks per worker. This can be further tuned with some added
|
||||||
#intelligence (like esteeming how many tasks will a worker complete
|
# intelligence (like esteeming how many tasks will a worker complete
|
||||||
#before the ticker reassign them around, but the gain is quite small
|
# 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
|
# 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
|
# NB: ticker reassign tasks every 5 cycles, so if a worker completes its
|
||||||
#50 tasks in less than heartbeat*5 seconds,
|
# 50 tasks in less than heartbeat*5 seconds,
|
||||||
#it won't pick new tasks until heartbeat*5 seconds pass.
|
# it won't pick new tasks until heartbeat*5 seconds pass.
|
||||||
|
|
||||||
#If a worker is currently elaborating a long task, its tasks needs to
|
# If a worker is currently elaborating a long task, its tasks needs to
|
||||||
#be reassigned to other workers
|
# be reassigned to other workers
|
||||||
#this shuffles up things a bit, in order to give a task equal chances
|
# this shuffles up things a bit, in order to give a task equal chances
|
||||||
#to be executed
|
# to be executed
|
||||||
|
|
||||||
#let's freeze it up
|
# let's freeze it up
|
||||||
db.commit()
|
db.commit()
|
||||||
x = 0
|
x = 0
|
||||||
for group in wkgroups.keys():
|
for group in wkgroups.keys():
|
||||||
tasks = all_available(st.group_name == group).select(
|
tasks = all_available(st.group_name == group).select(
|
||||||
limitby=(0, limit), orderby = st.next_run_time)
|
limitby=(0, limit), orderby = st.next_run_time)
|
||||||
#let's break up the queue evenly among workers
|
# let's break up the queue evenly among workers
|
||||||
for task in tasks:
|
for task in tasks:
|
||||||
x += 1
|
x += 1
|
||||||
gname = task.group_name
|
gname = task.group_name
|
||||||
@@ -1182,13 +1182,13 @@ class Scheduler(MetaScheduler):
|
|||||||
).update(**d)
|
).update(**d)
|
||||||
wkgroups[gname]['workers'][myw]['c'] += 1
|
wkgroups[gname]['workers'][myw]['c'] += 1
|
||||||
db.commit()
|
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.w_stats.empty_runs = 0
|
self.w_stats.empty_runs = 0
|
||||||
self.w_stats.queue = x
|
self.w_stats.queue = x
|
||||||
self.w_stats.distribution = wkgroups
|
self.w_stats.distribution = wkgroups
|
||||||
self.w_stats.workers = len(all_workers)
|
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
|
self.greedy = x >= limit
|
||||||
logger.info('TICKER: workers are %s', len(all_workers))
|
logger.info('TICKER: workers are %s', len(all_workers))
|
||||||
@@ -1201,7 +1201,7 @@ class Scheduler(MetaScheduler):
|
|||||||
# 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, worker_name=None):
|
exclude=None, limit=None, worker_name=None):
|
||||||
"""Internal function to set worker's status"""
|
"""Internal function to set worker's status"""
|
||||||
ws = self.db.scheduler_worker
|
ws = self.db.scheduler_worker
|
||||||
if not group_names:
|
if not group_names:
|
||||||
@@ -1220,10 +1220,9 @@ class Scheduler(MetaScheduler):
|
|||||||
).update(status=action)
|
).update(status=action)
|
||||||
else:
|
else:
|
||||||
for group in group_names:
|
for group in group_names:
|
||||||
workers = self.db(
|
workers = self.db((ws.group_names.contains(group)) &
|
||||||
(ws.group_names.contains(group)) &
|
(~ws.status.belongs(exclusion))
|
||||||
(~ws.status.belongs(exclusion))
|
)._select(ws.id, limitby=(0, limit))
|
||||||
)._select(ws.id, limitby=(0,limit))
|
|
||||||
self.db(ws.id.belongs(workers)).update(status=action)
|
self.db(ws.id.belongs(workers)).update(status=action)
|
||||||
|
|
||||||
def disable(self, group_names=None, limit=None, worker_name=None):
|
def disable(self, group_names=None, limit=None, worker_name=None):
|
||||||
|
|||||||
Reference in New Issue
Block a user