瀏覽代碼

HUE-8713 [jb] Fetch log name list dynamically for Yarn jobs.

jdesjean 6 年之前
父節點
當前提交
8542cef2ae

+ 5 - 8
apps/jobbrowser/src/jobbrowser/api.py

@@ -33,7 +33,7 @@ import hadoop.yarn.resource_manager_api as resource_manager_api
 import hadoop.yarn.spark_history_server_api as spark_history_server_api
 
 from jobbrowser.conf import SHARE_JOBS
-from jobbrowser.yarn_models import Application, OozieYarnJob, Job as YarnJob, KilledJob as KilledYarnJob, Container, SparkJob
+from jobbrowser.yarn_models import Application, YarnV2Job, Job as YarnJob, KilledJob as KilledYarnJob, Container, SparkJob
 from desktop.auth.backend import is_admin
 
 
@@ -144,20 +144,17 @@ class YarnApi(JobBrowserApi):
 
     try:
       app = self.resource_manager_api.app(app_id)['app']
-
-      if app['finalStatus'] in ('SUCCEEDED', 'FAILED', 'KILLED'):
+      if app['applicationType'] == 'Oozie Launcher' or app['applicationType'] == 'TEZ' or app['applicationType'] == 'yarn-service':
+        job = YarnV2Job(self.resource_manager_api, app)
+      elif 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'] in ('KILLED', 'FAILED'):
           job = KilledYarnJob(self.resource_manager_api, app)
-        elif app['applicationType'] == 'Oozie Launcher' or app['applicationType'] == 'TEZ':
-          job = OozieYarnJob(self.resource_manager_api, app)
         else:  # Job succeeded, attempt to fetch from JHS
           job = self._get_job_from_history_server(job_id)
       else:
-        if app['applicationType'] == 'Oozie Launcher' or app['applicationType'] == 'TEZ':
-          job = OozieYarnJob(self.resource_manager_api, app)
-        elif app['state'] == 'ACCEPTED':
+        if app['state'] == 'ACCEPTED':
           raise ApplicationNotRunning(app_id, app)
         # The MapReduce API only returns JSON when the application is in a RUNNING state
         elif app['state'] in ('NEW', 'SUBMITTED', 'RUNNING') and app['applicationType'] == 'MAPREDUCE':

+ 8 - 5
apps/jobbrowser/src/jobbrowser/apis/job_api.py

@@ -191,7 +191,7 @@ class YarnApi(Api):
       }
       if hasattr(job, 'metrics'):
         common['metrics'] = job.metrics
-    elif app['applicationType'] == 'Oozie Launcher' or app['applicationType'] == 'TEZ':
+    elif app['applicationType'] == 'YarnV2':
       common['properties'] = {
         'startTime': job.startTime,
         'finishTime': job.finishTime,
@@ -220,11 +220,14 @@ class YarnApi(Api):
 
   def logs(self, appid, app_type, log_name, is_embeddable=False):
     logs = ''
+    logs_list = []
     try:
-      if app_type == 'TEZ' or app_type == 'MAPREDUCE' or app_type == 'Oozie Launcher':
+      if app_type == 'YarnV2' or app_type == 'MAPREDUCE':
         if log_name == 'default':
           response = job_single_logs(MockDjangoRequest(self.user), job=appid)
-          logs = json.loads(response.content).get('logs')
+          parseResponse = json.loads(response.content)
+          logs = parseResponse.get('logs')
+          logs_list = parseResponse.get('logsList')
           if logs and len(logs) == 4:
             logs = logs[1]
         else:
@@ -237,7 +240,7 @@ class YarnApi(Api):
         logs = None
     except PopupException, e:
       LOG.warn('No task attempt found for logs: %s' % smart_str(e))
-    return {'logs': logs}
+    return {'logs': logs, 'logsList': logs_list}
 
 
   def profile(self, appid, app_type, app_property, app_filters):
@@ -257,7 +260,7 @@ class YarnApi(Api):
           'executor_list': NativeYarnApi(self.user).get_job(jobid=appid).get_executors(),
           'filter_text': ''
         }
