From 448e619186095976ca76e3aae3e344580c21766f Mon Sep 17 00:00:00 2001 From: Nate Coraor Date: Fri, 11 Jun 2010 09:36:41 -0400 Subject: [PATCH] Remove the user round robin scheduling policy (it didn't work as intended). The former job_scheduler_policy config option now has no effect. --- lib/galaxy/config.py | 1 - lib/galaxy/jobs/__init__.py | 48 +----- lib/galaxy/jobs/schedulingpolicy/__init__.py | 4 - .../jobs/schedulingpolicy/roundrobin.py | 146 ------------------ universe_wsgi.ini.sample | 6 - 5 files changed, 2 insertions(+), 203 deletions(-) delete mode 100644 lib/galaxy/jobs/schedulingpolicy/__init__.py delete mode 100644 lib/galaxy/jobs/schedulingpolicy/roundrobin.py diff --git a/lib/galaxy/config.py b/lib/galaxy/config.py index 3610f340fb2..a1b6ed21155 100644 --- a/lib/galaxy/config.py +++ b/lib/galaxy/config.py @@ -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 ) diff --git a/lib/galaxy/jobs/__init__.py b/lib/galaxy/jobs/__init__.py index 60003256d75..dfe735ae7f5 100644 --- a/lib/galaxy/jobs/__init__.py +++ b/lib/galaxy/jobs/__init__.py @@ -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 diff --git a/lib/galaxy/jobs/schedulingpolicy/__init__.py b/lib/galaxy/jobs/schedulingpolicy/__init__.py deleted file mode 100644 index e90b370e799..00000000000 --- a/lib/galaxy/jobs/schedulingpolicy/__init__.py +++ /dev/null @@ -1,4 +0,0 @@ -""" -This package contains job scheduling policy classes. - -""" \ No newline at end of file diff --git a/lib/galaxy/jobs/schedulingpolicy/roundrobin.py b/lib/galaxy/jobs/schedulingpolicy/roundrobin.py deleted file mode 100644 index 93a1d7d57b4..00000000000 --- a/lib/galaxy/jobs/schedulingpolicy/roundrobin.py +++ /dev/null @@ -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 diff --git a/universe_wsgi.ini.sample b/universe_wsgi.ini.sample index f4aa53af89f..d342bfc18f9 100644 --- a/universe_wsgi.ini.sample +++ b/universe_wsgi.ini.sample @@ -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