Merge pull request #10177 from mvdbeek/db_skip_locked

Add workflow invocation grabbing with db-skipped-lock
This commit is contained in:
Dannon
2020-09-21 21:39:29 -04:00
committed by GitHub
5 changed files with 204 additions and 81 deletions
+78 -59
View File
@@ -62,6 +62,71 @@ class JobHandler:
self.job_stop_queue.shutdown()
class ItemGrabber:
def __init__(self, app, grab_type='Job', handler_assignment_method=None, max_grab=None, self_handler_tags=None, handler_tags=None):
self.app = app
self.sa_session = app.model.context
self.grab_this = getattr(model, grab_type)
self.grab_type = grab_type
self._grab_conn_opts = {'autocommit': False}
subq = select([self.grab_this.id]) \
.where(and_(
self.grab_this.table.c.handler.in_(self_handler_tags),
self.grab_this.table.c.state == self.grab_this.states.NEW)) \
.order_by(self.grab_this.table.c.id)
if max_grab:
subq = subq.limit(max_grab)
if handler_assignment_method == HANDLER_ASSIGNMENT_METHODS.DB_SKIP_LOCKED:
subq = subq.with_for_update(skip_locked=True)
self._grab_query = self.grab_this.table.update() \
.returning(self.grab_this.table.c.id) \
.where(self.grab_this.table.c.id.in_(subq)) \
.values(handler=self.app.config.server_name)
if handler_assignment_method == HANDLER_ASSIGNMENT_METHODS.DB_TRANSACTION_ISOLATION:
self._grab_conn_opts['isolation_level'] = 'SERIALIZABLE'
log.info(
"Handler job grabber initialized with '%s' assignment method for handler '%s', tag(s): %s", handler_assignment_method,
self.app.config.server_name, ', '.join(str(x) for x in handler_tags)
)
@staticmethod
def get_grabbable_handler_assignment_method(handler_assignment_methods):
grabbable_methods = {
HANDLER_ASSIGNMENT_METHODS.DB_TRANSACTION_ISOLATION,
HANDLER_ASSIGNMENT_METHODS.DB_SKIP_LOCKED,
}
try:
return [m for m in handler_assignment_methods if m in grabbable_methods][0]
except IndexError:
return
def grab_unhandled_items(self):
"""
Attempts to assign unassigned jobs or invocaions to itself using DB serialization methods, if enabled. This
simply sets `Job.handler` or `WorkflowInvocation.handler` to the current server name, which causes the job to be picked up by
the appropriate handler.
"""
# an excellent discussion on PostgreSQL concurrency safety:
# https://blog.2ndquadrant.com/what-is-select-skip-locked-for-in-postgresql-9-5/
self.sa_session.expunge_all()
conn = self.sa_session.connection(execution_options=self._grab_conn_opts)
with conn.begin() as trans:
try:
rows = conn.execute(self._grab_query).fetchall()
if rows:
log.debug('Grabbed %s(s): %s', self.grab_type, ', '.join(str(row[0]) for row in rows))
trans.commit()
else:
trans.rollback()
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('Grabbing %s failed (serialization failures are ok): %s', self.grab_type, unicodify(e))
trans.rollback()
class JobHandlerQueue(Monitors):
"""
Job Handler's Internal Queue, this is what actually implements waiting for
@@ -91,38 +156,17 @@ class JobHandlerQueue(Monitors):
self.job_wrappers = {}
name = "JobHandlerQueue.monitor_thread"
self._init_monitor_thread(name, target=self.__monitor, config=app.config)
self.__grab_query = None
self.__grab_conn_opts = {'autocommit': False}
self.__initialize_job_grabbing()
def __initialize_job_grabbing(self):
grabbable_methods = {
HANDLER_ASSIGNMENT_METHODS.DB_TRANSACTION_ISOLATION,
HANDLER_ASSIGNMENT_METHODS.DB_SKIP_LOCKED,
}
try:
method = [m for m in self.app.job_config.handler_assignment_methods if m in grabbable_methods][0]
except IndexError:
return
subq = select([model.Job.id]) \
.where(and_(
model.Job.table.c.handler.in_(self.app.job_config.self_handler_tags),
model.Job.table.c.state == model.Job.states.NEW)) \
.order_by(model.Job.table.c.id)
if self.app.job_config.handler_max_grab:
subq = subq.limit(self.app.job_config.handler_max_grab)
if method == HANDLER_ASSIGNMENT_METHODS.DB_SKIP_LOCKED:
subq = subq.with_for_update(skip_locked=True)
self.__grab_query = model.Job.table.update() \
.returning(model.Job.table.c.id) \
.where(model.Job.table.c.id.in_(subq)) \
.values(handler=self.app.config.server_name)
if method == HANDLER_ASSIGNMENT_METHODS.DB_TRANSACTION_ISOLATION:
self.__grab_conn_opts['isolation_level'] = 'SERIALIZABLE'
log.info(
"Handler job grabber initialized with '%s' assignment method for handler '%s', tag(s): %s", method,
self.app.config.server_name, ', '.join(str(x) for x in self.app.job_config.handler_tags)
)
self.job_grabber = None
handler_assignment_method = ItemGrabber.get_grabbable_handler_assignment_method(self.app.job_config.handler_assignment_methods)
if handler_assignment_method:
self.job_grabber = ItemGrabber(
app=app,
grab_type='Job',
handler_assignment_method=handler_assignment_method,
max_grab=self.app.job_config.handler_max_grab,
self_handler_tags=self.app.job_config.self_handler_tags,
handler_tags=self.app.job_config.handler_tags,
)
def start(self):
"""
@@ -248,36 +292,11 @@ class JobHandlerQueue(Monitors):
'internal.galaxy.jobs.handlers.monitor_step',
'Job handler monitor step complete.'
)
if self.__grab_query is not None:
self.__grab_unhandled_jobs()
if self.job_grabber is not None:
self.job_grabber.grab_unhandled_items()
self.__handle_waiting_jobs()
log.trace(monitor_step_timer.to_str())
def __grab_unhandled_jobs(self):
"""
Attempts to "grab" jobs (assign unassigned jobs to itself) using DB serialization methods, if enabled. This
simply sets `Job.handler` to the current server name, which causes the job to be picked up by
`__handle_waiting_jobs()`.
"""
# an excellent discussion on PostgreSQL concurrency safety:
# https://blog.2ndquadrant.com/what-is-select-skip-locked-for-in-postgresql-9-5/
self.sa_session.expunge_all()
conn = self.sa_session.connection(execution_options=self.__grab_conn_opts)
with conn.begin() as trans:
try:
rows = conn.execute(self.__grab_query).fetchall()
if rows:
log.debug('Grabbed job(s): %s', ', '.join(str(row[0]) for row in rows))
trans.commit()
else:
trans.rollback()
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('Grabbing job failed (serialization failures are ok): %s', unicodify(e))
trans.rollback()
def __handle_waiting_jobs(self):
"""
Gets any new jobs (either from the database or from its own queue), then iterates over all new and waiting jobs
+20 -5
View File
@@ -4,6 +4,7 @@ from functools import partial
import galaxy.workflow.schedulers
from galaxy import model
from galaxy.exceptions import HandlerAssignmentError
from galaxy.jobs.handler import ItemGrabber
from galaxy.util import (
parse_xml,
plugin_config,
@@ -31,11 +32,9 @@ class WorkflowSchedulingManager(ConfiguresHandlers):
processes.
"""
DEFAULT_BASE_HANDLER_POOLS = ('workflow-schedulers', 'job-handlers')
UNSUPPORTED_HANDLER_ASSIGNMENT_METHODS = (
HANDLER_ASSIGNMENT_METHODS.DB_TRANSACTION_ISOLATION,
HANDLER_ASSIGNMENT_METHODS.DB_SKIP_LOCKED,
UNSUPPORTED_HANDLER_ASSIGNMENT_METHODS = {
HANDLER_ASSIGNMENT_METHODS.UWSGI_MULE_MESSAGE,
)
}
def __init__(self, app):
self.app = app
@@ -267,10 +266,27 @@ class WorkflowRequestMonitor(Monitors):
self.app = app
self.workflow_scheduling_manager = workflow_scheduling_manager
self._init_monitor_thread(name="WorkflowRequestMonitor.monitor_thread", target=self.__monitor, config=app.config)
self.invocation_grabber = None
if self.workflow_scheduling_manager.handler_assignment_methods_configured:
self_handler_tags = set(self.app.job_config.self_handler_tags)
self_handler_tags.add(self.workflow_scheduling_manager.default_handler_id)
handler_assignment_method = ItemGrabber.get_grabbable_handler_assignment_method(self.workflow_scheduling_manager.handler_assignment_methods)
if handler_assignment_method:
self.invocation_grabber = ItemGrabber(
app=app,
grab_type='WorkflowInvocation',
handler_assignment_method=handler_assignment_method,
max_grab=self.workflow_scheduling_manager.handler_max_grab,
self_handler_tags=self_handler_tags,
handler_tags=self_handler_tags,
)
def __monitor(self):
to_monitor = self.workflow_scheduling_manager.active_workflow_schedulers
while self.monitor_running:
if self.invocation_grabber:
self.invocation_grabber.grab_unhandled_items()
monitor_step_timer = self.app.execution_timer_factory.get_timer(
'internal.galaxy.workflows.scheduling_manager.monitor_step',
'Workflow scheduling manager monitor step complete.'
@@ -280,7 +296,6 @@ class WorkflowRequestMonitor(Monitors):
return
self.__schedule(workflow_scheduler_id, workflow_scheduler)
log.trace(monitor_step_timer.to_str())
self._monitor_sleep(1)
@@ -17,6 +17,7 @@ from .driver_util import GalaxyTestDriver
NO_APP_MESSAGE = "test_case._app called though no Galaxy has been configured."
# Following should be for Homebrew Rabbitmq and Docker on Mac "amqp://guest:guest@localhost:5672//"
AMQP_URL = os.environ.get("GALAXY_TEST_AMQP_URL", None)
POSTGRES_CONFIGURED = 'postgres' in os.environ.get("GALAXY_TEST_DBURI", '')
def _identity(func):
@@ -36,6 +37,12 @@ def skip_unless_amqp():
return pytest.mark.skip("AMQP_URL is not set, required for this test.")
def skip_unless_postgres():
if POSTGRES_CONFIGURED:
return _identity
return pytest.mark.skip("GALAXY_TEST_DBURI does not point to postgres database, required for this test.")
def skip_unless_executable(executable):
if which(executable):
return _identity
@@ -2,6 +2,9 @@
import os
import re
import string
import tempfile
import time
from json import dumps
from galaxy_test.base.populators import (
@@ -12,8 +15,42 @@ from galaxy_test.driver import integration_util
SCRIPT_DIRECTORY = os.path.abspath(os.path.dirname(__file__))
WORKFLOW_HANDLER_CONFIGURATION_JOB_CONF = os.path.join(SCRIPT_DIRECTORY, "workflow_handler_configuration_job_conf.xml")
WORKFLOW_HANDLER_CONFIGURATION_WORKFLOW_SCHEDULER_CONF = os.path.join(SCRIPT_DIRECTORY, "workflow_handler_configuration_workflow_scheduler_conf.xml")
WORKFLOW_HANDLER_JOB_CONFIG_TEMPLATE = string.Template("""
<job_conf>
<plugins>
<plugin id="local" type="runner" load="galaxy.jobs.runners.local:LocalJobRunner" workers="2"/>
</plugins>
<handlers default="handlers" $assign_with>
<handler id="handler0" tags="handlers"/>
<handler id="handler1" tags="handlers"/>
<handler id="handler2" tags="handlers" />
<handler id="handler3" tags="handlers" />
<handler id="handler4" tags="handlers" />
<handler id="handler5" tags="handlers" />
<handler id="handler6" tags="handlers" />
<handler id="handler7" tags="handlers" />
<handler id="handler8" tags="handlers" />
<handler id="handler9" tags="handlers" />
</handlers>
<destinations default="local">
<destination id="local" runner="local">
</destination>
</destinations>
</job_conf>
""")
WORKFLOW_SCHEDULERS_CONFIG_TEMPLATE = string.Template("""
<workflow_schedulers default="core">
<core id="core" />
<handlers default="workflow_handlers" $assign_with>
<handler id="work1" tags="workflow_handlers" />
<handler id="work2" tags="workflow_handlers" />
</handlers>
</workflow_schedulers>
""")
JOB_HANDLER_PATTERN = re.compile(r"handler\d")
WORKFLOW_SCHEDULER_HANDLER_PATTERN = re.compile(r"work\d")
@@ -30,9 +67,20 @@ steps:
"""
def config_file(template, assign_with=''):
fd, path = tempfile.mkstemp(suffix=".xml", prefix="workflow_handler_config_")
os.close(fd)
with open(path, 'w') as config:
if assign_with:
assign_with = 'assign_with="{}"'.format(assign_with)
config.write(template.substitute(assign_with=assign_with))
return path
class BaseWorkflowHandlerConfigurationTestCase(integration_util.IntegrationTestCase):
framework_tool_and_types = True
assign_with = ""
def setUp(self):
super(BaseWorkflowHandlerConfigurationTestCase, self).setUp()
@@ -42,7 +90,7 @@ class BaseWorkflowHandlerConfigurationTestCase(integration_util.IntegrationTestC
@classmethod
def handle_galaxy_config_kwds(cls, config):
config["job_config_file"] = WORKFLOW_HANDLER_CONFIGURATION_JOB_CONF
config["job_config_file"] = config_file(WORKFLOW_HANDLER_JOB_CONFIG_TEMPLATE, assign_with=cls.assign_with)
def _invoke_n_workflows(self, n):
workflow_id = self.workflow_populator.upload_yaml_workflow(PAUSE_WORKFLOW)
@@ -114,7 +162,7 @@ class WorkflowSchedulerHandlerAssignment(BaseWorkflowHandlerConfigurationTestCas
@classmethod
def handle_galaxy_config_kwds(cls, config):
BaseWorkflowHandlerConfigurationTestCase.handle_galaxy_config_kwds(config)
config["workflow_schedulers_config_file"] = WORKFLOW_HANDLER_CONFIGURATION_WORKFLOW_SCHEDULER_CONF
config["workflow_schedulers_config_file"] = config_file(WORKFLOW_SCHEDULERS_CONFIG_TEMPLATE, assign_with=cls.assign_with)
def test_handler_assignment(self):
self._invoke_n_workflows(1)
@@ -122,11 +170,14 @@ class WorkflowSchedulerHandlerAssignment(BaseWorkflowHandlerConfigurationTestCas
assert WORKFLOW_SCHEDULER_HANDLER_PATTERN.match(workflow_invocations[0].handler)
# Follow five classes test 5 different ways Galaxy processes can be workflow schedulers or not.
# Following 8 classes test 8 different ways Galaxy processes can be workflow schedulers or not.
# - In single process mode, the process is a workflow scheduler.
# - If no workflow schedulers conf is configured and the process is a job handler, it is a workflow scheduler as well.
# - If no workflow schedulers conf is configured and the process is a job handler, it is a workflow scheduler as well (with db-dkip-locked).
# - If no workflow schedulers conf is configured and the process is not a job handler, it is not a workflow scheduler as well.
# - If a workflow scheduler conf is defined and the process is listed as a handler, it is workflow scheduler.
# - If a workflow scheduler conf is defined and assign_with is set to db-skip-locked, invocation handler is correctly set
# - If a workflow scheduler conf is defined and assign_with is set to db-transaction-isolation, invocation handler is correctly set
# - If a workflow scheduler conf is defined and the process is not listed as a handler, it is not workflow scheduler.
class DefaultWorkflowHandlerOnTestCase(BaseWorkflowHandlerConfigurationTestCase):
@@ -150,6 +201,25 @@ class DefaultWorkflowHandlerIfJobHandlerOnTestCase(BaseWorkflowHandlerConfigurat
assert self.is_app_workflow_scheduler
class JobHandlerAsWorkflowHandlerWithDbSkipLocked(BaseWorkflowHandlerConfigurationTestCase):
assign_with = 'db-skip-locked'
@classmethod
def handle_galaxy_config_kwds(cls, config):
BaseWorkflowHandlerConfigurationTestCase.handle_galaxy_config_kwds(config)
config["server_name"] = "handler0"
def test_handler_assignment(self):
self._invoke_n_workflows(1)
time.sleep(2)
workflow_invocations = self._get_workflow_invocations()
assert JOB_HANDLER_PATTERN.match(workflow_invocations[0].handler)
def test_default_job_handler_is_workflow_handler(self):
assert self.is_app_workflow_scheduler
class DefaultWorkflowHandlerIfJobHandlerOffTestCase(BaseWorkflowHandlerConfigurationTestCase):
@classmethod
@@ -157,29 +227,49 @@ class DefaultWorkflowHandlerIfJobHandlerOffTestCase(BaseWorkflowHandlerConfigura
BaseWorkflowHandlerConfigurationTestCase.handle_galaxy_config_kwds(config)
config["server_name"] = "web0"
def test_default_job_handler_is_workflow_handler(self):
def test_default_job_handler_is_not_workflow_handler(self):
assert not self.is_app_workflow_scheduler
class ExplicitWorkflowHandlersOnTestCase(BaseWorkflowHandlerConfigurationTestCase):
assign_with = ""
@classmethod
def handle_galaxy_config_kwds(cls, config):
BaseWorkflowHandlerConfigurationTestCase.handle_galaxy_config_kwds(config)
config["workflow_schedulers_config_file"] = WORKFLOW_HANDLER_CONFIGURATION_WORKFLOW_SCHEDULER_CONF
config["workflow_schedulers_config_file"] = config_file(WORKFLOW_SCHEDULERS_CONFIG_TEMPLATE, assign_with=cls.assign_with)
config["server_name"] = "work1"
def test_workflows_spread_across_multiple_handlers(self):
def test_app_is_workflow_scheduler(self):
assert self.is_app_workflow_scheduler
@integration_util.skip_unless_postgres()
class WorkflowSchedulerHandlerAssignmentDbSkipLocked(ExplicitWorkflowHandlersOnTestCase):
assign_with = 'db-skip-locked'
def test_handler_assignment(self):
self._invoke_n_workflows(1)
time.sleep(2)
workflow_invocations = self._get_workflow_invocations()
assert WORKFLOW_SCHEDULER_HANDLER_PATTERN.match(workflow_invocations[0].handler)
@integration_util.skip_unless_postgres()
class WorkflowSchedulerHandlerAssignmentDbTransactionIsolation(WorkflowSchedulerHandlerAssignmentDbSkipLocked):
assign_with = 'db-transaction-isolation'
class ExplicitWorkflowHandlersOffTestCase(BaseWorkflowHandlerConfigurationTestCase):
@classmethod
def handle_galaxy_config_kwds(cls, config):
BaseWorkflowHandlerConfigurationTestCase.handle_galaxy_config_kwds(config)
config["workflow_schedulers_config_file"] = WORKFLOW_HANDLER_CONFIGURATION_WORKFLOW_SCHEDULER_CONF
config["workflow_schedulers_config_file"] = config_file(WORKFLOW_SCHEDULERS_CONFIG_TEMPLATE, assign_with=cls.assign_with)
config["server_name"] = "handler0" # Configured as a job handler but not a workflow handler.
def test_workflows_spread_across_multiple_handlers(self):
def test_app_is_not_workflow_scheduler(self):
assert not self.is_app_workflow_scheduler
@@ -1,8 +0,0 @@
<?xml version="1.0"?>
<workflow_schedulers default="core">
<core id="core" />
<handlers default="workflow_handlers">
<handler id="work1" tags="workflow_handlers" />
<handler id="work2" tags="workflow_handlers" />
</handlers>
</workflow_schedulers>