-    elif app_type == 'Oozie Launcher' or app_type == 'TEZ':
+    elif app_type == 'YarnV2':
       if app_property == 'attempts':
         return {
           'task_list': NativeYarnApi(self.user).get_job(jobid=appid).job_attempts['jobAttempt'],

+ 17 - 15
apps/jobbrowser/src/jobbrowser/templates/job_browser.mako

@@ -589,12 +589,12 @@ ${ commonheader("Job Browser", "jobbrowser", user, request) | n,unicode }
     <div data-bind="template: { name: 'job-yarn-page${ SUFFIX }', data: $root.job() }"></div>
   <!-- /ko -->
 
-  <!-- ko if: type() == 'Oozie Launcher' || type() == 'TEZ' -->
-    <div data-bind="template: { name: 'job-oozie-page${ SUFFIX }', data: $root.job() }"></div>
+  <!-- ko if: type() == 'YarnV2' -->
+    <div data-bind="template: { name: 'job-yarnv2-page${ SUFFIX }', data: $root.job() }"></div>
   <!-- /ko -->
 
-  <!-- ko if: type() == 'Oozie Launcher_ATTEMPT' -->
-    <div data-bind="template: { name: 'job-oozie-attempt-page${ SUFFIX }', data: $root.job() }"></div>
+  <!-- ko if: type() == 'YarnV2_ATTEMPT' -->
+    <div data-bind="template: { name: 'job-yarnv2-attempt-page${ SUFFIX }', data: $root.job() }"></div>
   <!-- /ko -->
 
   <!-- ko if: type() == 'SPARK' -->
@@ -857,7 +857,7 @@ ${ commonheader("Job Browser", "jobbrowser", user, request) | n,unicode }
   </div>
 </script>
 
-<script type="text/html" id="job-oozie-page${ SUFFIX }">
+<script type="text/html" id="job-yarnv2-page${ SUFFIX }">
   <div class="row-fluid">
     <div data-bind="css:{'span2': !$root.isMini(), 'span12': $root.isMini() }">
       <div class="sidebar-nav">
@@ -896,21 +896,19 @@ ${ commonheader("Job Browser", "jobbrowser", user, request) | n,unicode }
     <div data-bind="css: {'span10': !$root.isMini(), 'span12': $root.isMini() }">
 
       <ul class="nav nav-pills margin-top-20">
-        <li class="active"><a class="jb-logs-link" href="#job-oozie-page-logs${ SUFFIX }" data-toggle="tab">${ _('Logs') }</a></li>
-        <li><a href="#job-oozie-page-attempts${ SUFFIX }" data-bind="click: function(){ fetchProfile('attempts'); $('a[href=\'#job-oozie-page-attempts${ SUFFIX }\']').tab('show'); }">${ _('Attempts') }</a></li>
+        <li class="active"><a class="jb-logs-link" href="#job-yarnv2-page-logs${ SUFFIX }" data-toggle="tab">${ _('Logs') }</a></li>
+        <li><a href="#job-yarnv2-page-attempts${ SUFFIX }" data-bind="click: function(){ fetchProfile('attempts'); $('a[href=\'#job-yarnv2-page-attempts${ SUFFIX }\']').tab('show'); }">${ _('Attempts') }</a></li>
       </ul>
 
       <div class="tab-content">
-        <div class="tab-pane active" id="job-oozie-page-logs${ SUFFIX }">
-          <ul class="nav nav-tabs">
-          % for name in ['stdout', 'stderr', 'syslog']:
-            <li class="${ name == 'stdout' and 'active' or '' }"><a href="javascript:void(0)" data-bind="click: function(data, e) { $(e.currentTarget).parent().siblings().removeClass('active'); $(e.currentTarget).parent().addClass('active'); fetchLogs('${ name }'); logActive('${ name }'); }, text: '${ name }'"></a></li>
-          % endfor
+        <div class="tab-pane active" id="job-yarnv2-page-logs${ SUFFIX }">
+          <ul class="nav nav-tabs" data-bind="foreach: logsList">
+            <li data-bind="css: { 'active': $data == $parent.logActive() }"><a href="javascript:void(0)" data-bind="click: function(data, e) { $parent.fetchLogs($data); $parent.logActive($data); }, text: $data"></a></li>
           </ul>
           <pre data-bind="html: logs, logScroller: logs"></pre>
         </div>
 
-        <div class="tab-pane" id="job-oozie-page-attempts${ SUFFIX }">
+        <div class="tab-pane" id="job-yarnv2-page-attempts${ SUFFIX }">
           <table class="table table-condensed">
             <thead>
             <tr>
@@ -945,7 +943,7 @@ ${ commonheader("Job Browser", "jobbrowser", user, request) | n,unicode }
   </div>
 </script>
 
-<script type="text/html" id="job-oozie-attempt-page${ SUFFIX }">
+<script type="text/html" id="job-yarnv2-attempt-page${ SUFFIX }">
   <div class="row-fluid">
     <div data-bind="css:{'span2': !$root.isMini(), 'span12': $root.isMini() }">
       <div class="sidebar-nav">
@@ -2362,6 +2360,7 @@ ${ commonheader("Job Browser", "jobbrowser", user, request) | n,unicode }
 
       self.logActive = ko.observable('default');
       self.logsByName = ko.observable({});
+      self.logsList = ko.observable(['default', 'stdout', 'stderr', 'syslog']);
       self.logs = ko.pureComputed(function() {
         return self.logsByName()[self.logActive()];
       });
@@ -2444,7 +2443,7 @@ ${ commonheader("Job Browser", "jobbrowser", user, request) | n,unicode }
       self.rerunModalContent = ko.observable('');
 
       self.hasKill = ko.pureComputed(function() {
-        return self.type() && (['MAPREDUCE', 'SPARK', 'workflow', 'schedule', 'bundle', 'QUERY', 'TEZ', 'Oozie Launcher'].indexOf(self.type()) != -1 || self.type().indexOf('Data Warehouse') != -1 || self.type().indexOf('Altus') != -1);
+        return self.type() && (['MAPREDUCE', 'SPARK', 'workflow', 'schedule', 'bundle', 'QUERY', 'TEZ', 'YarnV2'].indexOf(self.type()) != -1 || self.type().indexOf('Data Warehouse') != -1 || self.type().indexOf('Altus') != -1);
       });
       self.killEnabled = ko.pureComputed(function() {
         // Impala can kill queries that are finished, but not yet terminated
@@ -2669,6 +2668,9 @@ ${ commonheader("Job Browser", "jobbrowser", user, request) | n,unicode }
             var result = self.logsByName();
             result[name] = data.logs.logs;
             self.logsByName(result);
+            if (data.logs.logsList && data.logs.logsList.length) {
+              self.logsList(['default'].concat(data.logs.logsList));
+            }
             if ($('.jb-panel pre:visible').length > 0){
               $('.jb-panel pre:visible').css('overflow-y', 'auto').height(Math.max(200, $(window).height() - $('.jb-panel pre:visible').offset().top - $('.page-content').scrollTop() - 75));
             }

+ 8 - 8
apps/jobbrowser/src/jobbrowser/views.py

@@ -357,7 +357,7 @@ def job_attempt_logs_json(request, job, attempt_index=0, name='syslog', offset=L
           log_link = log_link.replace(attempt['nodeHttpAddress'], attempt['nodeId'])
       elif app['state'] == 'RUNNING':
         log_link = app['amContainerLogs']
-    elif app['applicationType'] == 'Oozie Launcher' or app['applicationType'] == 'TEZ':
+    elif app.get('amContainerLogs'):
       log_link = app.get('amContainerLogs')
   except (KeyError, RestException), e:
     raise KeyError(_("Cannot find job attempt '%(id)s'.") % {'id': job.jobId}, e)
@@ -366,12 +366,10 @@ def job_attempt_logs_json(request, job, attempt_index=0, name='syslog', offset=L
 
   if log_link:
     link = '/%s/' % name
-    if app['applicationType'] == 'Oozie Launcher' and app['state'] != 'FINISHED': # Yarn currently dumps with 500 error with doas in running state
-      params = {}
-    else:
-      params = {
-        'doAs': request.user.username
-      }
+    params = {
+      'doAs': request.user.username
+    }
+      
     if offset != 0:
       params['start'] = offset
 
@@ -401,7 +399,7 @@ def job_attempt_logs_json(request, job, attempt_index=0, name='syslog', offset=L
 @check_job_permission
 def job_single_logs(request, job, offset=LOG_OFFSET_BYTES):
   """
-  Try to smartly detect the most useful task attempt (e.g. Oozie launcher, failed task) and get its MR logs.
+  Try to smartly detect the most useful task attempt (e.g. YarnV2, failed task) and get its MR logs.
   """
   def cmp_exec_time(task1, task2):
     return cmp(task1.execStartTimeMs, task2.execStartTimeMs)
@@ -547,6 +545,7 @@ def single_task_attempt_logs(request, job, taskid, attemptid, offset=LOG_OFFSET_
       "joblnk": job_link,
       "task": task,
       "logs": logs,
+      "logs_list": attempt.get_log_list(),
       "first_log_tab": first_log_tab,
   }
 
@@ -558,6 +557,7 @@ def single_task_attempt_logs(request, job, taskid, attemptid, offset=LOG_OFFSET_
   if request.GET.get('format') == 'json':
     response = {
       "logs": context['logs'],
+      "logsList": context['logs_list'],
       "isRunning": job.status.lower() in ('running', 'pending', 'prep')
     }
     return JsonResponse(response)

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

@@ -336,7 +336,7 @@ class Job(object):
       self._job_attempts = self.api.job_attempts(self.id)['jobAttempts']
     return self._job_attempts
 
-class OozieYarnJob(Job):
+class YarnV2Job(Job):
   def __init__(self, api, attrs):
     self.api = api
     for attr in attrs.keys():
@@ -355,6 +355,8 @@ class OozieYarnJob(Job):
       setattr(self, 'status', self.finalStatus)
     else:
       setattr(self, 'status', self.state)
+    setattr(self, 'type', self.applicationType)
+    setattr(self, 'applicationType', 'YarnV2')
     setattr(self, 'jobName', self.name)
     setattr(self, 'jobId', jobid)
     setattr(self, 'jobId_short', self.jobId.replace('job_', ''))
@@ -416,7 +418,7 @@ class YarnTask:
 
   def get_attempt(self, attempt_id):
     json = self.job.api.appattempts_attempt(self.job.id, attempt_id)
-    return YarnOozieAttempt(self, json)
+    return YarnV2Attempt(self, json)
 
 class KilledJob(Job):
 
@@ -539,12 +541,11 @@ class Attempt:
       self._counters = self.task.job.api.task_attempt_counters(self.task.jobId, self.task.id, self.id)['jobTaskAttemptCounters']
     return self._counters
 
-  def get_task_log(self, offset=0):
-    logs = []
+  def get_log_link(self):
     attempt = self.task.job.job_attempts['jobAttempt'][-1]
     log_link = attempt['logsLink']
     if not log_link:
-      return ['', '', '']
+      return log_link
 
     # Generate actual task log link from logsLink url
     if self.task.job.status in ('NEW', 'SUBMITTED', 'RUNNING'):
@@ -596,6 +597,29 @@ class Attempt:
         'user': user
       }
 
+    return log_link, user
+
+  def get_log_list(self):
+    log_link, user = self.get_log_link()
+    if not log_link:
+      return []
+    params = {
+      'doAs': user
+    }
+    log_link = re.sub('job_[^/]+', str(self.id), log_link)
+    root = Resource(get_log_client(log_link), urlparse.urlsplit(log_link)[2], urlencode=False)
+    response = root.get('/', params=params)
+    links = html.fromstring(response, parser=html.HTMLParser()).xpath('/html/body/table/tbody/tr/td[2]//a/@href')
+    parsed_links = map(lambda x: urlparse.urlsplit(x), links)
+    return map(lambda x: x and len(x) >= 2 and x[2].split('/')[-2] or '', parsed_links)
+
+  def get_task_log(self, offset=0):
+    logs = []
+
+    log_link, user = self.get_log_link()
+    if not log_link:
+      return ['', '', '']
+
     for name in ('stdout', 'stderr', 'syslog'):
       link = '/%s/' % name
       if self.type == 'Oozie Launcher' and not self.task.job.status == 'FINISHED': # Yarn currently dumps with 500 error with doas in running state
@@ -630,7 +654,7 @@ class Attempt:
 
     return logs + [''] * (3 - len(logs))
 
-class YarnOozieAttempt(Attempt):
+class YarnV2Attempt(Attempt):
   def __init__(self, task, attrs):
     self.task = task
     if attrs:
@@ -642,7 +666,7 @@ class YarnOozieAttempt(Attempt):
   def _fixup(self):
     if not hasattr(self, 'diagnostics'):
       self.diagnostics = ''
-    setattr(self, 'type', 'Oozie Launcher')
+    setattr(self, 'type', 'YarnV2')
     if self.finishedTime == 0:
       finishTime = int(time.time() * 1000)
     else: