From 41d3dc62f7e967ea3b409be410310563d0737bc4 Mon Sep 17 00:00:00 2001 From: John Chilton Date: Mon, 30 Nov 2015 13:45:22 +0000 Subject: [PATCH] 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. --- config/galaxy.ini.sample | 9 +++++++ lib/galaxy/config.py | 2 ++ lib/galaxy/managers/collections.py | 9 ++++++- lib/galaxy/model/__init__.py | 14 ++++++++--- lib/galaxy/security/__init__.py | 3 ++- lib/galaxy/tools/actions/__init__.py | 7 +++++- lib/galaxy/tools/execute.py | 37 ++++++++++++++++++++++++++-- 7 files changed, 72 insertions(+), 9 deletions(-) diff --git a/config/galaxy.ini.sample b/config/galaxy.ini.sample index f6fb9e4ef55..322ea506346 100644 --- a/config/galaxy.ini.sample +++ b/config/galaxy.ini.sample @@ -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 diff --git a/lib/galaxy/config.py b/lib/galaxy/config.py index 95457d21b25..dd3b16231ba 100644 --- a/lib/galaxy/config.py +++ b/lib/galaxy/config.py @@ -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 diff --git a/lib/galaxy/managers/collections.py b/lib/galaxy/managers/collections.py index ded32d1ba84..24dbdff1c13 100644 --- a/lib/galaxy/managers/collections.py +++ b/lib/galaxy/managers/collections.py @@ -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=} try: diff --git a/lib/galaxy/model/__init__.py b/lib/galaxy/model/__init__.py index e45702a41dd..be008e24f39 100644 --- a/lib/galaxy/model/__init__.py +++ b/lib/galaxy/model/__init__.py @@ -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 ): diff --git a/lib/galaxy/security/__init__.py b/lib/galaxy/security/__init__.py index 685bf3db20f..44ced674516 100644 --- a/lib/galaxy/security/__init__.py +++ b/lib/galaxy/security/__init__.py @@ -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: diff --git a/lib/galaxy/tools/actions/__init__.py b/lib/galaxy/tools/actions/__init__.py index bce476fe545..0408ec7c8ff 100644 --- a/lib/galaxy/tools/actions/__init__.py +++ b/lib/galaxy/tools/actions/__init__.py @@ -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) diff --git a/lib/galaxy/tools/execute.py b/lib/galaxy/tools/execute.py index 8f167c7fc4d..6b1512c2617 100644 --- a/lib/galaxy/tools/execute.py +++ b/lib/galaxy/tools/execute.py @@ -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