浏览代码

HUE-3356 [jb] Fall back to JHS if a job is not found in YARN RM

Jenny Kim 9 年之前
父节点
当前提交
b1866eb

+ 25 - 23
apps/jobbrowser/src/jobbrowser/api.py

@@ -19,6 +19,7 @@ import logging
 
 from desktop.lib.exceptions_renderable import PopupException
 from desktop.lib.paginator import Paginator
+from desktop.lib.rest.http_client import RestException
 
 from hadoop import cluster
 from hadoop.cluster import jt_ha, rm_ha
@@ -220,35 +221,36 @@ class YarnApi(JobBrowserApi):
 
   @rm_ha
   def get_job(self, jobid):
+
+    job_id = jobid.replace('application', 'job')
+    app_id = jobid.replace('job', 'application')
+
     try:
-      # App id
-      jobid = jobid.replace('job', 'application')
-      job = self.resource_manager_api.app(jobid)['app']
-
-      if job['state'] == 'ACCEPTED':
-        raise ApplicationNotRunning(jobid, job)
-      elif job.get('applicationType') == 'SPARK':
-        job = SparkJob(job, rm_api=self.resource_manager_api, hs_api=self.spark_history_server_api)
-      elif job['state'] == 'KILLED':
-        return KilledYarnJob(self.resource_manager_api, job)
-      elif job.get('applicationType') == 'MAPREDUCE':
-        jobid = jobid.replace('application', 'job')
-
-        if job['state'] in ('NEW', 'SUBMITTED', 'ACCEPTED', 'RUNNING'):
-          json = self.mapreduce_api.job(self.user, jobid)
-          job = YarnJob(self.mapreduce_api, json['job'])
+      app = self.resource_manager_api.app(app_id)['app']
+
+      if app['finalStatus'] in ('SUCCEEDED', 'FAILED', 'KILLED'):
+        if app['applicationType'] == 'SPARK':
+          job = SparkJob(app, rm_api=self.resource_manager_api, hs_api=self.spark_history_server_api)
+        elif app['state'] == 'KILLED':
+          job = KilledYarnJob(self.resource_manager_api, app)
         else:
-          json = self.history_server_api.job(self.user, jobid)
-          job = YarnJob(self.history_server_api, json['job'])
+          resp = self.history_server_api.job(self.user, job_id)
+          job = YarnJob(self.history_server_api, resp['job'])
       else:
-        job = Application(job, self.resource_manager_api)
+        if app['state'] == 'ACCEPTED':
+          raise ApplicationNotRunning(app_id, app)
+        else:
+          job = Application(app, self.resource_manager_api)
+    except RestException, e:
+      if e.code == 404:  # Job not found in RM so attempt to find job in History Server
+        resp = self.history_server_api.job(self.user, job_id)
+        job = YarnJob(self.history_server_api, resp['job'])
+      else:
+        raise JobExpired(app_id)
     except ApplicationNotRunning, e:
       raise e
     except Exception, e:
-      if 'NotFoundException' in str(e):
-        raise JobExpired(jobid)
-      else:
-        raise PopupException('Job %s could not be found: %s' % (jobid, e), detail=e)
+      raise PopupException('Job %s could not be found: %s' % (jobid, e), detail=e)
 
     return job
 

+ 0 - 2
apps/jobbrowser/src/jobbrowser/views.py

@@ -15,7 +15,6 @@
 # See the License for the specific language governing permissions and
 # limitations under the License.
 
-import json
 import re
 import time
 import logging
@@ -44,7 +43,6 @@ from hadoop.yarn.clients import get_log_client
 import hadoop.yarn.resource_manager_api as resource_manager_api
 
 from jobbrowser.conf import SHARE_JOBS
-from jobbrowser.conf import DISABLE_KILLING_JOBS
 from jobbrowser.api import get_api, ApplicationNotRunning, JobExpired
 from jobbrowser.models import Job, JobLinkage, Tracker, Cluster, can_view_job, can_modify_job, LinkJobLogs, can_kill_job
 from jobbrowser.yarn_models import Application

+ 7 - 0
apps/jobbrowser/src/jobbrowser/yarn_models.py

@@ -46,6 +46,13 @@ class Application(object):
 
     self._fixup()
 
+  @property
+  def logs_url(self):
+    url = self.trackingUrl
+    if self.applicationType == 'SPARK':
+      url = os.path.join(self.trackingUrl, 'executors')
+    return url
+
   def _fixup(self):
     self.is_mr2 = True
     jobid = self.id