diff --git a/config/job_conf.xml.sample_advanced b/config/job_conf.xml.sample_advanced index c0d702a0189..c7fbea173e7 100644 --- a/config/job_conf.xml.sample_advanced +++ b/config/job_conf.xml.sample_advanced @@ -332,9 +332,9 @@ For documentation on handler assignment methods, see the documentation under: https://docs.galaxyproject.org/en/latest/admin/scaling.html#job-handler-assignment-methods - The container tag takes two optional attributes: + The container tag takes three optional attributes: - + - `assign_with` - How jobs should be assigned to handlers. The value can be a single method or a comma-separated list that will be tried in order. The default depends on whether any handlers and a job @@ -348,6 +348,11 @@ - `db-transaction-isolation` - Database Transaction Isolation - `db-skip-locked` - Database SKIP LOCKED + - `max_grab` - Limit the number of jobs that a handler will self-assign per jobs-ready-to-run check loop + iteration. This only applies for methods that cause handlers to self-assign multiple jobs at once + (db-skip-locked, db-transaction-isolation) and the value is an integer > 0. Default is to grab as many + jobs ready to run as possible. + - `default` - An ID or tag of the handler(s) that should handle any jobs not assigned to a specific handler (which is probably most of them). If unset, the default is any untagged handlers plus any handlers in the `job-handlers` (no tag) pool. diff --git a/lib/galaxy/queue_worker.py b/lib/galaxy/queue_worker.py index 7d5b3302858..23ac03f7597 100644 --- a/lib/galaxy/queue_worker.py +++ b/lib/galaxy/queue_worker.py @@ -16,7 +16,10 @@ from kombu import ( uuid, ) from kombu.mixins import ConsumerProducerMixin -from kombu.pools import producers +from kombu.pools import ( + connections, + producers, +) from six.moves import reload_module import galaxy.queues @@ -308,6 +311,7 @@ class GalaxyQueueWorker(ConsumerProducerMixin, threading.Thread): # Force connection instead of lazy-connecting the first time it is required. # Fixes `'kombu.transport.sqlalchemy.Message' is not mapped` error. self.connection.connect() + self.connection.release() self.app = app self.task_mapping = task_mapping self.exchange_queue = None @@ -325,8 +329,9 @@ class GalaxyQueueWorker(ConsumerProducerMixin, threading.Thread): self.exchange_queue, self.direct_queue = galaxy.queues.control_queues_from_config(self.app.config) self.control_queues = [self.exchange_queue, self.direct_queue] # Delete messages for the current workers' control queues on startup - for q in self.control_queues: - q(self.connection).delete() + with connections[self.connection].acquire(block=True) as conn: + for q in self.control_queues: + q(conn).delete() self.start() def get_consumers(self, Consumer, channel):