From 4111be6a774e1329dd9a7292d2c784fe83fc6570 Mon Sep 17 00:00:00 2001 From: John Chilton Date: Wed, 1 Aug 2018 15:19:57 -0400 Subject: [PATCH 01/12] Flush all collection mapping jobs at once. For a simple example of 10 jobs each creating 10 datasets on my laptop the execution time for everything went from almost 5 seconds to 1.5 seconds. --- lib/galaxy/model/__init__.py | 2 ++ lib/galaxy/tools/__init__.py | 3 ++- lib/galaxy/tools/actions/__init__.py | 13 +++++++++---- lib/galaxy/tools/execute.py | 9 ++++++++- lib/galaxy/workflow/modules.py | 9 ++++----- test/unit/tools/test_execution.py | 3 +++ 6 files changed, 28 insertions(+), 11 deletions(-) diff --git a/lib/galaxy/model/__init__.py b/lib/galaxy/model/__init__.py index a5ab245a5a9..ec1905d416f 100644 --- a/lib/galaxy/model/__init__.py +++ b/lib/galaxy/model/__init__.py @@ -5300,6 +5300,8 @@ class WorkflowInvocation(UsesCreateAndUpdateTime, Dictifiable, RepresentById): for step in self.steps: if step.workflow_step.type == 'tool': for job in step.jobs: + if job is None: + continue for step_input in step.workflow_step.input_connections: output_step_type = step_input.output_step.type if output_step_type in ['data_input', 'data_collection_input']: diff --git a/lib/galaxy/tools/__init__.py b/lib/galaxy/tools/__init__.py index d3e697a3157..968b9efcd12 100755 --- a/lib/galaxy/tools/__init__.py +++ b/lib/galaxy/tools/__init__.py @@ -1584,7 +1584,7 @@ class Tool(Dictifiable): output_collections=execution_tracker.output_collections, implicit_collections=execution_tracker.implicit_collections) - def handle_single_execution(self, trans, rerun_remap_job_id, execution_slice, history, execution_cache=None, completed_job=None, collection_info=None): + def handle_single_execution(self, trans, rerun_remap_job_id, execution_slice, history, execution_cache=None, completed_job=None, collection_info=None, flush_job=True): """ Return a pair with whether execution is successful as well as either resulting output data or an error message indicating the problem. @@ -1599,6 +1599,7 @@ class Tool(Dictifiable): dataset_collection_elements=execution_slice.dataset_collection_elements, completed_job=completed_job, collection_info=collection_info, + flush_job=flush_job, ) except (webob.exc.HTTPFound, exceptions.MessageException) as e: # if it's a webob redirect exception, pass it up the stack diff --git a/lib/galaxy/tools/actions/__init__.py b/lib/galaxy/tools/actions/__init__.py index 2bd342f8c25..8d4f5385a29 100644 --- a/lib/galaxy/tools/actions/__init__.py +++ b/lib/galaxy/tools/actions/__init__.py @@ -276,7 +276,7 @@ class DefaultToolAction(object): preserved_tags[tag.value] = tag return history, inp_data, inp_dataset_collections, preserved_tags, all_permissions - def execute(self, tool, trans, incoming=None, return_job=False, set_output_hid=True, history=None, job_params=None, rerun_remap_job_id=None, execution_cache=None, dataset_collection_elements=None, completed_job=None, collection_info=None): + def execute(self, tool, trans, incoming=None, return_job=False, set_output_hid=True, history=None, job_params=None, rerun_remap_job_id=None, execution_cache=None, dataset_collection_elements=None, completed_job=None, collection_info=None, flush_job=True): """ Executes a tool, creating job and tool outputs, associating them, and submitting the job to the job queue. If history is not specified, use @@ -578,9 +578,14 @@ class DefaultToolAction(object): trans.sa_session.flush() trans.response.send_redirect(url_for(controller='tool_runner', action='redirect', redirect_url=redirect_url)) else: - # Dispatch to a job handler. enqueue() is responsible for flushing the job - app.job_manager.enqueue(job, tool=tool) - trans.log_event("Added job to the job queue, id: %s" % str(job.id), tool_id=job.tool_id) + if flush_job: + job_flush_timer = ExecutionTimer() + trans.sa_session.flush() + log.info("Flushed transaction for job %s %s" % (job.log_str(), job_flush_timer)) + + # Dispatch to a job handler. enqueue() is responsible for flushing the job + app.job_manager.enqueue(job, tool=tool) + trans.log_event("Added job to the job queue, id: %s" % str(job.id), tool_id=job.tool_id) return job, out_data def _remap_job_on_rerun(self, trans, galaxy_session, rerun_remap_job_id, current_job, out_data): diff --git a/lib/galaxy/tools/execute.py b/lib/galaxy/tools/execute.py index b746f224647..013453d2383 100644 --- a/lib/galaxy/tools/execute.py +++ b/lib/galaxy/tools/execute.py @@ -68,7 +68,7 @@ def execute(trans, tool, mapping_params, history, rerun_remap_job_id=None, colle del params['__workflow_resource_params__'] if validate_outputs: params['__validate_outputs__'] = True - job, result = tool.handle_single_execution(trans, rerun_remap_job_id, execution_slice, history, execution_cache, completed_job, collection_info) + job, result = tool.handle_single_execution(trans, rerun_remap_job_id, execution_slice, history, execution_cache, completed_job, collection_info, flush_job=False) if job: log.debug(job_timer.to_str(tool_id=tool.id, job_id=job.id)) execution_tracker.record_success(execution_slice, job, result) @@ -101,6 +101,13 @@ def execute(trans, tool, mapping_params, history, rerun_remap_job_id=None, colle else: execute_single_job(execution_slice, completed_jobs[i]) + full_flush_timer = ExecutionTimer() + trans.sa_session.flush() + for job in execution_tracker.successful_jobs: + # Put the job in the queue if tracking in memory + app.job_manager.enqueue(job, tool=tool) + trans.log_event("Added job to the job queue, id: %s" % str(job.id), tool_id=job.tool_id) + if has_remaining_jobs: raise PartialJobExecution(execution_tracker) else: diff --git a/lib/galaxy/workflow/modules.py b/lib/galaxy/workflow/modules.py index ce02662c64f..a96f47f20f8 100644 --- a/lib/galaxy/workflow/modules.py +++ b/lib/galaxy/workflow/modules.py @@ -1802,17 +1802,16 @@ class ToolModule(WorkflowModule): # Combine workflow and runtime post job actions into the effective post # job actions for this execution. - flush_required = False effective_post_job_actions = self._effective_post_job_actions(step) for pja in effective_post_job_actions: if pja.action_type in ActionBox.immediate_actions or isinstance(self.tool, DatabaseOperationTool): ActionBox.execute(self.trans.app, self.trans.sa_session, pja, job, replacement_dict) else: - pjaa = model.PostJobActionAssociation(pja, job_id=job.id) + if job.id: + pjaa = model.PostJobActionAssociation(pja, job_id=job.id) + else: + pjaa = model.PostJobActionAssociation(pja, job=job) self.trans.sa_session.add(pjaa) - flush_required = True - if flush_required: - self.trans.sa_session.flush() def __restore_step_meta_runtime_state(self, step_runtime_state): if RUNTIME_POST_JOB_ACTIONS_KEY in step_runtime_state: diff --git a/test/unit/tools/test_execution.py b/test/unit/tools/test_execution.py index fc547af9375..83185a15510 100644 --- a/test/unit/tools/test_execution.py +++ b/test/unit/tools/test_execution.py @@ -211,6 +211,9 @@ class MockTrans(object): def get_current_user_roles(self): return [] + def log_event(self, *args, **kwds): + pass + class MockCollectionService(object): From 29584f8ef275d1d1674e33dfa3c867b906de30f0 Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Sun, 12 Jul 2020 18:46:31 +0200 Subject: [PATCH 02/12] Fix John's commit for current Galaxy --- lib/galaxy/tools/actions/__init__.py | 2 +- lib/galaxy/tools/execute.py | 3 +-- 2 files changed, 2 insertions(+), 3 deletions(-) diff --git a/lib/galaxy/tools/actions/__init__.py b/lib/galaxy/tools/actions/__init__.py index 8d4f5385a29..3dde44e74e3 100644 --- a/lib/galaxy/tools/actions/__init__.py +++ b/lib/galaxy/tools/actions/__init__.py @@ -408,7 +408,7 @@ class DefaultToolAction(object): dataset_collection_elements[name].hda = data trans.sa_session.add(data) if not completed_job: - trans.app.security_agent.set_all_dataset_permissions(data.dataset, output_permissions, new=True) + trans.app.security_agent.set_all_dataset_permissions(data.dataset, output_permissions, new=True, flush=False) data.copy_tags_to(preserved_tags) if not completed_job and trans.app.config.legacy_eager_objectstore_initialization: diff --git a/lib/galaxy/tools/execute.py b/lib/galaxy/tools/execute.py index 013453d2383..f3f7c813e52 100644 --- a/lib/galaxy/tools/execute.py +++ b/lib/galaxy/tools/execute.py @@ -101,11 +101,10 @@ def execute(trans, tool, mapping_params, history, rerun_remap_job_id=None, colle else: execute_single_job(execution_slice, completed_jobs[i]) - full_flush_timer = ExecutionTimer() trans.sa_session.flush() for job in execution_tracker.successful_jobs: # Put the job in the queue if tracking in memory - app.job_manager.enqueue(job, tool=tool) + tool.app.job_manager.enqueue(job, tool=tool) trans.log_event("Added job to the job queue, id: %s" % str(job.id), tool_id=job.tool_id) if has_remaining_jobs: From bf6481193577da46681168b4bfd1b49fc4440e38 Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Sun, 12 Jul 2020 19:11:54 +0200 Subject: [PATCH 03/12] Move addition of datasets to history out of single tool execution At this point mapping `Concatenate datasets (with sleep)` over a collection of 10 datasets has this profile: ``` *** PROFILER RESULTS *** execute (/Users/mvandenb/src/galaxy/lib/galaxy/tools/execute.py:33) function called 1 times r 184c62007f Move addition of datasets to history out of single tool execution Move addition of datasets to history out of single tool execution 362378 function calls (353366 primitive calls) in 0.474 seconds Ordered by: cumulative time, internal time, call count List reduced from 1622 to 40 due to restriction <40> ncalls tottime percall cumtime percall filename:lineno(function) 1 0.000 0.000 0.475 0.475 execute.py:33(execute) 15 0.001 0.000 0.225 0.015 session.py:2489(flush) 311 0.001 0.000 0.225 0.001 base.py:952(execute) 14 0.001 0.000 0.224 0.016 session.py:2542(_flush) 311 0.000 0.000 0.224 0.001 elements.py:296(_execute_on_connection) 311 0.002 0.000 0.223 0.001 base.py:1088(_execute_clauseelement) 2724/2624 0.002 0.000 0.206 0.000 attributes.py:699(get) 1297/1292 0.001 0.000 0.203 0.000 attributes.py:279(__get__) 311 0.003 0.000 0.160 0.001 base.py:1195(_execute_context) 49 0.000 0.000 0.155 0.003 query.py:3417(one) 49 0.001 0.000 0.155 0.003 query.py:3381(one_or_none) 14 0.000 0.000 0.151 0.011 unitofwork.py:402(execute) 298 0.000 0.000 0.144 0.000 state.py:640(_load_expired) 48 0.000 0.000 0.143 0.003 loading.py:938(load_scalar_attributes) 48 0.000 0.000 0.141 0.003 loading.py:190(load_on_ident) 48 0.000 0.000 0.141 0.003 loading.py:211(load_on_pk_identity) 109 0.000 0.000 0.135 0.001 query.py:3501(_execute_and_instances) 10 0.000 0.000 0.132 0.013 manager.py:52(enqueue) 2 0.000 0.000 0.131 0.065 mapping.py:2827(db_next_hid) 1 0.000 0.000 0.127 0.127 __init__.py:1790(add_datasets) 1 0.000 0.000 0.126 0.126 __init__.py:1811(__add_datasets_optimized) 50 0.000 0.000 0.126 0.003 query.py:3476(__iter__) 16/15 0.000 0.000 0.124 0.008 session.py:899(begin) 16/15 0.000 0.000 0.123 0.008 session.py:221(__init__) 16/15 0.000 0.000 0.123 0.008 session.py:338(_take_snapshot) 308 0.000 0.000 0.116 0.000 default.py:589(do_execute) 308 0.114 0.000 0.116 0.000 {method 'execute' of 'psycopg2.extensions.cursor' objects} 10 0.000 0.000 0.110 0.011 execute.py:54(execute_single_job) 39 0.001 0.000 0.110 0.003 persistence.py:184(save_obj) 10 0.000 0.000 0.104 0.010 __init__.py:1587(handle_single_execution) 10 0.000 0.000 0.104 0.010 __init__.py:1680(execute) 10 0.001 0.000 0.104 0.010 __init__.py:279(execute) 37 0.000 0.000 0.095 0.003 unitofwork.py:585(execute) 10 0.000 0.000 0.091 0.009 handlers.py:423(assign_handler) 10 0.000 0.000 0.089 0.009 handlers.py:326(_assign_db_self_handler) 10 0.000 0.000 0.089 0.009 handlers.py:463(_timed_flush_obj) 122 0.000 0.000 0.083 0.001 unitofwork.py:520(execute_aggregate) 39 0.002 0.000 0.078 0.002 persistence.py:1039(_emit_insert_statements) 450 0.001 0.000 0.073 0.000 strategies.py:665(_load_for_state) 59 0.000 0.000 0.069 0.001 scoping.py:162(do) ``` --- lib/galaxy/tools/__init__.py | 6 +++++- lib/galaxy/tools/actions/__init__.py | 8 ++++---- lib/galaxy/tools/execute.py | 6 ++++++ 3 files changed, 15 insertions(+), 5 deletions(-) diff --git a/lib/galaxy/tools/__init__.py b/lib/galaxy/tools/__init__.py index 968b9efcd12..ac9daece1dd 100755 --- a/lib/galaxy/tools/__init__.py +++ b/lib/galaxy/tools/__init__.py @@ -1590,7 +1590,7 @@ class Tool(Dictifiable): resulting output data or an error message indicating the problem. """ try: - job, out_data = self.execute( + rval = self.execute( trans, incoming=execution_slice.param_combination, history=history, @@ -1601,6 +1601,10 @@ class Tool(Dictifiable): collection_info=collection_info, flush_job=flush_job, ) + job = rval[0] + out_data = rval[1] + if len(rval) == 3: + execution_slice.datasets_to_persist = rval[2] except (webob.exc.HTTPFound, exceptions.MessageException) as e: # if it's a webob redirect exception, pass it up the stack raise e diff --git a/lib/galaxy/tools/actions/__init__.py b/lib/galaxy/tools/actions/__init__.py index 3dde44e74e3..1b9c399057d 100644 --- a/lib/galaxy/tools/actions/__init__.py +++ b/lib/galaxy/tools/actions/__init__.py @@ -525,9 +525,6 @@ class DefaultToolAction(object): for name, data in out_data.items(): if name not in child_dataset_names and name not in incoming: # don't add children; or already existing datasets, i.e. async created datasets_to_persist.append(data) - # Set HID and add to history. - # This is brand new and certainly empty so don't worry about quota. - history.add_datasets(trans.sa_session, datasets_to_persist, set_hid=set_output_hid, quota=False, flush=False) # Add all the children to their parents for parent_name, child_name in parent_to_child_pairs: @@ -579,6 +576,9 @@ class DefaultToolAction(object): trans.response.send_redirect(url_for(controller='tool_runner', action='redirect', redirect_url=redirect_url)) else: if flush_job: + # Set HID and add to history. + # This is brand new and certainly empty so don't worry about quota. + history.add_datasets(trans.sa_session, datasets_to_persist, set_hid=set_output_hid, quota=False, flush=False) job_flush_timer = ExecutionTimer() trans.sa_session.flush() log.info("Flushed transaction for job %s %s" % (job.log_str(), job_flush_timer)) @@ -586,7 +586,7 @@ class DefaultToolAction(object): # Dispatch to a job handler. enqueue() is responsible for flushing the job app.job_manager.enqueue(job, tool=tool) trans.log_event("Added job to the job queue, id: %s" % str(job.id), tool_id=job.tool_id) - return job, out_data + return job, out_data, datasets_to_persist def _remap_job_on_rerun(self, trans, galaxy_session, rerun_remap_job_id, current_job, out_data): """ diff --git a/lib/galaxy/tools/execute.py b/lib/galaxy/tools/execute.py index f3f7c813e52..98384e5b5ae 100644 --- a/lib/galaxy/tools/execute.py +++ b/lib/galaxy/tools/execute.py @@ -93,6 +93,7 @@ def execute(trans, tool, mapping_params, history, rerun_remap_job_id=None, colle jobs_executed = 0 has_remaining_jobs = False + datasets_to_persist = [] for i, execution_slice in enumerate(execution_tracker.new_execution_slices()): if max_num_jobs and jobs_executed >= max_num_jobs: @@ -100,7 +101,11 @@ def execute(trans, tool, mapping_params, history, rerun_remap_job_id=None, colle break else: execute_single_job(execution_slice, completed_jobs[i]) + if execution_slice.datasets_to_persist: + datasets_to_persist.extend(execution_slice.datasets_to_persist) + if datasets_to_persist: + history.add_datasets(trans.sa_session, datasets_to_persist, set_hid=True, quota=False, flush=False) trans.sa_session.flush() for job in execution_tracker.successful_jobs: # Put the job in the queue if tracking in memory @@ -122,6 +127,7 @@ class ExecutionSlice(object): self.job_index = job_index self.param_combination = param_combination self.dataset_collection_elements = dataset_collection_elements + self.datasets_to_persist = None class ExecutionTracker(object): From e4ebef313c006e2083c933bafed5a62f418928ab Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Sun, 12 Jul 2020 19:31:03 +0200 Subject: [PATCH 04/12] Unify created_element_datasets and datasets_to_persist I think that might speed up creation of paired datasets --- lib/galaxy/tools/actions/__init__.py | 6 ++---- 1 file changed, 2 insertions(+), 4 deletions(-) diff --git a/lib/galaxy/tools/actions/__init__.py b/lib/galaxy/tools/actions/__init__.py index 1b9c399057d..d37c518e294 100644 --- a/lib/galaxy/tools/actions/__init__.py +++ b/lib/galaxy/tools/actions/__init__.py @@ -454,6 +454,7 @@ class DefaultToolAction(object): # Flush all datasets at once. return data + datasets_to_persist = [] for name, output in tool.outputs.items(): if not filter_output(output, incoming): handle_output_timer = ExecutionTimer() @@ -461,7 +462,6 @@ class DefaultToolAction(object): collections_manager = app.dataset_collections_service element_identifiers = [] known_outputs = output.known_outputs(input_collections, collections_manager.type_registry) - created_element_datasets = [] # Just to echo TODO elsewhere - this should be restructured to allow # nested collections. for output_part_def in known_outputs: @@ -488,7 +488,7 @@ class DefaultToolAction(object): effective_output_name = output_part_def.effective_output_name element = handle_output(effective_output_name, output_part_def.output_def, hidden=True) - created_element_datasets.append(element) + datasets_to_persist.append(element) # TODO: this shouldn't exist in the top-level of the history at all # but for now we are still working around that by hiding the contents # there. @@ -500,7 +500,6 @@ class DefaultToolAction(object): "name": output_part_def.element_identifier, }) - history.add_datasets(trans.sa_session, created_element_datasets, set_hid=set_output_hid, quota=False, flush=True) if output.dynamic_structure: assert not element_identifiers # known_outputs must have been empty element_kwds = dict(elements=collections_manager.ELEMENTS_UNINITIALIZED) @@ -521,7 +520,6 @@ class DefaultToolAction(object): 'Added output datasets to history', ) # Add all the top-level (non-child) datasets to the history unless otherwise specified - datasets_to_persist = [] for name, data in out_data.items(): if name not in child_dataset_names and name not in incoming: # don't add children; or already existing datasets, i.e. async created datasets_to_persist.append(data) From ce1ac452b21c9b2aafb18042c2a12f700eba5502 Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Sun, 12 Jul 2020 21:48:29 +0200 Subject: [PATCH 05/12] Since we may not have a job_id yet we need to use job.id --- lib/galaxy/tools/execute.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/lib/galaxy/tools/execute.py b/lib/galaxy/tools/execute.py index 98384e5b5ae..40d3f0ee025 100644 --- a/lib/galaxy/tools/execute.py +++ b/lib/galaxy/tools/execute.py @@ -373,7 +373,7 @@ class ExecutionTracker(object): job_assoc = model.ImplicitCollectionJobsJobAssociation() job_assoc.order_index = execution_slice.job_index job_assoc.implicit_collection_jobs = implicit_collection_jobs - job_assoc.job_id = job.id + job_assoc.job = job self.trans.sa_session.add(job_assoc) From 37962fa58072bf8742ea267e5ad68c8f07471165 Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Sun, 12 Jul 2020 23:12:57 +0200 Subject: [PATCH 06/12] Some docs on why we flush here --- lib/galaxy/tools/execute.py | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/lib/galaxy/tools/execute.py b/lib/galaxy/tools/execute.py index 40d3f0ee025..9e409cbd11e 100644 --- a/lib/galaxy/tools/execute.py +++ b/lib/galaxy/tools/execute.py @@ -106,7 +106,10 @@ def execute(trans, tool, mapping_params, history, rerun_remap_job_id=None, colle if datasets_to_persist: history.add_datasets(trans.sa_session, datasets_to_persist, set_hid=True, quota=False, flush=False) - trans.sa_session.flush() + # a side effect of history.add_datasets is a commit within db_next_hid (even with flush=False). + else: + # Make sure collections, implicit jobs etc are flushed even if there are no precreated output datasets + trans.sa_session.flush() for job in execution_tracker.successful_jobs: # Put the job in the queue if tracking in memory tool.app.job_manager.enqueue(job, tool=tool) From ef70994113d1ba146b19f11b3979d67c1b42bcc0 Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Mon, 13 Jul 2020 10:38:58 +0200 Subject: [PATCH 07/12] Don't flush once per mapped over output collection --- lib/galaxy/managers/collections.py | 12 +++++++----- lib/galaxy/model/__init__.py | 16 ++++++++-------- lib/galaxy/tools/__init__.py | 3 ++- lib/galaxy/tools/actions/__init__.py | 13 ++++++++++--- lib/galaxy/tools/execute.py | 3 ++- test/unit/tools/test_actions.py | 3 ++- 6 files changed, 31 insertions(+), 19 deletions(-) diff --git a/lib/galaxy/managers/collections.py b/lib/galaxy/managers/collections.py index 59eca934806..461252d1994 100644 --- a/lib/galaxy/managers/collections.py +++ b/lib/galaxy/managers/collections.py @@ -91,7 +91,8 @@ class DatasetCollectionManager(object): def create(self, trans, parent, name, collection_type, element_identifiers=None, elements=None, implicit_collection_info=None, trusted_identifiers=None, - hide_source_items=False, tags=None, copy_elements=False, history=None): + hide_source_items=False, tags=None, copy_elements=False, history=None, + set_hid=True, flush=True): """ PRECONDITION: security checks on ability to add to parent occurred during load. @@ -122,10 +123,10 @@ class DatasetCollectionManager(object): implicit_output_name = implicit_collection_info["implicit_output_name"] return self._create_instance_for_collection( - trans, parent, name, dataset_collection, implicit_inputs=implicit_inputs, implicit_output_name=implicit_output_name, tags=tags + trans, parent, name, dataset_collection, implicit_inputs=implicit_inputs, implicit_output_name=implicit_output_name, tags=tags, set_hid=set_hid, flush=flush, ) - def _create_instance_for_collection(self, trans, parent, name, dataset_collection, implicit_output_name=None, implicit_inputs=None, tags=None, flush=True): + def _create_instance_for_collection(self, trans, parent, name, dataset_collection, implicit_output_name=None, implicit_inputs=None, tags=None, set_hid=True, flush=True): if isinstance(parent, model.History): dataset_collection_instance = self.model.HistoryDatasetCollectionAssociation( collection=dataset_collection, @@ -139,8 +140,9 @@ class DatasetCollectionManager(object): dataset_collection_instance.implicit_output_name = implicit_output_name log.debug("Created collection with %d elements" % (len(dataset_collection_instance.collection.elements))) - # Handle setting hid - parent.add_dataset_collection(dataset_collection_instance) + + if set_hid: + parent.add_dataset_collection(dataset_collection_instance) elif isinstance(parent, model.LibraryFolder): dataset_collection_instance = self.model.LibraryDatasetCollectionAssociation( diff --git a/lib/galaxy/model/__init__.py b/lib/galaxy/model/__init__.py index ec1905d416f..f274f37229b 100644 --- a/lib/galaxy/model/__init__.py +++ b/lib/galaxy/model/__init__.py @@ -1765,9 +1765,10 @@ class History(HasTags, Dictifiable, UsesAnnotations, HasName, RepresentById): dataset = HistoryDatasetAssociation(dataset=dataset) object_session(self).add(dataset) object_session(self).flush() - elif not isinstance(dataset, HistoryDatasetAssociation): + elif not isinstance(dataset, (HistoryDatasetAssociation, HistoryDatasetCollectionAssociation)): raise TypeError("You can only add Dataset and HistoryDatasetAssociation instances to a history" + " ( you tried to add %s )." % str(dataset)) + is_dataset = is_hda(dataset) if parent_id: for data in self.datasets: if data.id == parent_id: @@ -1779,24 +1780,23 @@ class History(HasTags, Dictifiable, UsesAnnotations, HasName, RepresentById): else: if set_hid: dataset.hid = self._next_hid() - if quota and self.user: + if quota and is_dataset and self.user: self.user.adjust_total_disk_usage(dataset.quota_amount(self.user)) dataset.history = self - if genome_build not in [None, '?']: + if is_dataset and genome_build not in [None, '?']: self.genome_build = genome_build dataset.history_id = self.id return dataset def add_datasets(self, sa_session, datasets, parent_id=None, genome_build=None, set_hid=True, quota=True, flush=False): """ Optimized version of add_dataset above that minimizes database - interactions when adding many datasets to history at once. + interactions when adding many datasets and collections to history at once. """ - all_hdas = all(is_hda(_) for _ in datasets) - optimize = len(datasets) > 1 and parent_id is None and all_hdas and set_hid + optimize = len(datasets) > 1 and parent_id is None and set_hid if optimize: self.__add_datasets_optimized(datasets, genome_build=genome_build) if quota and self.user: - disk_usage = sum([d.get_total_size() for d in datasets]) + disk_usage = sum([d.get_total_size() for d in datasets if is_hda(d)]) self.user.adjust_total_disk_usage(disk_usage) sa_session.add_all(datasets) if flush: @@ -1821,7 +1821,7 @@ class History(HasTags, Dictifiable, UsesAnnotations, HasName, RepresentById): dataset.hid = base_hid + i dataset.history = self dataset.history_id = cached_id(self) - if set_genome: + if set_genome and is_hda(dataset): self.genome_build = genome_build return datasets diff --git a/lib/galaxy/tools/__init__.py b/lib/galaxy/tools/__init__.py index ac9daece1dd..128193cbe1e 100755 --- a/lib/galaxy/tools/__init__.py +++ b/lib/galaxy/tools/__init__.py @@ -1603,8 +1603,9 @@ class Tool(Dictifiable): ) job = rval[0] out_data = rval[1] - if len(rval) == 3: + if len(rval) == 4: execution_slice.datasets_to_persist = rval[2] + execution_slice.history = rval[3] except (webob.exc.HTTPFound, exceptions.MessageException) as e: # if it's a webob redirect exception, pass it up the stack raise e diff --git a/lib/galaxy/tools/actions/__init__.py b/lib/galaxy/tools/actions/__init__.py index d37c518e294..0a928e2b4bf 100644 --- a/lib/galaxy/tools/actions/__init__.py +++ b/lib/galaxy/tools/actions/__init__.py @@ -505,11 +505,15 @@ class DefaultToolAction(object): element_kwds = dict(elements=collections_manager.ELEMENTS_UNINITIALIZED) else: element_kwds = dict(element_identifiers=element_identifiers) - output_collections.create_collection( + hdca = output_collections.create_collection( output=output, name=name, + set_hid=True if flush_job else False, + flush=flush_job, **element_kwds ) + if hdca: + datasets_to_persist.append(hdca) log.info("Handled collection output named %s for tool %s %s" % (name, tool.id, handle_output_timer)) else: handle_output(name, output) @@ -584,7 +588,7 @@ class DefaultToolAction(object): # Dispatch to a job handler. enqueue() is responsible for flushing the job app.job_manager.enqueue(job, tool=tool) trans.log_event("Added job to the job queue, id: %s" % str(job.id), tool_id=job.tool_id) - return job, out_data, datasets_to_persist + return job, out_data, datasets_to_persist, history def _remap_job_on_rerun(self, trans, galaxy_session, rerun_remap_job_id, current_job, out_data): """ @@ -807,7 +811,7 @@ class OutputCollections(object): self.out_collection_instances = {} self.tags = tags - def create_collection(self, output, name, collection_type=None, **element_kwds): + def create_collection(self, output, name, collection_type=None, set_hid=True, flush=True, **element_kwds): input_collections = self.input_collections collections_manager = self.trans.app.dataset_collections_service collection_type = collection_type or output.structure.collection_type @@ -886,11 +890,14 @@ class OutputCollections(object): collection_type=collection_type, trusted_identifiers=True, tags=self.tags, + set_hid=set_hid, + flush=flush, **element_kwds ) # name here is name of the output element - not name # of the hdca. self.out_collection_instances[name] = hdca + return hdca def on_text_for_names(input_names): diff --git a/lib/galaxy/tools/execute.py b/lib/galaxy/tools/execute.py index 9e409cbd11e..806fb003ad8 100644 --- a/lib/galaxy/tools/execute.py +++ b/lib/galaxy/tools/execute.py @@ -105,7 +105,7 @@ def execute(trans, tool, mapping_params, history, rerun_remap_job_id=None, colle datasets_to_persist.extend(execution_slice.datasets_to_persist) if datasets_to_persist: - history.add_datasets(trans.sa_session, datasets_to_persist, set_hid=True, quota=False, flush=False) + execution_slice.history.add_datasets(trans.sa_session, datasets_to_persist, set_hid=True, quota=False, flush=False) # a side effect of history.add_datasets is a commit within db_next_hid (even with flush=False). else: # Make sure collections, implicit jobs etc are flushed even if there are no precreated output datasets @@ -131,6 +131,7 @@ class ExecutionSlice(object): self.param_combination = param_combination self.dataset_collection_elements = dataset_collection_elements self.datasets_to_persist = None + self.history = None class ExecutionTracker(object): diff --git a/test/unit/tools/test_actions.py b/test/unit/tools/test_actions.py index db948079459..d8dfc9f83f7 100644 --- a/test/unit/tools/test_actions.py +++ b/test/unit/tools/test_actions.py @@ -134,12 +134,13 @@ class DefaultToolActionTestCase(unittest.TestCase, tools_support.UsesApp, tools_ if incoming is None: incoming = dict(param1="moo") self._init_tool(contents) - return self.action.execute( + job, out_data, _, _ = self.action.execute( tool=self.tool, trans=self.trans, history=self.history, incoming=incoming, ) + return job, out_data def test_determine_output_format(): From ed4951e689b0a53a2536226c62a3780572513fff Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Mon, 13 Jul 2020 11:12:19 +0200 Subject: [PATCH 08/12] Fix remap jobs This makes sure we only add hids / add items to history where needed. --- lib/galaxy/tools/actions/__init__.py | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/lib/galaxy/tools/actions/__init__.py b/lib/galaxy/tools/actions/__init__.py index 0a928e2b4bf..ff115337a9b 100644 --- a/lib/galaxy/tools/actions/__init__.py +++ b/lib/galaxy/tools/actions/__init__.py @@ -549,11 +549,16 @@ class DefaultToolAction(object): # Now that we have a job id, we can remap any outputs if this is a rerun and the user chose to continue dependent jobs # This functionality requires tracking jobs in the database. if app.config.track_jobs_in_database and rerun_remap_job_id is not None: + # We need a flush here in order to rewrite jobs parameter, + # but remapping jobs should only affect single jobs anyway, so this is not too costly. + trans.sa_session.flush() self._remap_job_on_rerun(trans=trans, galaxy_session=galaxy_session, rerun_remap_job_id=rerun_remap_job_id, current_job=job, out_data=out_data) + # remap_job_on_rerun may assign hids, only add datasets to history that don't have hid yet. + datasets_to_persist = [d for d in datasets_to_persist if d.hid is None] log.info("Setup for job %s complete, ready to be enqueued %s" % (job.log_str(), job_setup_timer)) # Some tools are not really executable, but jobs are still created for them ( for record keeping ). From 0f42daea4289848fbd2596b629e1b2a3188c163b Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Mon, 13 Jul 2020 12:45:23 +0200 Subject: [PATCH 09/12] Don't delete unpersisted tags --- lib/galaxy/model/tags.py | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/lib/galaxy/model/tags.py b/lib/galaxy/model/tags.py index f5792040433..38ebf312763 100644 --- a/lib/galaxy/model/tags.py +++ b/lib/galaxy/model/tags.py @@ -132,7 +132,9 @@ class TagHandler(object): """Delete tags from an item.""" # Delete item-tag associations. for tag in item.tags: - self.sa_session.delete(tag) + if tag.id: + # Only can and need to delete tag if tag is persisted + self.sa_session.delete(tag) # Delete tags from item. del item.tags[:] From 29b6cbb7e8b1098ddbb727b15a6bc0415fac6c56 Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Mon, 13 Jul 2020 18:51:49 +0200 Subject: [PATCH 10/12] Fix hid order again --- lib/galaxy/tools/actions/__init__.py | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/lib/galaxy/tools/actions/__init__.py b/lib/galaxy/tools/actions/__init__.py index ff115337a9b..d95ff67f38d 100644 --- a/lib/galaxy/tools/actions/__init__.py +++ b/lib/galaxy/tools/actions/__init__.py @@ -549,16 +549,16 @@ class DefaultToolAction(object): # Now that we have a job id, we can remap any outputs if this is a rerun and the user chose to continue dependent jobs # This functionality requires tracking jobs in the database. if app.config.track_jobs_in_database and rerun_remap_job_id is not None: - # We need a flush here in order to rewrite jobs parameter, + # We need a flush here and get hids in order to rewrite jobs parameter, # but remapping jobs should only affect single jobs anyway, so this is not too costly. trans.sa_session.flush() + history.add_datasets(trans.sa_session, datasets_to_persist, set_hid=set_output_hid, quota=False, flush=False) self._remap_job_on_rerun(trans=trans, galaxy_session=galaxy_session, rerun_remap_job_id=rerun_remap_job_id, current_job=job, out_data=out_data) - # remap_job_on_rerun may assign hids, only add datasets to history that don't have hid yet. - datasets_to_persist = [d for d in datasets_to_persist if d.hid is None] + datasets_to_persist = [] log.info("Setup for job %s complete, ready to be enqueued %s" % (job.log_str(), job_setup_timer)) # Some tools are not really executable, but jobs are still created for them ( for record keeping ). From 14f40c6fb14be462bf4002d1aff20befd8951453 Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Mon, 13 Jul 2020 17:12:28 +0200 Subject: [PATCH 11/12] Adjust DataManagerToolAction to new return values --- lib/galaxy/tools/actions/data_manager.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/lib/galaxy/tools/actions/data_manager.py b/lib/galaxy/tools/actions/data_manager.py index bd26cdf3f7b..6b0a761d212 100644 --- a/lib/galaxy/tools/actions/data_manager.py +++ b/lib/galaxy/tools/actions/data_manager.py @@ -10,7 +10,7 @@ class DataManagerToolAction(DefaultToolAction): def execute(self, tool, trans, **kwds): rval = super(DataManagerToolAction, self).execute(tool, trans, **kwds) - if isinstance(rval, tuple) and len(rval) == 2 and isinstance(rval[0], trans.app.model.Job): + if isinstance(rval, tuple) and len(rval) >= 2 and isinstance(rval[0], trans.app.model.Job): assoc = trans.app.model.DataManagerJobAssociation(job=rval[0], data_manager_id=tool.data_manager_id) trans.sa_session.add(assoc) trans.sa_session.flush() From ee971242cdabbc200d0955714697d6001e8c7a5b Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Tue, 14 Jul 2020 12:25:58 +0200 Subject: [PATCH 12/12] Cache collection datasets when matching element Iterating over all collection datasets to find a match gets more expensive the larger the collection is. This only affects mapped over jobs for tools that use format_source, structured_like or other tags that relate to their input in their output section. For 1000 paired collections with fastp this shaves off another 20 seconds: ``` *** PROFILER RESULTS *** execute (/Users/mvandenb/src/galaxy/lib/galaxy/tools/execute.py:33) function called 1 times 52219151 function calls (51664659 primitive calls) in 60.581 seconds Ordered by: cumulative time, internal time, call count List reduced from 1981 to 40 due to restriction <40> ncalls tottime percall cumtime percall filename:lineno(function) 1 0.014 0.014 60.774 60.774 execute.py:33(execute) 1000 0.014 0.000 27.455 0.027 execute.py:54(execute_single_job) 6 0.131 0.022 27.234 4.539 session.py:2489(flush) 5 0.050 0.010 27.023 5.405 session.py:2542(_flush) 1000 0.010 0.000 26.885 0.027 __init__.py:1587(handle_single_execution) 1000 0.016 0.000 26.874 0.027 __init__.py:1681(execute) 1000 0.114 0.000 26.859 0.027 __init__.py:280(execute) 54048 0.082 0.000 23.254 0.000 base.py:952(execute) 54048 0.064 0.000 23.158 0.000 elements.py:296(_execute_on_connection) 54048 0.307 0.000 23.094 0.000 base.py:1088(_execute_clauseelement) 3208071/3208065 1.757 0.000 22.337 0.000 attributes.py:279(__get__) 579205/554189 0.490 0.000 21.206 0.000 attributes.py:699(get) 54048 0.448 0.000 18.370 0.000 base.py:1195(_execute_context) 5 0.004 0.001 17.935 3.587 unitofwork.py:402(execute) 3 0.000 0.000 17.395 5.798 mapping.py:2827(db_next_hid) 1 0.000 0.000 17.336 17.336 __init__.py:1791(add_datasets) 8/6 0.000 0.000 17.299 2.883 session.py:899(begin) 8/6 0.000 0.000 17.299 2.883 session.py:221(__init__) 8/6 0.000 0.000 17.299 2.883 session.py:338(_take_snapshot) 1 0.005 0.005 17.175 17.175 __init__.py:1811(__add_datasets_optimized) 38 0.020 0.001 12.747 0.335 persistence.py:184(save_obj) 54042 0.037 0.000 12.186 0.000 default.py:589(do_execute) 54042 11.858 0.000 12.149 0.000 {method 'execute' of 'psycopg2.extensions.cursor' objects} 17031 0.067 0.000 11.594 0.001 query.py:3501(_execute_and_instances) 11012 0.024 0.000 11.237 0.001 scoping.py:162(do) 1000 0.014 0.000 11.109 0.011 __init__.py:247(_collect_inputs) 168 0.000 0.000 10.789 0.064 unitofwork.py:520(execute_aggregate) 6013 0.016 0.000 10.570 0.002 loading.py:190(load_on_ident) 6013 0.050 0.000 10.554 0.002 loading.py:211(load_on_pk_identity) 35 0.002 0.000 10.486 0.300 unitofwork.py:585(execute) 6014 0.007 0.000 10.480 0.002 query.py:3417(one) 6014 0.126 0.000 10.473 0.002 query.py:3381(one_or_none) 38 0.341 0.009 10.297 0.271 persistence.py:1039(_emit_insert_statements) 109036 0.185 0.000 9.492 0.000 strategies.py:665(_load_for_state) 2000 0.004 0.000 9.419 0.005 __init__.py:1477(visit_inputs) 32000/2000 0.201 0.000 9.415 0.005 __init__.py:22(visit_input_values) 1000 0.003 0.000 9.128 0.009 __init__.py:66(_collect_input_datasets) 11016 0.021 0.000 9.049 0.001 :1() 68000 0.154 0.000 9.032 0.000 __init__.py:118(callback_helper) 11016 0.120 0.000 9.028 0.001 strategies.py:772(_emit_lazyload) ``` --- lib/galaxy/tools/actions/__init__.py | 12 ++++++++++-- 1 file changed, 10 insertions(+), 2 deletions(-) diff --git a/lib/galaxy/tools/actions/__init__.py b/lib/galaxy/tools/actions/__init__.py index d95ff67f38d..27eca8426ec 100644 --- a/lib/galaxy/tools/actions/__init__.py +++ b/lib/galaxy/tools/actions/__init__.py @@ -33,6 +33,7 @@ class ToolExecutionCache(object): self.trans = trans self.current_user_roles = trans.get_current_user_roles() self.chrom_info = {} + self.cached_collection_elements = {} def get_chrom_info(self, tool_id, input_dbkey): genome_builds = self.trans.app.genome_builds @@ -381,6 +382,7 @@ class DefaultToolAction(object): inp_dataset_collections, input_ext, python_template_version=tool.python_template_version, + execution_cache=execution_cache, ) create_datasets = True dataset = None @@ -949,7 +951,7 @@ def get_ext_or_implicit_ext(hda): return hda.ext -def determine_output_format(output, parameter_context, input_datasets, input_dataset_collections, random_input_ext, python_template_version='3'): +def determine_output_format(output, parameter_context, input_datasets, input_dataset_collections, random_input_ext, python_template_version='3', execution_cache=None): """ Determines the output format for a dataset based on an abstract description of the output (galaxy.tool_util.parser.ToolOutput), the parameter wrappers, a map of the input datasets (name => HDA), and the last input @@ -990,7 +992,13 @@ def determine_output_format(output, parameter_context, input_datasets, input_dat try: input_element = input_collection_collection[element_index] except KeyError: - for element in input_collection_collection.dataset_elements: + if execution_cache: + dataset_elements = execution_cache.cached_collection_elements.get(input_collection_collection.id) + if dataset_elements is None: + dataset_elements = execution_cache.cached_collection_elements[input_collection_collection.id] = input_collection_collection.dataset_elements + else: + dataset_elements = input_collection_collection.dataset_elements + for element in dataset_elements: if element.element_identifier == element_index: input_element = element break