mirror of
https://github.com/galaxyproject/galaxy.git
synced 2026-09-24 16:30:27 +08:00
Merge with latest galaxy-central to resolve conflict introduced with 4bd4197.
This commit is contained in:
@@ -89,9 +89,16 @@ 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 )
|
||||
|
||||
def get_parallelism(self):
|
||||
return self.tool.parallelism
|
||||
|
||||
# legacy naming
|
||||
get_job_runner = get_job_runner_url
|
||||
|
||||
@@ -941,6 +948,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 )
|
||||
@@ -1160,3 +1172,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'
|
||||
|
||||
@@ -460,7 +460,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]
|
||||
@@ -478,7 +478,7 @@ class DefaultJobDispatcher( object ):
|
||||
job_wrapper.fail( failure_message )
|
||||
return
|
||||
try:
|
||||
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:
|
||||
|
||||
@@ -72,12 +72,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
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -9,7 +9,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")
|
||||
@@ -92,7 +92,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")
|
||||
|
||||
@@ -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 *
|
||||
@@ -800,19 +801,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( object ):
|
||||
"""
|
||||
Represents a computational tool that can be executed through Galaxy.
|
||||
@@ -992,7 +980,7 @@ class Tool( object ):
|
||||
# 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'.
|
||||
|
||||
Reference in New Issue
Block a user