Преглед изворни кода

HUE-8713 [jb] Support TEZ jobs.

jdesjean пре 6 година
родитељ
комит
403c6f8

+ 2 - 2
apps/jobbrowser/src/jobbrowser/api.py

@@ -150,12 +150,12 @@ class YarnApi(JobBrowserApi):
           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':
+        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':
+        if app['applicationType'] == 'Oozie Launcher' or app['applicationType'] == 'TEZ':
           job = OozieYarnJob(self.resource_manager_api, app)
         elif app['state'] == 'ACCEPTED':
           raise ApplicationNotRunning(app_id, app)

+ 85 - 6
apps/jobbrowser/src/jobbrowser/apis/job_api.py

@@ -22,6 +22,7 @@ from django.utils.encoding import smart_str
 from django.utils.translation import ugettext as _
 from hadoop.yarn import resource_manager_api
 
+from desktop.lib.django_util import JsonResponse
 from desktop.lib.exceptions import MessageException
 from desktop.lib.exceptions_renderable import PopupException
 from jobbrowser.conf import MAX_JOB_FETCH, LOG_OFFSET
@@ -68,8 +69,10 @@ class JobApi(Api):
       return self.yarn_api
     elif appid.startswith('task_'):
       return YarnMapReduceTaskApi(self.user, appid)
-    elif appid.startswith('attempt_') or appid.startswith('appattempt_'):
+    elif appid.startswith('attempt_'):
       return YarnMapReduceTaskAttemptApi(self.user, appid)
+    elif appid.startswith('appattempt_'):
+      return YarnAttemptApi(self.user, appid)
     elif appid.find('_executor_') > 0:
       return SparkExecutorApi(self.user, appid)
     else:
@@ -188,12 +191,13 @@ class YarnApi(Api):
       }
       if hasattr(job, 'metrics'):
         common['metrics'] = job.metrics
-    elif app['applicationType'] == 'Oozie Launcher':
+    elif app['applicationType'] == 'Oozie Launcher' or app['applicationType'] == 'TEZ':
       common['properties'] = {
         'startTime': job.startTime,
         'finishTime': job.finishTime,
         'elapsedTime': job.duration,
-        'attempts': []
+        'attempts': [],
+        'diagnostics': job.diagnostics
       }
 
     return common
@@ -204,7 +208,9 @@ class YarnApi(Api):
       kills = []
       for app_id in app_ids:
         try:
-          kill_job(MockDjangoRequest(self.user), job=app_id)
+          response = kill_job(MockDjangoRequest(self.user), job=app_id)
+          if isinstance(response, JsonResponse) and json.loads(response.content).get('status') == 0:
+             kills.append(app_id)
         except MessageException:
           kills.append(app_id)
       return {'kills': kills, 'status': len(app_ids) - len(kills), 'message': _('Stop signal sent to %s') % kills}
@@ -215,7 +221,7 @@ class YarnApi(Api):
   def logs(self, appid, app_type, log_name, is_embeddable=False):
     logs = ''
     try:
-      if app_type == 'MAPREDUCE' or app_type == 'Oozie Launcher':
+      if app_type == 'TEZ' or app_type == 'MAPREDUCE' or app_type == 'Oozie Launcher':
         if log_name == 'default':
           response = job_single_logs(MockDjangoRequest(self.user), job=appid)
           logs = json.loads(response.content).get('logs')
@@ -251,7 +257,7 @@ class YarnApi(Api):
           'executor_list': NativeYarnApi(self.user).get_job(jobid=appid).get_executors(),
           'filter_text': ''
         }
