Use SKIP LOCKED only if database supports it

This commit is contained in:
mvdbeek
2021-04-16 16:27:39 +02:00
parent 44b988ca2f
commit a3f59a25ef
4 changed files with 31 additions and 17 deletions
+3 -3
View File
@@ -1321,7 +1321,7 @@ class JobWrapper(HasResourceParameters):
# Pause any dependent jobs (and those jobs' outputs)
for dep_job_assoc in dataset.dependent_jobs:
self.pause(dep_job_assoc.job, "Execution of this dataset's job is paused because its input datasets are in an error state.")
job.set_final_state(job.states.ERROR)
job.set_final_state(job.states.ERROR, supports_skip_locked=self.app.application_stack.supports_skip_locked())
job.command_line = unicodify(self.command_line)
job.info = message
# TODO: Put setting the stdout, stderr, and exit code in one place
@@ -1408,7 +1408,7 @@ class JobWrapper(HasResourceParameters):
job.info = info
job.set_state(state)
self.sa_session.add(job)
job.update_output_states()
job.update_output_states(self.app.application_stack.supports_skip_locked())
if flush:
self.sa_session.flush()
@@ -1776,7 +1776,7 @@ class JobWrapper(HasResourceParameters):
# Finally set the job state. This should only happen *after* all
# dataset creation, and will allow us to eliminate force_history_refresh.
job.set_final_state(final_job_state)
job.set_final_state(final_job_state, supports_skip_locked=self.app.application_stack.supports_skip_locked())
if not job.tasks:
# If job was composed of tasks, don't attempt to recollect statistics
self._collect_metrics(job, job_metrics_directory)
+1 -1
View File
@@ -924,7 +924,7 @@ class JobHandlerStopQueue(Monitors):
if error_msg is not None:
final_state = job.states.ERROR
job.info = error_msg
job.set_final_state(final_state)
job.set_final_state(final_state, supports_skip_locked=self.app.application_stack.supports_skip_locked())
self.sa_session.add(job)
self.sa_session.flush()
+23 -9
View File
@@ -37,6 +37,7 @@ from sqlalchemy import (
true,
type_coerce,
types)
from sqlalchemy.exc import OperationalError
from sqlalchemy.ext import hybrid
from sqlalchemy.orm import (
aliased,
@@ -1177,22 +1178,36 @@ class Job(JobLike, UsesCreateAndUpdateTime, Dictifiable, RepresentById):
return rval
def update_hdca_update_time_for_job(self, update_time, sa_session):
def update_hdca_update_time_for_job(self, update_time, sa_session, supports_skip_locked):
subq = sa_session.query(HistoryDatasetCollectionAssociation.id) \
.join(ImplicitCollectionJobs) \
.join(ImplicitCollectionJobsJobAssociation) \
.filter(ImplicitCollectionJobsJobAssociation.job_id == self.id) \
.with_for_update(skip_locked=True).subquery()
.filter(ImplicitCollectionJobsJobAssociation.job_id == self.id)
if supports_skip_locked:
subq = subq.with_for_update(skip_locked=True).subquery()
implicit_statement = HistoryDatasetCollectionAssociation.table.update() \
.where(HistoryDatasetCollectionAssociation.table.c.id.in_(subq)) \
.values(update_time=update_time)
explicit_statement = HistoryDatasetCollectionAssociation.table.update() \
.where(HistoryDatasetCollectionAssociation.table.c.job_id == self.id) \
.values(update_time=update_time)
sa_session.execute(implicit_statement)
sa_session.execute(explicit_statement)
if supports_skip_locked:
sa_session.execute(implicit_statement)
else:
conn = sa_session.connection(execution_options={'isolation_level': 'SERIALIZABLE'})
with conn.begin() as trans:
try:
conn.execute(implicit_statement)
trans.commit()
except OperationalError as e:
# If this is a serialization failure on PostgreSQL, then e.orig is a psycopg2 TransactionRollbackError
# and should have attribute `code`. Other engines should just report the message and move on.
if int(getattr(e.orig, 'pgcode', -1)) != 40001:
log.debug(f"Updating implicit collection uptime_time for job {self.id} failed (this is expected for large collections and not a problem): {unicodify(e)}")
trans.rollback()
def set_final_state(self, final_state):
def set_final_state(self, final_state, supports_skip_locked):
self.set_state(final_state)
# TODO: migrate to where-in subqueries?
statement = '''
@@ -1202,8 +1217,7 @@ class Job(JobLike, UsesCreateAndUpdateTime, Dictifiable, RepresentById):
'''
sa_session = object_session(self)
update_time = galaxy.model.orm.now.now()
log.debug(f'Final state update time: {update_time}')
self.update_hdca_update_time_for_job(update_time=update_time, sa_session=sa_session)
self.update_hdca_update_time_for_job(update_time=update_time, sa_session=sa_session, supports_skip_locked=supports_skip_locked)
params = {
'job_id': self.id,
'update_time': update_time
@@ -1230,7 +1244,7 @@ class Job(JobLike, UsesCreateAndUpdateTime, Dictifiable, RepresentById):
for dataset_assoc in self.output_datasets:
return dataset_assoc.dataset.tool_version
def update_output_states(self):
def update_output_states(self, supports_skip_locked):
# TODO: migrate to where-in subqueries?
statements = ['''
UPDATE dataset
@@ -1275,7 +1289,7 @@ class Job(JobLike, UsesCreateAndUpdateTime, Dictifiable, RepresentById):
''']
sa_session = object_session(self)
update_time = galaxy.model.orm.now.now()
self.update_hdca_update_time_for_job(update_time=update_time, sa_session=sa_session)
self.update_hdca_update_time_for_job(update_time=update_time, sa_session=sa_session, supports_skip_locked=supports_skip_locked)
params = {
'job_id': self.id,
'state': self.state,
@@ -241,15 +241,15 @@ class HistoryAsContainerTestCase(HistoryAsContainerBaseTestCase):
contents.append(self.add_list_collection_to_history(history, contents[4:6]))
self.log("should allow filtering by update_time")
# in the case of collections we have to change the collection.collection (ugh) to change the update_time
contents[3].collection.populated_state = 'big ball of mud'
# change the update_time by updating the name
contents[3].name = 'big ball of mud'
self.app.model.context.flush()
update_time = contents[3].collection.update_time
update_time = contents[3].update_time
def get_update_time(item):
update_time = getattr(item, 'update_time', None)
if not update_time:
update_time = item.collection.update_time
update_time = item.update_time
return update_time
results = self.contents_manager.contents(history, filters=[parsed_filter("orm", column('update_time') >= update_time)])