From df0042aced39a2e7863f4c8384eb2255e145a0bf Mon Sep 17 00:00:00 2001 From: John Chilton Date: Thu, 18 Oct 2012 23:32:32 -0500 Subject: [PATCH 1/3] Improved encapsulation of job splitting logic, setting the stage for implicit splitting. --- lib/galaxy/jobs/__init__.py | 9 +++++++++ lib/galaxy/jobs/handler.py | 4 ++-- 2 files changed, 11 insertions(+), 2 deletions(-) diff --git a/lib/galaxy/jobs/__init__.py b/lib/galaxy/jobs/__init__.py index 292d558e6ae..436a818a7e8 100644 --- a/lib/galaxy/jobs/__init__.py +++ b/lib/galaxy/jobs/__init__.py @@ -89,6 +89,10 @@ class JobWrapper( object ): self.__user_system_pwent = None self.__galaxy_system_pwent = None + def can_split( self ): + # Should the job handler split this job up? + return self.app.config.use_tasked_jobs and self.tool.parallelism + def get_job_runner_url( self ): return self.job_runner_mapper.get_job_runner_url( self.params ) @@ -922,6 +926,11 @@ class TaskWrapper(JobWrapper): self.prepare_input_files_cmds = None self.status = task.states.NEW + def can_split( self ): + # Should the job handler split this job up? TaskWrapper should + # always return False as the job has already been split. + return False + def get_job( self ): if self.job_id: return self.sa_session.query( model.Job ).get( self.job_id ) diff --git a/lib/galaxy/jobs/handler.py b/lib/galaxy/jobs/handler.py index 04be9687c57..fddbfc6c2e3 100644 --- a/lib/galaxy/jobs/handler.py +++ b/lib/galaxy/jobs/handler.py @@ -449,7 +449,7 @@ class DefaultJobDispatcher( object ): log.debug( 'Loaded job runner: %s' % display_name ) def __get_runner_name( self, job_wrapper ): - if self.app.config.use_tasked_jobs and job_wrapper.tool.parallelism is not None and not isinstance(job_wrapper, TaskWrapper): + if job_wrapper.can_split(): runner_name = "tasks" else: runner_name = ( job_wrapper.get_job_runner_url().split(":", 1) )[0] @@ -458,7 +458,7 @@ class DefaultJobDispatcher( object ): def put( self, job_wrapper ): try: runner_name = self.__get_runner_name( job_wrapper ) - if self.app.config.use_tasked_jobs and job_wrapper.tool.parallelism is not None and isinstance(job_wrapper, TaskWrapper): + if isinstance(job_wrapper, TaskWrapper): #DBTODO Refactor log.debug( "dispatching task %s, of job %d, to %s runner" %( job_wrapper.task_id, job_wrapper.job_id, runner_name ) ) else: From 2794861511e5c0f997c4438a076f69f24a473246 Mon Sep 17 00:00:00 2001 From: John Chilton Date: Sun, 11 Nov 2012 15:36:14 -0600 Subject: [PATCH 2/3] Replace access pattern 'job_wrapper.tool.parallelism' with 'job_wrapper.get_parallelism()' (a newly implemented method on JobWrapper) as a step toward enabling of per job parallelism (as opposed to per tool parallelism). --- lib/galaxy/jobs/__init__.py | 3 +++ lib/galaxy/jobs/runners/tasks.py | 7 ++++--- lib/galaxy/jobs/splitters/basic.py | 5 +++-- lib/galaxy/jobs/splitters/multi.py | 4 ++-- 4 files changed, 12 insertions(+), 7 deletions(-) diff --git a/lib/galaxy/jobs/__init__.py b/lib/galaxy/jobs/__init__.py index 436a818a7e8..62a31406e9d 100644 --- a/lib/galaxy/jobs/__init__.py +++ b/lib/galaxy/jobs/__init__.py @@ -96,6 +96,9 @@ class JobWrapper( object ): def get_job_runner_url( self ): return self.job_runner_mapper.get_job_runner_url( self.params ) + def get_parallelism(self): + return self.tool.parallelism + # legacy naming get_job_runner = get_job_runner_url diff --git a/lib/galaxy/jobs/runners/tasks.py b/lib/galaxy/jobs/runners/tasks.py index 52372cb2d7f..05829826e7e 100644 --- a/lib/galaxy/jobs/runners/tasks.py +++ b/lib/galaxy/jobs/runners/tasks.py @@ -71,12 +71,13 @@ class TaskedJobRunner( object ): try: job_wrapper.change_state( model.Job.states.RUNNING ) self.sa_session.flush() - # Split with the tool-defined method. + # Split with the defined method. + parallelism = job_wrapper.get_parallelism() try: - splitter = getattr(__import__('galaxy.jobs.splitters', globals(), locals(), [job_wrapper.tool.parallelism.method]), job_wrapper.tool.parallelism.method) + splitter = getattr(__import__('galaxy.jobs.splitters', globals(), locals(), [parallelism.method]), parallelism.method) except: job_wrapper.change_state( model.Job.states.ERROR ) - job_wrapper.fail("Job Splitting Failed, no match for '%s'" % job_wrapper.tool.parallelism) + job_wrapper.fail("Job Splitting Failed, no match for '%s'" % parallelism) return tasks = splitter.do_split(job_wrapper) # Not an option for now. Task objects don't *do* anything diff --git a/lib/galaxy/jobs/splitters/basic.py b/lib/galaxy/jobs/splitters/basic.py index 5ec1bf24dfb..50d7151691e 100644 --- a/lib/galaxy/jobs/splitters/basic.py +++ b/lib/galaxy/jobs/splitters/basic.py @@ -5,8 +5,9 @@ log = logging.getLogger( __name__ ) def set_basic_defaults(job_wrapper): parent_job = job_wrapper.get_job() - job_wrapper.tool.parallelism.attributes['split_inputs'] = parent_job.input_datasets[0].name - job_wrapper.tool.parallelism.attributes['merge_outputs'] = job_wrapper.get_output_hdas_and_fnames().keys()[0] + parallelism = job_wrapper.get_parallelism() + parallelism.attributes['split_inputs'] = parent_job.input_datasets[0].name + parallelism.attributes['merge_outputs'] = job_wrapper.get_output_hdas_and_fnames().keys()[0] def do_split (job_wrapper): if len(job_wrapper.get_input_fnames()) > 1 or len(job_wrapper.get_output_fnames()) > 1: diff --git a/lib/galaxy/jobs/splitters/multi.py b/lib/galaxy/jobs/splitters/multi.py index 95f30b7f3b7..1618ceee21f 100644 --- a/lib/galaxy/jobs/splitters/multi.py +++ b/lib/galaxy/jobs/splitters/multi.py @@ -8,7 +8,7 @@ def do_split (job_wrapper): parent_job = job_wrapper.get_job() working_directory = os.path.abspath(job_wrapper.working_directory) - parallel_settings = job_wrapper.tool.parallelism.attributes + parallel_settings = job_wrapper.get_parallelism().attributes # Syntax: split_inputs="input1,input2" shared_inputs="genome" # Designates inputs to be split or shared split_inputs=parallel_settings.get("split_inputs") @@ -91,7 +91,7 @@ def do_split (job_wrapper): def do_merge( job_wrapper, task_wrappers): - parallel_settings = job_wrapper.tool.parallelism.attributes + parallel_settings = job_wrapper.get_parallelism().attributes # Syntax: merge_outputs="export" pickone_outputs="genomesize" # Designates outputs to be merged, or selected from as a representative merge_outputs = parallel_settings.get("merge_outputs") From e549a284ba08b5d030872ca4f291d6eb338ee0cd Mon Sep 17 00:00:00 2001 From: John Chilton Date: Thu, 15 Nov 2012 09:19:08 -0600 Subject: [PATCH 3/3] Rename ToolParallelismInfo to ParallelismInfo and move to jobs module to reflect the fact per-job (instead of per-tool) parallelism could also be a possibility (even if only in downstream Galaxy forks). Also, allow building these objects from dictionaries (in addition to traditional XML-based creation). --- lib/galaxy/jobs/__init__.py | 17 +++++++++++++++++ lib/galaxy/tools/__init__.py | 16 ++-------------- 2 files changed, 19 insertions(+), 14 deletions(-) diff --git a/lib/galaxy/jobs/__init__.py b/lib/galaxy/jobs/__init__.py index 62a31406e9d..4cdaa071fb7 100644 --- a/lib/galaxy/jobs/__init__.py +++ b/lib/galaxy/jobs/__init__.py @@ -1153,3 +1153,20 @@ class NoopQueue( object ): return def shutdown( self ): return + +class ParallelismInfo(object): + """ + Stores the information (if any) for running multiple instances of the tool in parallel + on the same set of inputs. + """ + def __init__(self, tag): + self.method = tag.get('method') + if isinstance(tag, dict): + items = tag.iteritems() + else: + items = tag.attrib.items() + self.attributes = dict([item for item in items if item[0] != 'method' ]) + if len(self.attributes) == 0: + # legacy basic mode - provide compatible defaults + self.attributes['split_size'] = 20 + self.attributes['split_mode'] = 'number_of_parts' diff --git a/lib/galaxy/tools/__init__.py b/lib/galaxy/tools/__init__.py index 3a33d84a4e3..6922f842a6b 100755 --- a/lib/galaxy/tools/__init__.py +++ b/lib/galaxy/tools/__init__.py @@ -15,6 +15,7 @@ from galaxy.util.odict import odict from galaxy.util.bunch import Bunch from galaxy.util.template import fill_template from galaxy import util, jobs, model +from galaxy.jobs import ParallelismInfo from elementtree import ElementTree from parameters import * from parameters.grouping import * @@ -797,19 +798,6 @@ class ToolRequirement( object ): self.type = type self.version = version -class ToolParallelismInfo(object): - """ - Stores the information (if any) for running multiple instances of the tool in parallel - on the same set of inputs. - """ - def __init__(self, tag): - self.method = tag.get('method') - self.attributes = dict([item for item in tag.attrib.items() if item[0] != 'method' ]) - if len(self.attributes) == 0: - # legacy basic mode - provide compatible defaults - self.attributes['split_size'] = 20 - self.attributes['split_mode'] = 'number_of_parts' - class Tool: """ Represents a computational tool that can be executed through Galaxy. @@ -989,7 +977,7 @@ class Tool: # Parallelism for tasks, read from tool config. parallelism = root.find("parallelism") if parallelism is not None and parallelism.get("method"): - self.parallelism = ToolParallelismInfo(parallelism) + self.parallelism = ParallelismInfo(parallelism) else: self.parallelism = None # Set job handler(s). Each handler is a dict with 'url' and, optionally, 'params'.