Merge pull request #2559 from phnmnl/feature/failingPodsNotDetected

Kubernetes runner: improves handling of failed jobs
This commit is contained in:
John Chilton
2016-07-11 12:28:48 -04:00
committed by GitHub
+61 -13
View File
@@ -75,7 +75,7 @@ class KubernetesJobRunner(AsynchronousJobRunner):
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_name = self.__produce_unique_k8s_job_name(job_wrapper.get_id_tag())
k8s_job_obj = {
"apiVersion": "extensions/v1beta1",
"kind": "Job",
@@ -91,7 +91,14 @@ class KubernetesJobRunner(AsynchronousJobRunner):
"spec": self.__get_k8s_job_spec(job_wrapper)
}
# Checks if job exists
job = Job(self._pykube_api, k8s_job_obj)
if job.exists():
job.delete()
# Creates the Kubernetes Job
# TODO if a job with that ID exists, what should we do?
# TODO do we trust that this is the same job and use that?
# TODO or create a new job as we cannot make sure
Job(self._pykube_api, k8s_job_obj).create()
# define job attributes in the AsyncronousJobState for follow-up
@@ -103,9 +110,9 @@ class KubernetesJobRunner(AsynchronousJobRunner):
external_runjob_script = None
return external_runjob_script
def __produce_unique_k8s_job_name(self, job_wrapper):
def __produce_unique_k8s_job_name(self, galaxy_internal_job_id):
# wrapper.get_id_tag() instead of job_id for compatibility with TaskWrappers.
return "galaxy-" + job_wrapper.get_id_tag()
return "galaxy-" + galaxy_internal_job_id
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."""
@@ -118,7 +125,7 @@ class KubernetesJobRunner(AsynchronousJobRunner):
(see pod selector) and an appropriate restart policy."""
k8s_spec_template = {
"metadata": {
"labels": {"app": self.__produce_unique_k8s_job_name(job_wrapper)}
"labels": {"app": self.__produce_unique_k8s_job_name(job_wrapper.get_id_tag())}
},
"spec": {
"volumes": self.__get_k8s_mountable_volumes(job_wrapper),
@@ -204,9 +211,17 @@ class KubernetesJobRunner(AsynchronousJobRunner):
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])
job_destination = job_state.job_wrapper.job_destination
succeeded = 0
active = 0
failed = 0
max_pod_retrials = 1
if 'k8s_pod_retrials' in self.runner_params:
max_pod_retrials = int(self.runner_params['k8s_pod_retrials'])
if 'max_pod_retrials' in job_destination.params:
max_pod_retrials = int(job_destination.params['max_pod_retrials'])
if 'succeeded' in job.obj['status']:
succeeded = job.obj['status']['succeeded']
if 'active' in job.obj['status']:
@@ -223,16 +238,16 @@ class KubernetesJobRunner(AsynchronousJobRunner):
job_state.running = False
self.mark_as_finished(job_state)
return None
elif active > 0 or succeeded + active + failed == 0:
elif active > 0 and failed <= max_pod_retrials:
job_state.running = True
return job_state
elif failed > job_state.job_destination.params['max_pod_retrials']:
elif failed > 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
job_state.fail_message = "More pods failed than allowed. See stdout for pods details."
self.mark_as_failed(job_state)
job.scale(replicas=0)
return None
@@ -260,14 +275,46 @@ class KubernetesJobRunner(AsynchronousJobRunner):
self.mark_as_failed(job_state)
return job_state
def fail_job(self, job_state):
"""
Kubernetes runner overrides fail_job (called by mark_as_failed) to rescue the pod's log files which are left as
stdout (pods logs are the natural stdout and stderr of the running processes inside the pods) and are
deleted in the parent implementation as part of the failing the job process.
:param job_state:
:return:
"""
# First we rescue the pods logs
with open(job_state.output_file, 'r') as outfile:
stdout_content = outfile.read()
if getattr(job_state, 'stop_job', True):
self.stop_job(self.sa_session.query(self.app.model.Job).get(job_state.job_wrapper.job_id))
self._handle_runner_state('failure', job_state)
# Not convinced this is the best way to indicate this state, but
# something necessary
if not job_state.runner_state_handled:
job_state.job_wrapper.fail(
message=getattr(job_state, 'fail_message', 'Job failed'),
stdout=stdout_content, stderr='See stdout for pod\'s stderr.'
)
if job_state.job_wrapper.cleanup_job == "always":
job_state.cleanup()
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 ===="
try:
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 ===="
except Exception as detail:
log.info("Could not write pod\'s " + pod_obj['metadata']['name'] +
" log file due to HTTPError " + str(detail))
logs_file_path = job_state.output_file
logs_file = open(logs_file_path, mode="w")
if isinstance(logs, text_type):
@@ -279,8 +326,9 @@ class KubernetesJobRunner(AsynchronousJobRunner):
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:
jobs = Job.objects(self._pykube_api).filter(selector="app=" +
self.__produce_unique_k8s_job_name(job.get_id_tag()))
if len(jobs.response['items']) >= 0:
job_to_delete = Job(self._pykube_api, jobs.response['items'][0])
job_to_delete.scale(replicas=0)
# TODO assert whether job parallelism == 0