Merge pull request #2314 from phnmnl/feature/k8s_job_runner

Kubernetes job runner
This commit is contained in:
John Chilton
2016-06-05 09:25:30 -04:00
6 changed files with 486 additions and 7 deletions
+137
View File
@@ -117,6 +117,83 @@
-->
<!-- <param id="pulsar_config">path/to/pulsar/app.yml</param> -->
</plugin>
<plugin id="k8s" type="runner" load="galaxy.jobs.runners.kubernetes:kubernetes">
<!-- The Kubernetes (k8s) plugin allows to send jobs to a k8s cluster which shares filesystem with Galaxy.
The shared file system needs to be exposed to k8s through a Persistent Volume (rw) and a Persistent
Volume Claim. An example of a Persistent Volume could be, in yaml (access modes, reclaim policy and
path are relevant) (persistent_volume.yaml):
kind: PersistentVolume
apiVersion: v1
metadata:
name: pv-galaxy-nfs
labels:
type: nfs
spec:
capacity:
storage: 10Gi
accessModes:
- ReadWriteMany
persistentVolumeReclaimPolicy: Retained
nfs:
path: /scratch1/galaxy_data
server: 192.168.64.1
The path (nfs:path: in the example) set needs to be a parent directory of the directories used for
variables “file_path” and “new_file_path” on the galaxy.ini files. Clearly, for this particular example
to work, there needs to be a NFS server serving that directory on that ip. Please make sure that you
use reasonable storage size for your set up (possibly larger that the 10Gi written).
An example of the volume claim should be (this needs to be followed more closely) (pv_claim.yaml):
kind: PersistentVolumeClaim
apiVersion: v1
metadata:
name: galaxy-pvc
spec:
accessModes:
- ReadWriteMany
volumeName: pv-galaxy-nfs
resources:
requests:
storage: 2Gi
The volume claim needs to reference the name of the volume in spec:volumeName. The name of the claim
(metadat:name) is referenced in the plugin definition (see below), through param
"k8s_persistent_volume_claim_name". These two k8s object need to be created before galaxy can use them:
kubectl create -f <path/to/persistent_volume.yaml>
kubectl create -f <path/to/pv_claim.yaml>
pointing of course to the same Kubernetes cluster that you intend to use.
-->
<param id="k8s_config_path">/path/to/kubeconfig</param>
<!--- This is the path to the kube config file, which is normally on ~/.kube/config, but that will depend on
your installation. This is the file that tells the plugin where the k8s cluster is, access credentials,
etc. -->
<param id="k8s_persistent_volume_claim_name">galaxy_pvc</param>
<!-- The name of the Persisten Volume Claim (PVC) to be used, details above, needs to match the PVC's
metadata:name -->
<param id="k8s_persistent_volume_claim_mount_path">/scratch1/galaxy_data</param>
<!-- The mount path needs to be parent directory of the "file_path" and "new_file_path" paths
set in universe_wsgi.ini (or equivalent general galaxy config file). This is the mount path of the
PVC within the docker container that will be actually running the tool -->
<param id="k8s_namespace">galaxy-instanceA</param>
<!-- The namespace to be used on the Kubernetes cluster, if different from default, this needs to be set
accordingly in the PV and PVC detailed above -->
<param id="k8s_pod_retrials">4</param>
<!-- Allows pods to retry up to this number of times, before marking the galaxy job failed. k8s is a state
setter essentially, so by default it will try to take a job submitted to successful completion. A job
submits pods, until the number of successes (1 in this use case) is achieved, assuming that whatever is
making the pods fail will be fixed (such as a stale disk or a dead node that it is being restarted).
This option sets a limit of retrials, so that after that number of failed pods, the job is re-scaled to
zero (no execution) and the stderr/stdout of the k8s job is reported in galaxy (and the galaxy job set
to failed) -->
</plugin>
</plugins>
<handlers default="handlers">
@@ -495,6 +572,61 @@
<param id="docker_enabled" from_config="use_docker">false</param>
</destination>
<destination id="my-tool-container" runner="k8s">
<!-- For the kubernetes (k8s) runner, each container is a destination.
Make sure that the container is able to execute the calls that will be passed by the galaxy built
command. Most notably, containers that execute scripts through an interpreter in the form
Rscript my-script.R <arguments>
should have this wrapped as the container set working directory won't be the one actually used by
galaxy (galaxy creates a new working director and moves to it). Recommendation is hence to wrap this
type of calls on a shell script, and leave that script with execution privileges on the PATH of the
container:
RUN echo '#!/bin/bash' > /usr/local/bin/myScriptExec
RUN echo 'Rscript /path/to/my-script.r "$@"' >> /usr/local/bin/myScriptExec
RUN chmod a+x /usr/local/bin/myScriptExec
-->
<!-- The following four fields assemble the container's full name:
docker pull <repo>/<owner>/<image>:tag
-->
<param id="docker_repo_override">my-docker-registry.org</param>
<param id="docker_owner_override">superbioinfo</param>
<param id="docker_image_override">my-tool</param>
<param id="docker_tag_override">latest</param>
<!-- Alternatively you could specify a different type of container, such as rkt (not tested with Kubernetes)
<param id="rkt_repo_override">my-docker-registry.org</param>
<param id="rkt_owner_override">superbioinfo</param>
<param id="rkt_image_override">my-tool</param>
<param id="rkt_tag_override">latest</param>
-->
<!-- You can also allow the destination to accept the docker container set in the tool, and only fall into
the docker image set by this destination if the tool doesn't set a docker container, by using the
"default" suffix instead of "override".
<param id="docker_repo_default">my-docker-registry.org</param>
<param id="docker_owner_default">superbioinfo</param>
<param id="docker_image_default">my-tool</param>
<param id="docker_tag_default">latest</param>
-->
<param id="max_pod_retrials">3</param>
<!-- Allows pods to retry up to this number of times, before marking the galaxy job failed. k8s is a state
setter essentially, so by default it will try to take a job submitted to successful completion. A job
submits pods, until the number of successes (1 in this use case) is achieved, assuming that whatever is
making the pods fail will be fixed (such as a stale disk or a dead node that it is being restarted).
This option sets a limit of retrials, so that after that number of failed pods, the job is re-scaled to
zero (no execution) and the stderr/stdout of the k8s job is reported in galaxy (and the galaxy job set
to failed).
Overrides the runner config. (Not implemented yet)
-->
<!-- REQUIRED: To play nicely with the existing galaxy setup for containers. This could be set though
internally by the runner. -->
<param id="docker_enabled">true</param>
</destination>
<!-- Templatized destinations - macros can be used to create templated
destinations with reduced XML duplication. Here we are creating 4 destinations in 4 lines instead of 28 using the macros defined below.
-->
@@ -526,6 +658,11 @@
and pass to dynamic destination (as resource_params argument). -->
<tool id="longbar" destination="dynamic" resources="all" />
<tool id="baz" handler="special_handlers" destination="bigmem"/>
<!-- Finally for Kubernetes runner, the following connects a particular tool to be executed with
the container of choice in Kubernetes.
-->
<tool id="my-tool" destination="my-tool-container"/>
</tools>
<limits>
<!-- Certain limits can be defined. The 'concurrent_jobs' limits all
@@ -67,3 +67,7 @@ ecdsa==0.13
# Flexible BAM index naming
pysam==0.8.4+gx1
# Kubernetes runner pykube requirement
-e git+https://github.com/pcm32/pykube.git@feature/allMergedFeatures#egg=pykube
+4 -3
View File
@@ -20,6 +20,7 @@ def build_command(
runner,
job_wrapper,
container=None,
modify_command_for_container=True,
include_metadata=False,
include_work_dir_outputs=True,
create_tool_working_directory=True,
@@ -57,14 +58,14 @@ def build_command(
if not container:
__handle_dependency_resolution(commands_builder, job_wrapper, remote_command_params)
if container or job_wrapper.commands_in_new_shell:
if container:
if (container and modify_command_for_container) or job_wrapper.commands_in_new_shell:
if container and modify_command_for_container:
# Many Docker containers do not have /bin/bash.
external_command_shell = "/bin/sh"
else:
external_command_shell = shell
externalized_commands = __externalize_commands(job_wrapper, external_command_shell, commands_builder, remote_command_params)
if container:
if container and modify_command_for_container:
# Stop now and build command before handling metadata and copying
# working directory files back. These should always happen outside
# of docker container - no security implications when generating
+6 -2
View File
@@ -145,7 +145,8 @@ class BaseJobRunner( object ):
"""
raise NotImplementedError()
def prepare_job(self, job_wrapper, include_metadata=False, include_work_dir_outputs=True):
def prepare_job(self, job_wrapper, include_metadata=False, include_work_dir_outputs=True,
modify_command_for_container=True):
"""Some sanity checks that all runners' queue_job() methods are likely to want to do
"""
job_id = job_wrapper.get_id_tag()
@@ -171,6 +172,7 @@ class BaseJobRunner( object ):
job_wrapper,
include_metadata=include_metadata,
include_work_dir_outputs=include_work_dir_outputs,
modify_command_for_container=modify_command_for_container
)
except Exception as e:
log.exception("(%s) Failure preparing job" % job_id)
@@ -193,13 +195,15 @@ class BaseJobRunner( object ):
def recover(self, job, job_wrapper):
raise NotImplementedError()
def build_command_line( self, job_wrapper, include_metadata=False, include_work_dir_outputs=True ):
def build_command_line( self, job_wrapper, include_metadata=False, include_work_dir_outputs=True,
modify_command_for_container=True ):
container = self._find_container( job_wrapper )
return build_command(
self,
job_wrapper,
include_metadata=include_metadata,
include_work_dir_outputs=include_work_dir_outputs,
modify_command_for_container=modify_command_for_container,
container=container
)
+311
View File
@@ -0,0 +1,311 @@
"""
Offload jobs to a Kubernetes cluster.
"""
import logging
from galaxy import model
from galaxy.jobs.runners import AsynchronousJobState, AsynchronousJobRunner
from os import environ as os_environ
# pykube imports:
try:
import operator
from pykube.config import KubeConfig
from pykube.http import HTTPClient
from pykube.objects import Job
from pykube.objects import Pod
except ImportError as exc:
operator = None
K8S_IMPORT_MESSAGE = ('The Python pykube package is required to use '
'this feature, please install it or correct the '
'following error:\nImportError %s' % str(exc))
log = logging.getLogger(__name__)
__all__ = ['KubernetesJobRunner']
class KubernetesJobRunner(AsynchronousJobRunner):
"""
Job runner backed by a finite pool of worker threads. FIFO scheduling
"""
runner_name = "KubernetesRunner"
def __init__(self, app, nworkers, **kwargs):
# Check if pykube was importable, fail if not
assert operator is not None, K8S_IMPORT_MESSAGE
runner_param_specs = dict(
k8s_config_path=dict(map=str, default=os_environ.get('KUBECONFIG', None)),
k8s_persistent_volume_claim_name=dict(map=str),
k8s_persistent_volume_claim_mount_path=dict(map=str),
k8s_namespace=dict(map=str, default="default"),
k8s_pod_retrials=dict(map=int, valid=lambda x: int > 0, default=3))
if 'runner_param_specs' not in kwargs:
kwargs['runner_param_specs'] = dict()
kwargs['runner_param_specs'].update(runner_param_specs)
"""Start the job runner parent object """
super(KubernetesJobRunner, self).__init__(app, nworkers, **kwargs)
# self.cli_interface = CliInterface()
# here we need to fetch the default kubeconfig path from the plugin defined in job_conf...
self._pykube_api = HTTPClient(KubeConfig.from_file(self.runner_params["k8s_config_path"]))
self._galaxy_vol_name = "pvc-galaxy" # TODO this needs to be read from params!!
self._init_monitor_thread()
self._init_worker_threads()
def queue_job(self, job_wrapper):
"""Create job script and submit it to Kubernetes cluster"""
# prepare the job
# We currently don't need to include_metadata or include_work_dir_outputs, as working directory is the same
# were galaxy will expect results.
log.debug("Starting queue_job for job " + job_wrapper.get_id_tag())
if not self.prepare_job(job_wrapper, include_metadata=False, include_work_dir_outputs=False,
modify_command_for_container=False):
return
job_destination = job_wrapper.job_destination
# Construction of the Kubernetes Job object follows: http://kubernetes.io/docs/user-guide/persistent-volumes/
k8s_job_name = self.__produce_unique_k8s_job_name(job_wrapper)
k8s_job_obj = {
"apiVersion": "extensions/v1beta1",
"kind": "Job",
"metadata":
# metadata.name is the name of the pod resource created, and must be unique
# http://kubernetes.io/docs/user-guide/configuring-containers/
{
"name": k8s_job_name,
"namespace": "default", # TODO this should be set
"labels": {"app": k8s_job_name},
}
,
"spec": self.__get_k8s_job_spec(job_wrapper)
}
# Creates the Kubernetes Job
Job(self._pykube_api, k8s_job_obj).create()
# define job attributes in the AsyncronousJobState for follow-up
ajs = AsynchronousJobState(files_dir=job_wrapper.working_directory, job_wrapper=job_wrapper,
job_id=k8s_job_name, job_destination=job_destination)
self.monitor_queue.put(ajs)
# external_runJob_script can be None, in which case it's not used.
external_runjob_script = None
return external_runjob_script
def __produce_unique_k8s_job_name(self, job_wrapper):
# wrapper.get_id_tag() instead of job_id for compatibility with TaskWrappers.
return "galaxy-" + job_wrapper.get_id_tag()
def __get_k8s_job_spec(self, job_wrapper):
"""Creates the k8s Job spec. For a Job spec, the only requirement is to have a .spec.template."""
k8s_job_spec = {"template": self.__get_k8s_job_spec_template(job_wrapper)}
return k8s_job_spec
def __get_k8s_job_spec_template(self, job_wrapper):
"""The k8s spec template is nothing but a Pod spec, except that it is nested and does not have an apiversion
nor kind. In addition to required fields for a Pod, a pod template in a job must specify appropriate labels
(see pod selector) and an appropriate restart policy."""
k8s_spec_template = {
"metadata": {
"labels": {"app": self.__produce_unique_k8s_job_name(job_wrapper)}
},
"spec": {
"volumes": self.__get_k8s_mountable_volumes(job_wrapper),
"restartPolicy": self.__get_k8s_restart_policy(job_wrapper),
"containers": self.__get_k8s_containers(job_wrapper)
}
}
# TODO include other relevant elements that people might want to use from
# TODO http://kubernetes.io/docs/api-reference/v1/definitions/#_v1_podspec
return k8s_spec_template
def __get_k8s_restart_policy(self, job_wrapper):
"""The default Kubernetes restart policy for Jobs"""
return "Never"
def __get_k8s_mountable_volumes(self, job_wrapper):
"""Provides the required volumes that the containers in the pod should be able to mount. This should be using
the new persistent volumes and persistent volumes claim objects. This requires that both a PersistentVolume and
a PersistentVolumeClaim are created before starting galaxy (starting a k8s job).
"""
# TODO on this initial version we only support a single volume to be mounted.
k8s_mountable_volume = {
"name": self._galaxy_vol_name,
"persistentVolumeClaim": {
"claimName": self.runner_params['k8s_persistent_volume_claim_name']
}
}
return [k8s_mountable_volume]
def __get_k8s_containers(self, job_wrapper):
"""Fills in all required for setting up the docker containers to be used."""
k8s_container = {
"name": self.__get_k8s_container_name(job_wrapper),
"image": self._find_container(job_wrapper).container_id,
# this form of command overrides the entrypoint and allows multi command
# command line execution, separated by ;, which is what Galaxy does
# to assemble the command.
# TODO possibly shell needs to be set by job_wrapper
"command": ["/bin/bash", "-c", job_wrapper.runner_command_line],
"volumeMounts": [{
"mountPath": self.runner_params['k8s_persistent_volume_claim_mount_path'],
"name": self._galaxy_vol_name
}]
}
# if self.__requires_ports(job_wrapper):
# k8s_container['ports'] = self.__get_k8s_containers_ports(job_wrapper)
return [k8s_container]
# def __get_k8s_containers_ports(self, job_wrapper):
# for k,v self.runner_params:
# if k.startswith("container_port_"):
def __assemble_k8s_container_image_name(self, job_wrapper):
"""Assembles the container image name as repo/owner/image:tag, where repo, owner and tag are optional"""
job_destination = job_wrapper.job_destination
# Determine the job's Kubernetes destination (context, namespace) and options from the job destination
# definition
repo = ""
owner = ""
if 'repo' in job_destination.params:
repo = job_destination.params['repo'] + "/"
if 'owner' in job_destination.params:
owner = job_destination.params['owner'] + "/"
k8s_cont_image = repo + owner + job_destination.params['image']
if 'tag' in job_destination.params:
k8s_cont_image += ":" + job_destination.params['tag']
return k8s_cont_image
def __get_k8s_container_name(self, job_wrapper):
# TODO check if this is correct
return job_wrapper.job_destination.id
def check_watched_item(self, job_state):
"""Checks the state of a job already submitted on k8s. Job state is a AsynchronousJobState"""
jobs = Job.objects(self._pykube_api).filter(selector="app=" + job_state.job_id)
if len(jobs.response['items']) == 1:
job = Job(self._pykube_api, jobs.response['items'][0])
succeeded = 0
active = 0
failed = 0
if 'succeeded' in job.obj['status']:
succeeded = job.obj['status']['succeeded']
if 'active' in job.obj['status']:
active = job.obj['status']['active']
if 'failed' in job.obj['status']:
failed = job.obj['status']['failed']
# This assumes jobs dependent on a single pod, single container
if succeeded > 0:
self.__produce_log_file(job_state)
error_file = open(job_state.error_file, 'w')
error_file.write("")
error_file.close()
job_state.running = False
self.mark_as_finished(job_state)
return None
elif active > 0 or succeeded + active + failed == 0:
job_state.running = True
return job_state
elif failed > job_state.job_destination.params['max_pod_retrials']:
self.__produce_log_file(job_state)
error_file = open(job_state.error_file, 'w')
error_file.write("Exceeded max number of Kubernetes pod retrials allowed for job\n")
error_file.close()
job_state.running = False
self.mark_as_failed(job_state)
job.scale(replicas=0)
return None
# We should not get here
log.debug(
"Reaching unexpected point for Kubernetes job, where it is not classified as succ., active nor failed.")
return job_state
elif len(jobs.response['items']) == 0:
# there is no job responding to this job_id, it is either lost or something happened.
log.error("No Jobs are available under expected selector app=" + job_state.job_id)
error_file = open(job_state.error_file, 'w')
error_file.write("No Kubernetes Jobs are available under expected selector app=" + job_state.job_id + "\n")
error_file.close()
self.mark_as_failed(job_state)
return job_state
else:
# there is more than one job associated to the expected unique job id used as selector.
log.error("There is more than one Kubernetes Job associated to job id " + job_state.job_id)
self.__produce_log_file(job_state)
error_file = open(job_state.error_file, 'w')
error_file.write("There is more than one Kubernetes Job associated to job id " + job_state.job_id + "\n")
error_file.close()
self.mark_as_failed(job_state)
return job_state
def __produce_log_file(self, job_state):
pod_r = Pod.objects(self._pykube_api).filter(selector="app=" + job_state.job_id)
logs = ""
for pod_obj in pod_r.response['items']:
pod = Pod(self._pykube_api, pod_obj)
logs += "\n\n==== Pod " + pod.name + " log start ====\n\n"
logs += pod.logs(timestamps=True)
logs += "\n\n==== Pod " + pod.name + " log end ===="
logs_file_path = job_state.output_file
logs_file = open(logs_file_path, mode="w")
logs_file.write(logs)
logs_file.close()
return logs_file_path
def stop_job(self, job):
"""Attempts to delete a dispatched job to the k8s cluster"""
try:
jobs = Job.objects(self._pykube_api).filter(selector="app=" + job.job_runner_external_id)
if jobs.response['items'].len() >= 0:
job_to_delete = Job(self._pykube_api, jobs.response['items'][0])
job_to_delete.scale(replicas=0)
# TODO assert whether job parallelism == 0
# assert not job_to_delete.exists(), "Could not delete job,"+job.job_runner_external_id+" it still exists"
log.debug("(%s/%s) Terminated at user's request" % (job.id, job.job_runner_external_id))
except Exception as e:
log.debug("(%s/%s) User killed running job, but error encountered during termination: %s" % (
job.id, job.job_runner_external_id, e))
def recover(self, job, job_wrapper):
"""Recovers jobs stuck in the queued/running state when Galaxy started"""
# TODO this needs to be implemented to override unimplemented base method
job_id = job.get_job_runner_external_id()
if job_id is None:
self.put(job_wrapper)
return
ajs = AsynchronousJobState(files_dir=job_wrapper.working_directory, job_wrapper=job_wrapper)
ajs.job_id = str(job_id)
ajs.command_line = job.command_line
ajs.job_wrapper = job_wrapper
ajs.job_destination = job_wrapper.job_destination
if job.state == model.Job.states.RUNNING:
log.debug("(%s/%s) is still in running state, adding to the runner monitor queue" % (
job.id, job.job_runner_external_id))
ajs.old_state = model.Job.states.RUNNING
ajs.running = True
self.monitor_queue.put(ajs)
elif job.state == model.Job.states.QUEUED:
log.debug("(%s/%s) is still in queued state, adding to the runner monitor queue" % (
job.id, job.job_runner_external_id))
ajs.old_state = model.Job.states.QUEUED
ajs.running = False
self.monitor_queue.put(ajs)
+24 -2
View File
@@ -102,7 +102,25 @@ class ContainerFinder(object):
def __overridden_container_id(self, container_type, destination_info):
if not self.__container_type_enabled(container_type, destination_info):
return None
return destination_info.get("%s_container_id_override" % container_type)
if "%s_container_id_override" % container_type in destination_info:
return destination_info.get("%s_container_id_override" % container_type)
if "%s_image_override" % container_type in destination_info:
return self.__build_container_id_from_parts(container_type, destination_info, mode="override")
def __build_container_id_from_parts(self, container_type, destination_info, mode):
repo = ""
owner = ""
repo_key = "%s_repo_%s" % (container_type, mode)
owner_key = "%s_owner_%s" % (container_type, mode)
if repo_key in destination_info:
repo = destination_info[repo_key] + "/"
if owner_key in destination_info:
owner = destination_info[owner_key] + "/"
cont_id = repo + owner + destination_info["%s_image_%s" % (container_type, mode)]
tag_key = "%s_tag_%s" % (container_type, mode)
if tag_key in destination_info:
cont_id += ":" + destination_info[tag_key]
return cont_id
def __default_container_id(self, container_type, destination_info):
if not self.__container_type_enabled(container_type, destination_info):
@@ -111,7 +129,11 @@ class ContainerFinder(object):
# Also allow docker_image...
if key not in destination_info:
key = "%s_image" % container_type
return destination_info.get(key)
if key in destination_info:
return destination_info.get(key)
elif "%s_image_default" in destination_info:
return self.__build_container_id_from_parts(container_type, destination_info, mode="default")
return None
def __destination_container(self, container_id, container_type, tool_info, destination_info, job_info):
# TODO: ensure destination_info is dict-like