Merge pull request #2057 from jmchilton/pulsar_overhaul

Implement an Embedded Pulsar Job Runner
This commit is contained in:
Nate Coraor
2016-04-01 17:19:40 -04:00
2 changed files with 108 additions and 24 deletions
+20
View File
@@ -98,6 +98,26 @@
deprecated and will disappear with a future release of Galaxy.
-->
</plugin>
<plugin id="pulsar_embedded" type="runner" load="galaxy.jobs.runners.pulsar:PulsarEmbeddedJobRunner">
<!-- The embedded Pulsar runner starts a Pulsar app
internal to Galaxy and communicates it directly.
This maybe be useful for instance when Pulsar
staging is important but a Pulsar server is
unneeded (most obviously for instance if compute
servers cannot mount Galaxy's files but Galaxy
can mount a scratch directory available on
compute). -->
<!-- Specify a complete description of the Pulsar app
to create. Currently this configuration (if set)
must create exactly on job manager. For more
information on configuring a Pulsar app see:
https://github.com/galaxyproject/pulsar/blob/master/app.yml.sample
http://pulsar.readthedocs.org/en/latest/configure.html
-->
<!-- <param id="pulsar_conf">path/to/pulsar/app.yml</param> -->
</plugin>
</plugins>
<handlers default="handlers">
<!-- Additional job handlers - the id should match the name of a
+88 -24
View File
@@ -1,6 +1,26 @@
"""Job runner used to execute Galaxy jobs through Pulsar.
More infromation on Pulsar can be found at http://pulsar.readthedocs.org/.
"""
from __future__ import absolute_import # Need to import pulsar_client absolutely.
import errno
import logging
import os
from time import sleep
from pulsar.client import build_client_manager
from pulsar.client import url_to_destination_params
from pulsar.client import finish_job as pulsar_finish_job
from pulsar.client import submit_job as pulsar_submit_job
from pulsar.client import ClientJobDescription
from pulsar.client import PulsarOutputs
from pulsar.client import ClientOutputs
from pulsar.client import PathMapper
import pulsar.core
import yaml
from galaxy import model
from galaxy.jobs.runners import AsynchronousJobState, AsynchronousJobRunner
@@ -12,22 +32,15 @@ from galaxy.util import string_as_bool_or_none
from galaxy.util.bunch import Bunch
from galaxy.util import specs
import errno
from time import sleep
import os
from pulsar.client import build_client_manager
from pulsar.client import url_to_destination_params
from pulsar.client import finish_job as pulsar_finish_job
from pulsar.client import submit_job as pulsar_submit_job
from pulsar.client import ClientJobDescription
from pulsar.client import PulsarOutputs
from pulsar.client import ClientOutputs
from pulsar.client import PathMapper
log = logging.getLogger( __name__ )
__all__ = [ 'PulsarLegacyJobRunner', 'PulsarRESTJobRunner', 'PulsarMQJobRunner' ]
__all__ = [
'PulsarLegacyJobRunner',
'PulsarRESTJobRunner',
'PulsarMQJobRunner',
'PulsarEmbeddedJobRunner',
]
NO_REMOTE_GALAXY_FOR_METADATA_MESSAGE = "Pulsar misconfiguration - Pulsar client configured to set metadata remotely, but remote Pulsar isn't properly configured with a galaxy_home directory."
NO_REMOTE_DATATYPES_CONFIG = "Pulsar client is configured to use remote datatypes configuration when setting metadata externally, but Pulsar is not configured with this information. Defaulting to datatypes_conf.xml."
@@ -57,6 +70,10 @@ PULSAR_PARAM_SPECS = dict(
map=specs.to_str_or_none,
default=None,
),
pulsar_config=dict(
map=specs.to_str_or_none,
default=None,
),
manager=dict(
map=specs.to_str_or_none,
default=None,
@@ -133,13 +150,13 @@ PARAMETER_SPECIFICATION_IGNORED = object()
class PulsarJobRunner( AsynchronousJobRunner ):
"""
Pulsar Job Runner
"""
"""Base class for pulsar job runners."""
runner_name = "PulsarJobRunner"
default_build_pulsar_app = False
def __init__( self, app, nworkers, **kwds ):
"""Start the job runner """
"""Start the job runner."""
super( PulsarJobRunner, self ).__init__( app, nworkers, runner_param_specs=PULSAR_PARAM_SPECS, **kwds )
self._init_worker_threads()
galaxy_url = self.runner_params.galaxy_url
@@ -156,16 +173,42 @@ class PulsarJobRunner( AsynchronousJobRunner ):
self._init_monitor_thread()
def __init_client_manager( self ):
pulsar_conf = self.runner_params.get('pulsar_conf', None)
self.__init_pulsar_app(pulsar_conf)
client_manager_kwargs = {}
for kwd in 'manager', 'cache', 'transport', 'persistence_directory':
client_manager_kwargs[ kwd ] = self.runner_params[ kwd ]
if self.pulsar_app is not None:
# TODO: Make this more generic and configurable - client_manager
# should define an app and client (destination) should reference
# a job manager.
job_manager = self.pulsar_app.only_manager
client_manager_kwargs[ "job_manager" ] = job_manager
# TODO: Hack remove this following line pulsar lib update
# that includes https://github.com/galaxyproject/pulsar/commit/ce0636a5b64fae52d165bcad77b2caa3f0e9c232
client_manager_kwargs[ "file_cache" ] = None
for kwd in self.runner_params.keys():
if kwd.startswith( 'amqp_' ):
client_manager_kwargs[ kwd ] = self.runner_params[ kwd ]
self.client_manager = build_client_manager(**client_manager_kwargs)
def __init_pulsar_app( self, pulsar_conf_path ):
if pulsar_conf_path is None and not self.default_build_pulsar_app:
self.pulsar_app = None
return
conf = {}
if pulsar_conf_path is None:
log.info("Creating a Pulsar app with default configuration (no pulsar_conf specified).")
else:
log.info("Loading Pulsar app configuration from %s" % pulsar_conf_path)
with open(pulsar_conf_path, "r") as f:
conf.update(yaml.load(f) or {})
self.pulsar_app = pulsar.core.PulsarApp(**conf)
def url_to_destination( self, url ):
"""Convert a legacy URL to a job destination"""
"""Convert a legacy URL to a job destination."""
return JobDestination( runner="pulsar", params=url_to_destination_params( url ) )
def check_watched_item(self, job_state):
@@ -243,7 +286,7 @@ class PulsarJobRunner( AsynchronousJobRunner ):
self.monitor_job(pulsar_job_state)
def __prepare_job(self, job_wrapper, job_destination):
""" Build command-line and Pulsar client for this job. """
"""Build command-line and Pulsar client for this job."""
command_line = None
client = None
remote_job_config = None
@@ -424,9 +467,7 @@ class PulsarJobRunner( AsynchronousJobRunner ):
job_wrapper.fail("Unable to finish job", exception=True)
def fail_job( self, job_state, message=GENERIC_REMOTE_ERROR ):
"""
Seperated out so we can use the worker threads for it.
"""
"""Seperated out so we can use the worker threads for it."""
self.stop_job( self.sa_session.query( self.app.model.Job ).get( job_state.job_wrapper.job_id ) )
job_state.job_wrapper.fail( getattr( job_state, "fail_message", message ) )
@@ -475,7 +516,7 @@ class PulsarJobRunner( AsynchronousJobRunner ):
client.kill()
def recover( self, job, job_wrapper ):
"""Recovers jobs stuck in the queued/running state when Galaxy started"""
"""Recover jobs stuck in the queued/running state when Galaxy started."""
job_state = self._job_state( job, job_wrapper )
job_wrapper.command_line = job.get_command_line()
state = job.get_state()
@@ -538,7 +579,9 @@ class PulsarJobRunner( AsynchronousJobRunner ):
@staticmethod
def __use_remote_datatypes_conf( pulsar_client ):
""" When setting remote metadata, use integrated datatypes from this
"""Use remote metadata datatypes instead of Galaxy's.
When setting remote metadata, use integrated datatypes from this
Galaxy instance or use the datatypes config configured via the remote
Pulsar.
@@ -604,6 +647,8 @@ class PulsarJobRunner( AsynchronousJobRunner ):
class PulsarLegacyJobRunner( PulsarJobRunner ):
"""Flavor of Pulsar job runner mimicking behavior of old LWR runner."""
destination_defaults = dict(
rewrite_parameters="false",
dependency_resolution="local",
@@ -611,6 +656,8 @@ class PulsarLegacyJobRunner( PulsarJobRunner ):
class PulsarMQJobRunner( PulsarJobRunner ):
"""Flavor of Pulsar job runner with sensible defaults for message queue communication."""
destination_defaults = dict(
default_file_action="remote_transfer",
rewrite_parameters="true",
@@ -640,6 +687,8 @@ class PulsarMQJobRunner( PulsarJobRunner ):
class PulsarRESTJobRunner( PulsarJobRunner ):
"""Flavor of Pulsar job runner with sensible defaults for RESTful usage."""
destination_defaults = dict(
default_file_action="transfer",
rewrite_parameters="true",
@@ -648,6 +697,21 @@ class PulsarRESTJobRunner( PulsarJobRunner ):
)
class PulsarEmbeddedJobRunner(PulsarJobRunner):
"""Flavor of Puslar job runnner that runs Pulsar's server code directly within Galaxy.
This is an appropriate job runner for when the desire is to use Pulsar staging
but their is not need to run a remote service.
"""
destination_defaults = dict(
default_file_action="copy",
rewrite_parameters="true",
dependency_resolution="remote",
)
default_build_pulsar_app = True
class PulsarComputeEnvironment( ComputeEnvironment ):
def __init__( self, pulsar_client, job_wrapper, remote_job_config ):