mirror of
https://github.com/galaxyproject/galaxy.git
synced 2026-09-24 16:30:27 +08:00
Remove the user round robin scheduling policy (it didn't work as intended). The former job_scheduler_policy config option now has no effect.
This commit is contained in:
@@ -59,7 +59,6 @@ class Configuration( object ):
|
||||
self.template_cache = resolve_path( kwargs.get( "template_cache_path", "database/compiled_templates" ), self.root )
|
||||
self.local_job_queue_workers = int( kwargs.get( "local_job_queue_workers", "5" ) )
|
||||
self.cluster_job_queue_workers = int( kwargs.get( "cluster_job_queue_workers", "3" ) )
|
||||
self.job_scheduler_policy = kwargs.get("job_scheduler_policy", "FIFO")
|
||||
self.job_queue_cleanup_interval = int( kwargs.get("job_queue_cleanup_interval", "5") )
|
||||
self.cluster_files_directory = os.path.abspath( kwargs.get( "cluster_files_directory", "database/pbs" ) )
|
||||
self.job_working_directory = resolve_path( kwargs.get( "job_working_directory", "database/job_working_directory" ), self.root )
|
||||
|
||||
@@ -75,30 +75,6 @@ class JobQueue( object ):
|
||||
self.sa_session = app.model.context
|
||||
# Should we read jobs form the database, or use an in memory queue
|
||||
self.track_jobs_in_database = app.config.get_bool( 'track_jobs_in_database', False )
|
||||
# Check if any special scheduling policy should be used. If not, default is FIFO.
|
||||
sched_policy = app.config.get('job_scheduler_policy', 'FIFO')
|
||||
# Parse the scheduler policy string. The policy class implements a special queue.
|
||||
# Ready-to-run jobs are inserted into this queue
|
||||
if sched_policy != 'FIFO' :
|
||||
try :
|
||||
self.use_policy = True
|
||||
if ":" in sched_policy :
|
||||
modname , policy_class = sched_policy.split(":")
|
||||
modfields = modname.split(".")
|
||||
module = __import__(modname)
|
||||
for mod in modfields[1:] : module = getattr( module, mod)
|
||||
# instantiate the policy class
|
||||
self.squeue = getattr( module , policy_class )(self.app)
|
||||
else :
|
||||
self.use_policy = False
|
||||
log.info("Scheduler policy not defined as expected, defaulting to FIFO")
|
||||
except AttributeError, detail: # try may throw AttributeError
|
||||
self.use_policy = False
|
||||
log.exception("Error while loading scheduler policy class, defaulting to FIFO")
|
||||
else :
|
||||
self.use_policy = False
|
||||
|
||||
log.info("job scheduler policy is %s" %sched_policy)
|
||||
# Keep track of the pid that started the job manager, only it
|
||||
# has valid threads
|
||||
self.parent_pid = os.getpid()
|
||||
@@ -221,13 +197,8 @@ class JobQueue( object ):
|
||||
elif job_state == JOB_INPUT_DELETED:
|
||||
log.info( "job %d unable to run: one or more inputs deleted" % job.job_id )
|
||||
elif job_state == JOB_READY:
|
||||
# If special queuing is enabled, put the ready jobs in the special queue
|
||||
if self.use_policy :
|
||||
self.squeue.put( job )
|
||||
log.debug( "job %d put in policy queue" % job.job_id )
|
||||
else: # or dispatch the job directly
|
||||
self.dispatcher.put( job )
|
||||
log.debug( "job %d dispatched" % job.job_id)
|
||||
self.dispatcher.put( job )
|
||||
log.debug( "job %d dispatched" % job.job_id)
|
||||
elif job_state == JOB_DELETED:
|
||||
msg = "job %d deleted by user while still queued" % job.job_id
|
||||
job.info = msg
|
||||
@@ -244,20 +215,6 @@ class JobQueue( object ):
|
||||
log.exception( "failure running job %d" % job.job_id )
|
||||
# Update the waiting list
|
||||
self.waiting = new_waiting
|
||||
# If special (e.g. fair) scheduling is enabled, dispatch all jobs
|
||||
# currently in the special queue
|
||||
if self.use_policy :
|
||||
while 1:
|
||||
try:
|
||||
sjob = self.squeue.get()
|
||||
self.dispatcher.put( sjob )
|
||||
log.debug( "job %d dispatched" % sjob.job_id )
|
||||
except Empty:
|
||||
# squeue is empty, so stop dispatching
|
||||
break
|
||||
except Exception, e: # if something else breaks while dispatching
|
||||
job.fail( "failure running job %d: %s" % ( sjob.job_id, str( e ) ) )
|
||||
log.exception( "failure running job %d" % sjob.job_id )
|
||||
# Done with the session
|
||||
self.sa_session.remove()
|
||||
|
||||
@@ -319,7 +276,6 @@ class JobWrapper( object ):
|
||||
"""
|
||||
def __init__(self, job, tool, queue ):
|
||||
self.job_id = job.id
|
||||
# This is immutable, we cache it for the scheduling policy to use if needed
|
||||
self.session_id = job.session_id
|
||||
self.tool = tool
|
||||
self.queue = queue
|
||||
|
||||
@@ -1,4 +0,0 @@
|
||||
"""
|
||||
This package contains job scheduling policy classes.
|
||||
|
||||
"""
|
||||
@@ -1,146 +0,0 @@
|
||||
import logging, threading, time
|
||||
from Queue import Queue, Empty
|
||||
|
||||
from galaxy import config, util, model, tools, jobs
|
||||
|
||||
log = logging.getLogger( __name__ )
|
||||
|
||||
# This class is designed to provide job dispatch fairness for Galaxy users. It uses a
|
||||
# dictionary of session (key) Queue (value) pairs to ensure no user (session)
|
||||
# can hog the Galaxy run queue.
|
||||
# Each get() call will return a job from a consecutive session
|
||||
# in the dict or throw an Empty exception if there are no jobs in any queue
|
||||
class UserRoundRobin( object ):
|
||||
"""
|
||||
Class UserRoundRobin provides per-user job scheduling fairness. It uses a dictionary of session-Queue() pairs
|
||||
to ensure no user/session can hog the Galaxy run queue. Each get() call will either return a job from a consecutive session
|
||||
in the dict or throw an Empty exception if the dict has no jobs in any queues.
|
||||
"""
|
||||
# private: dictionary of queues (its not really private in Python)
|
||||
__DOQ = {}
|
||||
|
||||
def __init__(self, app):
|
||||
self.app = app
|
||||
self.__DOQ = {}
|
||||
self.keylist = []
|
||||
self.iterator = None
|
||||
# these are used to decide when to do a DOQ cleanup
|
||||
self.cleanup_tstamp = time.time()
|
||||
# get from ini/config and convert to secs
|
||||
self.cleanup_mininterval = self.app.config.job_queue_cleanup_interval * 60
|
||||
# Don't allow cleanup interval less than 5 minutes (should this be hardcoded ?)
|
||||
if self.cleanup_mininterval < 300 :
|
||||
self.cleanup_mininterval = 300
|
||||
# locks for get and put methods
|
||||
self.putlock = threading.Lock()
|
||||
self.getlock = threading.Lock()
|
||||
log.info("RoundRobin policy: initialized ")
|
||||
|
||||
# Insert a job in the dict of queues
|
||||
def put(self, job):
|
||||
self.putlock.acquire()
|
||||
try :
|
||||
# get this job's user/session id
|
||||
sessid = job.get_session_id()
|
||||
# Check if a queue already exists for user/session and add one if not
|
||||
if self.__DOQ.has_key(sessid) :
|
||||
self.__DOQ[sessid].put(job)
|
||||
log.debug("RoundRobin queue: inserted new job for user/session = %d" % sessid)
|
||||
else :
|
||||
self.__DOQ[sessid] = Queue()
|
||||
self.__DOQ[sessid].put(job)
|
||||
log.debug("RoundRobin queue: user/session did not exist, created new jobqueue for session = %d" % sessid)
|
||||
finally :
|
||||
self.putlock.release()
|
||||
|
||||
# Return a job from the dictionary of queues. Each get() call tries to
|
||||
# return a job from another queue in the dictionary .
|
||||
# Throws Empty if it cant find any jobs in the dictionary of queues
|
||||
# This method also does a cleanup of the dict regularly (specified)
|
||||
def get(self) :
|
||||
self.getlock.acquire()
|
||||
try :
|
||||
# get the next user/session in the dict
|
||||
sessionid = self.__get_next_session()
|
||||
if sessionid is not None :
|
||||
log.debug("RoundRobin queue: retrieving job from job queue for session = %d" % sessionid)
|
||||
return self.__DOQ[sessionid].get()
|
||||
else :
|
||||
# sessionid = None implies empty dictionary, throw back to caller
|
||||
raise Empty
|
||||
finally :
|
||||
# Clean up DOQ
|
||||
self.__timed_clean_up() #cleanup will happen at specified intervals
|
||||
self.getlock.release()
|
||||
|
||||
|
||||
# In case the Queue.get_nowait() method is used somewhere
|
||||
def get_nowait(self):
|
||||
return self.get()
|
||||
|
||||
# Returns the total number of jobs in the dict (counts all queues)
|
||||
# Locks the get() and put() methods during calculation.
|
||||
# Analogous to qsize() in Queue.Queue. Not guaranteed to be correct
|
||||
def qsize(self):
|
||||
try :
|
||||
count = 0
|
||||
self.getlock.acquire()
|
||||
self.putlock.acquire()
|
||||
for sessid in self.__DOQ :
|
||||
count += self.__DOQ[sessid].qsize()
|
||||
return count
|
||||
finally :
|
||||
self.putlock.release()
|
||||
self.getlock.release()
|
||||
|
||||
# Internal method - get the next user/session (key) from any nonempty queue in the dict.
|
||||
# Returns None if dictionary is empty.
|
||||
# Note: We use a separate list to hold the DOQ keys and create an iterator from this list
|
||||
# because an iterator created directly from DOQ barfs if DOQ size changes while iterating
|
||||
def __get_next_session(self):
|
||||
try:
|
||||
# Check if this is either startup or previous iteration ended
|
||||
if self.iterator is None :
|
||||
self.keylist = self.__DOQ.keys()
|
||||
self.iterator = iter(self.keylist)
|
||||
# get the first available job from any nonempty queue
|
||||
while 1 :
|
||||
tmpsid = self.iterator.next()
|
||||
if not self.__DOQ[tmpsid].empty() :
|
||||
break
|
||||
return tmpsid
|
||||
except StopIteration :
|
||||
# StopIteration implies we hit the end of the dict.
|
||||
# re-initalize iterator, start from beginning (can we improve this ?)
|
||||
try:
|
||||
self.keylist = self.__DOQ.keys()
|
||||
self.iterator = iter(self.keylist)
|
||||
# get the first available job from any nonempty queue
|
||||
while 1 :
|
||||
tmpsid = self.iterator.next()
|
||||
if not self.__DOQ[tmpsid].empty() :
|
||||
break
|
||||
return tmpsid
|
||||
except StopIteration :
|
||||
# this 2nd exception means there are no users queues at this moment
|
||||
# returning None implies DOQ is empty
|
||||
#log.debug("RoundRobin queue: there are no user queues at this time")
|
||||
self.iterator = None # this will cause the iterator to be recreated next call
|
||||
return None
|
||||
|
||||
|
||||
# Cleanup method. Does a timestamp check to see if the preset minutes
|
||||
# have passed and does a dict cleanup.
|
||||
def __timed_clean_up(self):
|
||||
# Note: Here again, we first create a temp list of keys from DOQ and
|
||||
# then an iterator from the temp list because an iterator from DOQ breaks
|
||||
# when deleting entries.
|
||||
tmp_tstamp = time.time()
|
||||
if ( (tmp_tstamp - self.cleanup_tstamp) >
|
||||
self.cleanup_mininterval ) :
|
||||
tmpkeylist = self.__DOQ.keys()
|
||||
for each in tmpkeylist :
|
||||
if self.__DOQ[each].empty() :
|
||||
del(self.__DOQ[each])
|
||||
log.debug("RoundRobin queue clean up: Removed job queue entry from dictionary for session = %d" % each)
|
||||
self.cleanup_tstamp = tmp_tstamp
|
||||
@@ -221,12 +221,6 @@ new_user_dataset_access_role_default_private = False
|
||||
# running more than one Galaxy server using the same database.
|
||||
#enable_job_recovery = True
|
||||
|
||||
# Job scheduling policy to be used.
|
||||
# module/package name and classname should be in "module:classname" format.
|
||||
# Comment / uncomment the following policies depending upon which is to be used.
|
||||
#job_scheduler_policy = FIFO
|
||||
job_scheduler_policy = galaxy.jobs.schedulingpolicy.roundrobin:UserRoundRobin
|
||||
|
||||
# Job queue cleanup interval in minutes. Currently only used by RoundRobin
|
||||
job_queue_cleanup_interval = 30
|
||||
|
||||
|
||||
Reference in New Issue
Block a user