Parallelize executing tool parameter combinations.

Should work when mapping over collections or for big muli-run tool submissions.

Because of database tension with sqlalchemy it is not strictly a linear increase, but the end user walltime experience for a 24 dataset collection being submitted using 4 threads instead 1 drops execution time from 68 seconds to 35.

Rebased with fixes thanks to @nsoranzo - https://github.com/jmchilton/galaxy/commit/7f6514a21222ea787a0676764e4a4fe825891b27#commitcomment-14712999.
This commit is contained in:
John Chilton
2015-12-03 10:16:12 +00:00
parent 876c28f88d
commit 41d3dc62f7
7 changed files with 72 additions and 9 deletions
+9
View File
@@ -832,6 +832,15 @@ use_interactive = True
#enable_beta_tool_command_isolation = False
# Set the following to a number of threads greater than 1 to spawn
# a Python task queue for dealing with large tool submissions (either
# through the tool form or as part of an individual workflow step across
# large collection). The size of a "large" tool request is controlled by
# the second parameter below and defaults to 10. This affects workflow
# scheduling and web processes, not job handlers.
#tool_submission_burst_threads = 1
#tool_submission_burst_at = 10
# Enable beta workflow modules that should not yet be considered part of Galaxy's
# stable API.
#enable_beta_workflow_modules = False
+2
View File
@@ -215,6 +215,8 @@ class Configuration( object ):
self.use_tasked_jobs = string_as_bool( kwargs.get( 'use_tasked_jobs', False ) )
self.local_task_queue_workers = int(kwargs.get("local_task_queue_workers", 2))
self.commands_in_new_shell = string_as_bool( kwargs.get( 'enable_beta_tool_command_isolation', "False" ) )
self.tool_submission_burst_threads = int( kwargs.get( 'tool_submission_burst_threads', '1' ) )
self.tool_submission_burst_at = int( kwargs.get( 'tool_submission_burst_at', '10' ) )
# The transfer manager and deferred job queue
self.enable_beta_job_managers = string_as_bool( kwargs.get( 'enable_beta_job_managers', 'False' ) )
# These workflow modules should not be considered part of Galaxy's
+8 -1
View File
@@ -77,6 +77,8 @@ class DatasetCollectionManager( object ):
for input_name, input_collection in implicit_collection_info[ "implicit_inputs" ]:
dataset_collection_instance.add_implicit_input_collection( input_name, input_collection )
for output_dataset in implicit_collection_info.get( "outputs" ):
if output_dataset not in trans.sa_session:
output_dataset = trans.sa_session.query( type( output_dataset ) ).get( output_dataset.id )
if isinstance( output_dataset, model.HistoryDatasetAssociation ):
output_dataset.hidden_beneath_collection_instance = dataset_collection_instance
elif isinstance( output_dataset, model.HistoryDatasetCollectionAssociation ):
@@ -265,7 +267,12 @@ class DatasetCollectionManager( object ):
# Previously created collection already found in request, just pass
# through as is.
if "__object__" in element_identifier:
return element_identifier[ "__object__" ]
the_object = element_identifier[ "__object__" ]
if the_object is not None and the_object.id:
context = self.model.context
if the_object not in context:
the_object = context.query( type(the_object) ).get(the_object.id)
return the_object
# dateset_identifier is dict {src=hda|ldda|hdca|new_collection, id=<encoded_id>}
try:
+10 -4
View File
@@ -556,8 +556,11 @@ class Job( object, JobLike, Dictifiable ):
def add_parameter( self, name, value ):
self.parameters.append( JobParameter( name, value ) )
def add_input_dataset( self, name, dataset ):
self.input_datasets.append( JobToInputDatasetAssociation( name, dataset ) )
def add_input_dataset( self, name, dataset=None, dataset_id=None ):
assoc = JobToInputDatasetAssociation( name, dataset )
if dataset is None and dataset_id is not None:
assoc.dataset_id = dataset_id
self.input_datasets.append( assoc )
def add_output_dataset( self, name, dataset ):
self.output_datasets.append( JobToOutputDatasetAssociation( name, dataset ) )
@@ -1484,10 +1487,13 @@ class DefaultQuotaAssociation( Quota, Dictifiable ):
class DatasetPermissions( object ):
def __init__( self, action, dataset, role ):
def __init__( self, action, dataset, role=None, role_id=None ):
self.action = action
self.dataset = dataset
self.role = role
if role is not None:
self.role = role
else:
self.role_id = role_id
class LibraryPermissions( object ):
+2 -1
View File
@@ -898,7 +898,8 @@ class GalaxyRBACAgent( RBACAgent ):
for action, roles in permissions.items():
if isinstance( action, Action ):
action = action.action
for dp in [ self.model.DatasetPermissions( action, dataset, role ) for role in roles ]:
for role in roles:
dp = self.model.DatasetPermissions( action, dataset, role_id=role.id )
self.sa_session.add( dp )
flush_needed = True
if flush_needed:
+6 -1
View File
@@ -162,6 +162,8 @@ class DefaultToolAction( object ):
# Set history.
if not history:
history = tool.get_default_history_by_trans( trans, create=True )
if history not in trans.sa_session:
history = trans.sa_session.query( trans.app.model.History ).get( history.id )
out_data = odict()
out_collections = {}
@@ -425,7 +427,10 @@ class DefaultToolAction( object ):
if dataset:
if not trans.app.security_agent.can_access_dataset( current_user_roles, dataset.dataset ):
raise Exception("User does not have permission to use a dataset (%s) provided for input." % data.id)
job.add_input_dataset( name, dataset )
if dataset in trans.sa_session:
job.add_input_dataset( name, dataset=dataset )
else:
job.add_input_dataset( name, dataset_id=dataset.id )
else:
job.add_input_dataset( name, None )
log.info("Verified access to datasets %s" % access_timer)
+35 -2
View File
@@ -7,6 +7,8 @@ import collections
import galaxy.tools
from galaxy.util import ExecutionTimer
from galaxy.tools.actions import on_text_for_names
from threading import Thread
from Queue import Queue
import logging
log = logging.getLogger( __name__ )
@@ -20,8 +22,9 @@ def execute( trans, tool, param_combinations, history, rerun_remap_job_id=None,
failures, etc...).
"""
execution_tracker = ToolExecutionTracker( tool, param_combinations, collection_info )
all_jobs_timer = ExecutionTimer()
for params in execution_tracker.param_combinations:
app = trans.app
def execute_single_job(params):
job_timer = ExecutionTimer()
if workflow_invocation_uuid:
params[ '__workflow_invocation_uuid__' ] = workflow_invocation_uuid
@@ -36,6 +39,34 @@ def execute( trans, tool, param_combinations, history, rerun_remap_job_id=None,
execution_tracker.record_success( job, result )
else:
execution_tracker.record_error( result )
all_jobs_timer = ExecutionTimer()
config = app.config
burst_at = getattr( config, 'tool_submission_burst_at', 10 )
burst_threads = getattr( config, 'tool_submission_burst_threads', 1 )
if len(execution_tracker.param_combinations) < burst_at or burst_threads < 2:
for params in execution_tracker.param_combinations:
execute_single_job(params)
else:
q = Queue()
def worker():
while True:
params = q.get()
execute_single_job(params)
q.task_done()
for i in range(burst_threads):
t = Thread(target=worker)
t.daemon = True
t.start()
for params in execution_tracker.param_combinations:
q.put(params)
q.join()
log.debug("Executed all jobs for tool request: %s" % all_jobs_timer)
if collection_info:
history = history or tool.get_default_history_by_trans( trans )
@@ -140,6 +171,8 @@ class ToolExecutionTracker( object ):
# TODO: Think through this, may only want this for output
# collections - or we may be already recording data in some
# other way.
if job not in trans.sa_session:
job = trans.sa_session.query( trans.app.model.Job ).get( job.id )
job.add_output_dataset_collection( output_name, collection )
collections[ output_name ] = collection