Merge pull request #8311 from telukir/f-kt-r19.05-jwt1

walltime_limit added for the k8s jobs
This commit is contained in:
Enis Afgan
2019-07-12 08:30:29 +02:00
committed by GitHub
3 changed files with 56 additions and 12 deletions
@@ -216,6 +216,11 @@
zero (no execution) and the stderr/stdout of the k8s job is reported in galaxy (and the galaxy job set
to failed) -->
<!-- <param id="k8s_walltime_limit">172800</param> -->
<!-- Controls the maximum time before kubernetes terminates the job. Sets activeDeadlineSeconds
in the kubernetes job spec. Once a Job reaches activeDeadlineSeconds, all of its running Pods are
terminated and the Job status will become `type: Failed` with `reason: DeadlineExceeded`. -->
<!-- <param id="k8s_galaxy_instance_id">my-instance</param> -->
<!-- Identifies the Galaxy instance where this runner belongs. Setting this variable means that the runner
will trust k8s Jobs with the structure galaxy-my-instance-<number> to be its own. This variable needs
+21 -12
View File
@@ -58,7 +58,8 @@ class KubernetesJobRunner(AsynchronousJobRunner):
k8s_default_limits_cpu=dict(map=str, default=None),
k8s_default_limits_memory=dict(map=str, default=None),
k8s_pod_retries=dict(map=int, valid=lambda x: int >= 0, default=3),
k8s_pod_retrials=dict(map=int, valid=lambda x: int >= 0, default=3))
k8s_pod_retrials=dict(map=int, valid=lambda x: int >= 0, default=3),
k8s_walltime_limit=dict(map=int, valid=lambda x: int(x) >= 0, default=172800))
if 'runner_param_specs' not in kwargs:
kwargs['runner_param_specs'] = dict()
@@ -220,8 +221,10 @@ class KubernetesJobRunner(AsynchronousJobRunner):
return produce_unique_k8s_job_name(app_prefix='galaxy', instance_id=instance_id, job_id=galaxy_internal_job_id)
def __get_k8s_job_spec(self, ajs):
"""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(ajs)}
"""Creates the k8s Job spec. For a Job spec, the only requirement is to have a .spec.template.
If the job hangs around unlimited it will be ended after k8s wall time limit, which sets activeDeadlineSeconds"""
k8s_job_spec = {"template": self.__get_k8s_job_spec_template(ajs),
"activeDeadlineSeconds": int(self.runner_params['k8s_walltime_limit'])}
return k8s_job_spec
def __get_k8s_job_spec_template(self, ajs):
@@ -427,23 +430,18 @@ class KubernetesJobRunner(AsynchronousJobRunner):
job_state.running = False
self.mark_as_finished(job_state)
return None
elif failed > 0 and self.__job_failed_due_to_low_memory(job_state):
return self._handle_job_failure(job, job_state, reason="OOM")
elif active > 0 and failed <= max_pod_retries:
if not job_state.running:
job_state.running = True
job_state.job_wrapper.change_state(model.Job.states.RUNNING)
return job_state
elif failed > max_pod_retries:
return self._handle_job_failure(job, job_state)
elif job_state.job_wrapper.get_job().state == model.Job.states.DELETED:
# Job has been deleted via stop_job, cleanup and remove from watched_jobs by returning `None`
if job_state.job_wrapper.cleanup_job in ("always", "onsuccess"):
job_state.job_wrapper.cleanup()
return None
else:
# We really shouldn't reach this point, but we might if the job has been killed by the kubernetes admin
log.info("Kubernetes job '%s' not classified as succ., active or failed. Full Job object: \n%s", job.name, job.obj)
return self._handle_job_failure(job, job_state)
elif len(jobs.response['items']) == 0:
# there is no job responding to this job_id, it is either lost or something happened.
@@ -460,12 +458,17 @@ class KubernetesJobRunner(AsynchronousJobRunner):
self.mark_as_failed(job_state)
return job_state
def _handle_job_failure(self, job, job_state, reason=None):
def _handle_job_failure(self, job, job_state):
# Figure out why job has failed
with open(job_state.error_file, 'a') as error_file:
if reason == "OOM":
if self.__job_failed_due_to_low_memory(job_state):
error_file.write("Job killed after running out of memory. Try with more memory.\n")
job_state.fail_message = "Tool failed due to insufficient memory. Try with more memory."
job_state.runner_state = JobState.runner_states.MEMORY_LIMIT_REACHED
elif self.__job_failed_due_to_walltime_limit(job):
error_file.write("DeadlineExceeded")
job_state.fail_message = "Job was active longer than specified deadline"
job_state.runner_state = JobState.runner_states.WALLTIME_REACHED
else:
error_file.write("Exceeded max number of Kubernetes pod retrials allowed for job\n")
job_state.fail_message = "More pods failed than allowed. See stdout for pods details."
@@ -474,6 +477,10 @@ class KubernetesJobRunner(AsynchronousJobRunner):
job.scale(replicas=0)
return None
def __job_failed_due_to_walltime_limit(self, job):
conditions = job.obj['status']['conditions']
return any(True for c in conditions if c['type'] == 'Failed' and c['reason'] == 'DeadlineExceeded')
def __job_failed_due_to_low_memory(self, job_state):
"""
checks the state of the pod to see if it was killed
@@ -483,8 +490,10 @@ class KubernetesJobRunner(AsynchronousJobRunner):
pods = Pod.objects(self._pykube_api).filter(selector="app=%s" % job_state.job_id,
namespace=self.runner_params['k8s_namespace'])
pod = Pod(self._pykube_api, pods.response['items'][0])
if not pods.response['items']:
return False
pod = Pod(self._pykube_api, pods.response['items'][0])
if pod.obj['status']['phase'] == "Failed" and \
pod.obj['status']['containerStatuses'][0]['state']['terminated']['reason'] == "OOMKilled":
return True
@@ -90,6 +90,12 @@ def job_config(jobs_directory):
<param id="k8s_config_path">$k8s_config_path</param>
<param id="k8s_galaxy_instance_id">gx-short-id</param>
</plugin>
<plugin id="k8s_walltime_short" type="runner" load="galaxy.jobs.runners.kubernetes:KubernetesJobRunner">
<param id="k8s_persistent_volume_claims">jobs-directory-claim:$jobs_directory,tool-directory-claim:$tool_directory</param>
<param id="k8s_config_path">$k8s_config_path</param>
<param id="k8s_galaxy_instance_id">gx-short-id</param>
<param id="k8s_walltime_limit">10</param>
</plugin>
</plugins>
<destinations default="k8s_destination">
<destination id="k8s_destination" runner="k8s">
@@ -99,11 +105,19 @@ def job_config(jobs_directory):
<param id="docker_default_container_id">busybox:ubuntu-14.04</param>
<env id="SOME_ENV_VAR">42</env>
</destination>
<destination id="k8s_destination_walltime_short" runner="k8s_walltime_short">
<param id="limits_cpu">1.9</param>
<param id="limits_memory">10M</param>
<param id="docker_enabled">true</param>
<param id="docker_default_container_id">busybox:ubuntu-14.04</param>
<env id="SOME_ENV_VAR">42</env>
</destination>
<destination id="local_dest" runner="local">
</destination>
</destinations>
<tools>
<tool id="upload1" destination="local_dest"/>
<tool id="create_2" destination="k8s_destination_walltime_short"/>
</tools>
</job_conf>
""")
@@ -255,3 +269,19 @@ class BaseKubernetesIntegrationTestCase(BaseJobEnvironmentIntegrationTestCase, M
MEM = '10'
MEM_PER_SLOT = '5'
assert [CPU, MEM, MEM_PER_SLOT] == dataset_content.split('\n'), dataset_content
@skip_without_tool('create_2')
def test_walltime_limit(self):
running_response = self.dataset_populator.run_tool(
'create_2',
{'sleep_time': 60},
self.history_id,
assert_ok=False
)
result = self.dataset_populator.wait_for_tool_run(run_response=running_response,
history_id=self.history_id,
assert_ok=False).json()
details = self.dataset_populator.get_job_details(result['jobs'][0]['id'], full=True).json()
assert details['state'] == 'error'
hda_details = self.dataset_populator.get_history_dataset_details(self.history_id, assert_ok=False)
assert hda_details['misc_info'] == 'Job was active longer than specified deadline'