-    elif app_type == 'Oozie Launcher':
+    elif app_type == 'Oozie Launcher' or app_type == 'TEZ':
       if app_property == 'attempts':
         return {
           'task_list': NativeYarnApi(self.user).get_job(jobid=appid).job_attempts['jobAttempt'],
@@ -267,6 +273,79 @@ class YarnApi(Api):
     else:
       return 'FAILED' # FAILED, KILLED
 
+class YarnAttemptApi(Api):
+
+  def __init__(self, user, app_id):
+    Api.__init__(self, user)
+    start = 'appattempt_' if app_id.startswith('appattempt_') else 'attempt_'
+    self.app_id = '_'.join(app_id.replace('task_', 'application_').replace(start, 'application_').split('_')[:3])
+    self.task_id = '_'.join(app_id.replace(start, 'task_').split('_')[:5])
+    self.attempt_id = app_id.split('_')[3]
+
+
+  def apps(self):
+    attempts = NativeYarnApi(self.user).get_task(jobid=self.app_id, task_id=self.task_id).attempts
+
+    return {
+      'apps': [self._massage_task(task) for task in attempts],
+      'total': len(attempts)
+    }
+
+
+  def app(self, appid):
+    task = NativeYarnApi(self.user).get_task(jobid=self.app_id, task_id=self.task_id).get_attempt(self.attempt_id)
+
+    common = self._massage_task(task)
+    common['properties'] = {
+        'metadata': [],
+        'counters': []
+    }
+    common['properties'].update(self._massage_task(task))
+
+    return common
+
+
+  def logs(self, appid, app_type, log_name, is_embeddable=False):
+    if log_name == 'default':
+      log_name = 'stdout'
+
+    task = NativeYarnApi(self.user).get_task(jobid=self.app_id, task_id=self.task_id).get_attempt(self.attempt_id)
+    stdout, stderr, syslog = task.get_task_log()
+
+    return {'progress': 0, 'logs': syslog if log_name == 'syslog' else stderr if log_name == 'stderr' else stdout}
+
+
+  def profile(self, appid, app_type, app_property, app_filters):
+    if app_property == 'counters':
+      return NativeYarnApi(self.user).get_task(jobid=self.app_id, task_id=self.task_id).get_attempt(self.attempt_id).counters
+
+    return {}
+
+
+  def _massage_task(self, task):
+    return {
+        #"elapsedMergeTime" : task.elapsedMergeTime,
+        #"shuffleFinishTime" : task.shuffleFinishTime,
+        'id': task.appAttemptId if hasattr(task, 'appAttemptId') else '',
+        'appAttemptId': task.appAttemptId if hasattr(task, 'appAttemptId') else '',
+        'blacklistedNodes': task.blacklistedNodes if hasattr(task, 'blacklistedNodes') else '',
+        'containerId' : task.containerId if hasattr(task, 'containerId') else '',
+        'diagnostics': task.diagnostics if hasattr(task, 'diagnostics') else '',
+        "startTimeFormatted" : task.startTimeFormatted if hasattr(task, 'startTimeFormatted') else '',
+        "startTime" : long(task.startTime) if hasattr(task, 'startTime') else '',
+        "finishTime" : long(task.finishedTime) if hasattr(task, 'finishedTime') else '',
+        "finishTimeFormatted" : task.finishTimeFormatted if hasattr(task, 'finishTimeFormatted') else '',
+        "type" : task.type + '_ATTEMPT' if hasattr(task, 'type') else '',
+        'nodesBlacklistedBySystem': task.nodesBlacklistedBySystem if hasattr(task, 'nodesBlacklistedBySystem') else '',
+        'nodeId': task.nodeId if hasattr(task, 'nodeId') else '',
+        'nodeHttpAddress': task.nodeHttpAddress if hasattr(task, 'nodeHttpAddress') else '',
+        'logsLink': task.logsLink if hasattr(task, 'logsLink') else '',
+        "app_id": self.app_id,
+        "task_id": self.task_id,
+        'duration' : task.duration if hasattr(task, 'duration') else '',
+        'durationFormatted' : task.duration if hasattr(task, 'durationFormatted') else '',
+        'state': task.status if hasattr(task, 'status') else ''
+    }
 
 class YarnMapReduceTaskApi(Api):
 

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

@@ -595,7 +595,7 @@ ${ 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' -->
+  <!-- ko if: type() == 'Oozie Launcher' || type() == 'TEZ' -->
     <div data-bind="template: { name: 'job-oozie-page${ SUFFIX }', data: $root.job() }"></div>
   <!-- /ko -->
 
@@ -610,8 +610,8 @@ ${ commonheader("Job Browser", "jobbrowser", user, request) | n,unicode }
   <!-- ko if: type() == 'SPARK_EXECUTOR' -->
     <div data-bind="template: { name: 'job-spark-executor-page${ SUFFIX }', data: $root.job() }"></div>
   <!-- /ko -->
-</script>
 
+</script>
 
 <script type="text/html" id="job-yarn-page${ SUFFIX }">
   <div class="row-fluid">
@@ -889,6 +889,10 @@ ${ commonheader("Job Browser", "jobbrowser", user, request) | n,unicode }
             <li><span data-bind="moment: {data: finishTime, format: 'LLL'}"></span></li>
             <li class="nav-header">${ _('Elapsed time') }</li>
             <li><span data-bind="text: elapsedTime().toHHMMSS()"></span></li>
+            <!-- ko if: diagnostics -->
+              <li class="nav-header">${ _('Diagnostics') }</li>
+              <li><span data-bind="text: diagnostics, attr: { title: diagnostics }"></span>%</li>
+            <!-- /ko -->
           <!-- /ko -->
           <!-- /ko -->
         </ul>
@@ -905,7 +909,7 @@ ${ commonheader("Job Browser", "jobbrowser", user, request) | n,unicode }
       <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']:
+          % 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
           </ul>
@@ -927,7 +931,7 @@ ${ commonheader("Job Browser", "jobbrowser", user, request) | n,unicode }
             </tr>
             </thead>
             <tbody data-bind="foreach: properties['attempts']()['task_list']">
-              <tr class="pointer" data-bind="click: function() { $root.job().id(id); $root.job().fetchJob(); }">
+              <tr class="pointer" data-bind="click: function() { $root.job().id(appAttemptId); $root.job().fetchJob(); }">
                 <td data-bind="text: containerId"></td>
                 <td data-bind="text: nodeId"></td>
                 <td data-bind="text: id"></td>
@@ -952,30 +956,23 @@ ${ commonheader("Job Browser", "jobbrowser", user, request) | n,unicode }
     <div data-bind="css:{'span2': !$root.isMini(), 'span12': $root.isMini() }">
       <div class="sidebar-nav">
         <ul class="nav nav-list">
-          <li class="nav-header">${ _('Id') }</li>
-          <li class="break-word"><span data-bind="text: id"></span></li>
-          <li class="nav-header">${ _('Type') }</li>
-          <li><span data-bind="text: type"></span></li>
           <!-- ko with: properties -->
+          <li class="nav-header">${ _('Attempt Id') }</li>
+          <li class="break-word"><span data-bind="text: appAttemptId"></span></li>
           <li class="nav-header">${ _('State') }</li>
           <li><span data-bind="text: state"></span></li>
           <!-- ko if: !$root.isMini() -->
-          <li class="nav-header">${ _('Finish time') }</li>
-          <li><span data-bind="moment: {data: finishTime, format: 'LLL'}"></span></li>
-          <li class="nav-header">${ _('Assigned Container ID') }</li>
-          <li><span data-bind="text: assignedContainerId"></span></li>
-          <li class="nav-header">${ _('Host') }</li>
-          <li><span data-bind="text: host"></span></li>
-          <li class="nav-header">${ _('RPC Port') }</li>
-          <li><span data-bind="text: rpcPort"></span></li>
-          <li class="nav-header">${ _('Diagnostics Info') }</li>
-          <li><span data-bind="text: diagnosticsInfo"></span></li>
+          <li class="nav-header">${ _('Start time') }</li>
+          <li><span data-bind="moment: {data: startTime, format: 'LLL'}"></span></li>
+          <li class="nav-header">${ _('Node Http Address') }</li>
+          <li><span data-bind="text: nodeHttpAddress"></span></li>
+          <li class="nav-header">${ _('Elapsed time') }</li>
+          <li><span data-bind="text: duration().toHHMMSS()"></span></li>
           <!-- /ko -->
           <!-- /ko -->
         </ul>
       </div>
     </div>
-
     <div data-bind="css: {'span10': !$root.isMini(), 'span12': $root.isMini() }">
     </div>
   </div>
@@ -2453,7 +2450,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'].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', 'Oozie Launcher'].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

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

@@ -357,8 +357,8 @@ 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':
-      log_link = app['amContainerLogs']
+    elif app['applicationType'] == 'Oozie Launcher' or app['applicationType'] == 'TEZ':
+      log_link = app.get('amContainerLogs')
   except (KeyError, RestException), e:
     raise KeyError(_("Cannot find job attempt '%(id)s'.") % {'id': job.jobId}, e)
   except Exception, e:

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

@@ -543,6 +543,8 @@ class Attempt:
     logs = []
     attempt = self.task.job.job_attempts['jobAttempt'][-1]
     log_link = attempt['logsLink']
+    if not log_link:
+      return ['', '', '']
 
     # Generate actual task log link from logsLink url
     if self.task.job.status in ('NEW', 'SUBMITTED', 'RUNNING'):
@@ -641,6 +643,22 @@ class YarnOozieAttempt(Attempt):
     if not hasattr(self, 'diagnostics'):
       self.diagnostics = ''
     setattr(self, 'type', 'Oozie Launcher')
+    if self.finishedTime == 0:
+      finishTime = int(time.time() * 1000)
+    else:
+      finishTime = self.finishedTime
+    if self.startTime == 0:
+      durationInMillis = None
+    else:
+      durationInMillis = finishTime - self.startTime
+
+    setattr(self, 'duration', durationInMillis)
+    setattr(self, 'durationInMillis', durationInMillis)
+    setattr(self, 'durationFormatted', self.duration and format_duration_in_millis(self.duration))
+    setattr(self, 'finishTimeFormatted', format_unixtime_ms(finishTime))
+    setattr(self, 'startTimeFormatted', format_unixtime_ms(self.startTime))
+    setattr(self, 'status', 'RUNNING' if self.finishedTime == 0 else 'SUCCEEDED')
+    setattr(self, 'properties', {})
 
 class Container:
 

+ 1 - 1
desktop/libs/hadoop/src/hadoop/yarn/resource_manager_api.py

@@ -131,7 +131,7 @@ class ResourceManagerApi(object):
   def appattempts_attempt(self, app_id, attempt_id):
     attempts = self.appattempts(app_id)
     for attempt in attempts['appAttempts']['appAttempt']:
-      if attempt['id'] == attempt_id:
+      if attempt['id'] == attempt_id or attempt.get('appAttemptId'):
         return attempt
     raise PopupException('Application {} does not have application attempt with id {}'.format(app_id, attempt_id))