mirror of
https://github.com/galaxyproject/galaxy.git
synced 2026-09-24 16:30:27 +08:00
Optimize initial queueing of jobs.
- Don't refresh job when this is the only thread fetching it. - Avoid a bunch of unnecessary flushes, just flush once essentially during process. Brings this process from taking over 1 second per job on average on sqlite for cat1 on my laptop to around 400 ms on average.
This commit is contained in:
+20
-12
@@ -1029,26 +1029,33 @@ class JobWrapper( object ):
|
||||
self.sa_session.add( job )
|
||||
self.sa_session.flush()
|
||||
|
||||
def change_state( self, state, info=False ):
|
||||
job = self.get_job()
|
||||
self.sa_session.refresh( job )
|
||||
def change_state( self, state, info=False, flush=True, job=None ):
|
||||
job_supplied = job is not None
|
||||
if not job_supplied:
|
||||
job = self.get_job()
|
||||
self.sa_session.refresh( job )
|
||||
# Else:
|
||||
# If this is a new job (e.g. initially queued) - we are in the same
|
||||
# thread and no other threads are working on the job yet - so don't refresh.
|
||||
|
||||
if job.state in model.Job.terminal_states:
|
||||
log.warning( "(%s) Ignoring state change from '%s' to '%s' for job "
|
||||
"that is already terminal", job.id, job.state, state )
|
||||
return
|
||||
for dataset_assoc in job.output_datasets + job.output_library_datasets:
|
||||
dataset = dataset_assoc.dataset
|
||||
self.sa_session.refresh( dataset )
|
||||
dataset.state = state
|
||||
if not job_supplied:
|
||||
self.sa_session.refresh( dataset )
|
||||
dataset.raw_set_dataset_state( state )
|
||||
if info:
|
||||
dataset.info = info
|
||||
self.sa_session.add( dataset )
|
||||
self.sa_session.flush()
|
||||
if info:
|
||||
job.info = info
|
||||
job.set_state( state )
|
||||
self.sa_session.add( job )
|
||||
self.sa_session.flush()
|
||||
if flush:
|
||||
self.sa_session.flush()
|
||||
|
||||
def get_state( self ):
|
||||
job = self.get_job()
|
||||
@@ -1059,22 +1066,23 @@ class JobWrapper( object ):
|
||||
log.warning('set_runner() is deprecated, use set_job_destination()')
|
||||
self.set_job_destination(self.job_destination, external_id)
|
||||
|
||||
def set_job_destination( self, job_destination, external_id=None ):
|
||||
def set_job_destination( self, job_destination, external_id=None, flush=True, job=None ):
|
||||
"""
|
||||
Persist job destination params in the database for recovery.
|
||||
|
||||
self.job_destination is not used because a runner may choose to rewrite
|
||||
parts of the destination (e.g. the params).
|
||||
"""
|
||||
job = self.get_job()
|
||||
self.sa_session.refresh(job)
|
||||
if job is None:
|
||||
job = self.get_job()
|
||||
log.debug('(%s) Persisting job destination (destination id: %s)' % (job.id, job_destination.id))
|
||||
job.destination_id = job_destination.id
|
||||
job.destination_params = job_destination.params
|
||||
job.job_runner_name = job_destination.runner
|
||||
job.job_runner_external_id = external_id
|
||||
self.sa_session.add(job)
|
||||
self.sa_session.flush()
|
||||
if flush:
|
||||
self.sa_session.flush()
|
||||
|
||||
def finish( self, stdout, stderr, tool_exit_code=None, remote_working_directory=None ):
|
||||
"""
|
||||
@@ -1759,7 +1767,7 @@ class TaskWrapper(JobWrapper):
|
||||
self.status = 'error'
|
||||
# How do we want to handle task failure? Fail the job and let it clean up?
|
||||
|
||||
def change_state( self, state, info=False ):
|
||||
def change_state( self, state, info=False, flush=True, job=None ):
|
||||
task = self.get_task()
|
||||
self.sa_session.refresh( task )
|
||||
if info:
|
||||
|
||||
@@ -106,10 +106,12 @@ class BaseJobRunner( object ):
|
||||
"""Add a job to the queue (by job identifier), indicate that the job is ready to run.
|
||||
"""
|
||||
put_timer = ExecutionTimer()
|
||||
job = job_wrapper.get_job()
|
||||
# Change to queued state before handing to worker thread so the runner won't pick it up again
|
||||
job_wrapper.change_state( model.Job.states.QUEUED )
|
||||
job_wrapper.change_state( model.Job.states.QUEUED, flush=False, job=job )
|
||||
# Persist the destination so that the job will be included in counts if using concurrency limits
|
||||
job_wrapper.set_job_destination( job_wrapper.job_destination, None )
|
||||
job_wrapper.set_job_destination( job_wrapper.job_destination, None, flush=False, job=job )
|
||||
self.sa_session.flush()
|
||||
self.mark_as_queued(job_wrapper)
|
||||
log.debug("Job [%s] queued %s" % (job_wrapper.job_id, put_timer))
|
||||
|
||||
|
||||
@@ -1789,10 +1789,17 @@ class DatasetInstance( object ):
|
||||
return self._state
|
||||
return self.dataset.state
|
||||
|
||||
def raw_set_dataset_state( self, state ):
|
||||
if state != self.dataset.state:
|
||||
self.dataset.state = state
|
||||
return True
|
||||
else:
|
||||
return False
|
||||
|
||||
def set_dataset_state( self, state ):
|
||||
self.dataset.state = state
|
||||
object_session( self ).add( self.dataset )
|
||||
object_session( self ).flush() # flush here, because hda.flush() won't flush the Dataset object
|
||||
if self.raw_set_dataset_state( state ):
|
||||
object_session( self ).add( self.dataset )
|
||||
object_session( self ).flush() # flush here, because hda.flush() won't flush the Dataset object
|
||||
state = property( get_dataset_state, set_dataset_state )
|
||||
|
||||
def get_file_name( self ):
|
||||
|
||||
Reference in New Issue
Block a user