Explorar o código

HUE-1176 [jb2] Add check app status and progress to the API

Romain Rigaux %!s(int64=8) %!d(string=hai) anos
pai
achega
73307307b5

+ 1 - 3
apps/jobbrowser/src/jobbrowser/apis/base_api.py

@@ -51,12 +51,10 @@ class Api(object):
 
   def apps(self): return []
 
-  def app(self, appid): return {}
+  def app(self, appid): return {} # Also contains progress (0-100) and status [RUNNING, FINISHED, PAUSED]
 
   def action(self, appid, operation): return {}
 
-  def status(self, appid): return {'status': 'RUNNING'}
-
   def logs(self, appid, app_type): return {'progress': 0, 'logs': {'default': ''}}
 
   def profile(self, appid, app_type, app_property): return {} # Tasks, XML, counters...

+ 11 - 3
apps/jobbrowser/src/jobbrowser/apis/job_api.py

@@ -81,8 +81,9 @@ class YarnApi(Api):
         'name': app.name,
         'type': app.applicationType,
         'status': app.status,
+        'apiStatus': self._api_status(app.status),
         'user': self.user.username,
-        'progress': 100,
+        'progress': app.progress,
         'duration': 10 * 3600,
         'submitted': 10 * 3600
     } for app in jobs]
@@ -96,8 +97,9 @@ class YarnApi(Api):
         'name': app.name,
         'type': app.applicationType,
         'status': app.status,
+        'apiStatus': self._api_status(app.status),
         'user': self.user.username,
-        'progress': 100,
+        'progress': app.progress,
         'duration': 10 * 3600,
         'submitted': 10 * 3600
     }
@@ -135,7 +137,7 @@ class YarnApi(Api):
       logs = json.loads(response.content)['log']
     else:
       logs = None
-    return {'progress': 0, 'logs': {'default': logs}}
+    return {'logs': {'default': logs}}
 
 
   def profile(self, appid, app_type, app_property):
@@ -151,6 +153,12 @@ class YarnApi(Api):
 
     return {}
 
+  def _api_status(self, status):
+    if status in ['NEW', 'NEW_SAVING', 'SUBMITTED', 'ACCEPTED', 'RUNNING']:
+      return 'RUNNING'
+    else:
+      return 'FINISHED' # FINISHED, FAILED, KILLED
+
 
 class YarnMapReduceTaskApi(Api):
 

+ 12 - 2
apps/jobbrowser/src/jobbrowser/apis/schedule_api.py

@@ -48,9 +48,10 @@ class ScheduleApi(Api):
         'id': app.id,
         'name': app.appName,
         'status': app.status,
+        'apiStatus': self._api_status(app.status),
         'type': 'schedule',
         'user': app.user,
-        'progress': 100,
+        'progress': app.get_progress(),
         'duration': 10 * 3600,
         'submitted': 10 * 3600
     } for app in wf_list.jobs]
@@ -67,6 +68,8 @@ class ScheduleApi(Api):
         'id': coordinator.coordJobId,
         'name': coordinator.coordJobName,
         'status': coordinator.status,
+        'apiStatus': self._api_status(coordinator.status),
+        'progress': coordinator.get_progress(),
         'type': 'schedule',
         'startTime': format_time(coordinator.startTime),
     }
@@ -80,7 +83,7 @@ class ScheduleApi(Api):
     request = MockDjangoRequest(self.user)
     data = get_oozie_job_log(request, job_id=appid)
 
-    return {'progress': 0, 'logs': {'default': json.loads(data.content)['log']}}
+    return {'logs': {'default': json.loads(data.content)['log']}}
 
 
   def profile(self, appid, app_type, app_property):
@@ -97,6 +100,13 @@ class ScheduleApi(Api):
         'properties': workflow.conf_dict,
       }
 
+  def _api_status(self, status):
+    if status in ['PREP', 'RUNNING', 'RUNNINGWITHERROR']:
+      return 'RUNNING'
+    elif status in ['PREPSUSPENDED', 'SUSPENDED', 'SUSPENDEDWITHERROR', 'PREPPAUSED', 'PAUSED', 'PAUSEDWITHERROR']:
+      return 'PAUSED'
+    else:
+      return 'FINISHED' # SUCCEEDED, DONEWITHERROR, KILLED, FAILED
 
 
 class MockGet():

+ 19 - 3
apps/jobbrowser/src/jobbrowser/apis/workflow_api.py

@@ -45,9 +45,10 @@ class WorkflowApi(Api):
         'id': app.id,
         'name': app.appName,
         'status': app.status,
+        'apiStatus': self._api_status(app.status),
         'type': 'workflow',
         'user': app.user,
-        'progress': 100,
+        'progress': app.get_progress(),
         'duration': 10 * 3600,
         'submitted': 10 * 3600
     } for app in wf_list.jobs]
@@ -60,7 +61,14 @@ class WorkflowApi(Api):
     oozie_api = get_oozie(self.user)
     workflow = oozie_api.get_job(jobid=appid)
 
-    common = {'id': workflow.id, 'name': workflow.appName, 'status': workflow.status, 'type': 'workflow'}
+    common = {
+        'id': workflow.id,
+        'name': workflow.appName,
+        'status': workflow.status,
+        'apiStatus': self._api_status(workflow.status),
+        'progress': workflow.get_progress(),
+        'type': 'workflow'
+    }
 
     request = MockDjangoRequest(self.user)
     response = list_oozie_workflow(request, job_id=appid)
@@ -89,7 +97,7 @@ class WorkflowApi(Api):
     request = MockDjangoRequest(self.user)
     data = get_oozie_job_log(request, job_id=appid)
 
-    return {'progress': 0, 'logs': {'default': json.loads(data.content)['log']}}
+    return {'logs': {'default': json.loads(data.content)['log']}}
 
 
   def profile(self, appid, app_type, app_property):
@@ -111,6 +119,14 @@ class WorkflowApi(Api):
 
     return {}
 
+  def _api_status(self, status):
+    if status in ['PREP', 'RUNNING']:
+      return 'RUNNING'
+    elif status == 'SUSPENDED':
+      return 'PAUSED'
+    else:
+      return 'FINISHED' # SUCCEEDED , KILLED and FAILED
+
 
 class WorkflowActionApi(Api):
 

+ 37 - 2
apps/jobbrowser/src/jobbrowser/templates/job_browser.mako

@@ -815,10 +815,14 @@ ${ commonheader("Job Browser", "jobbrowser", user, request) | n,unicode }
       self.id = ko.observableDefault(job.id);
       self.name = ko.observableDefault(job.name);
       self.type = ko.observableDefault(job.type);
+
       self.status = ko.observableDefault(job.status);
+      self.apiStatus = ko.observableDefault(job.apiStatus);
+      self.progress = ko.observableDefault(job.progress);
+      self.checkStatusTimeout = null;
+
       self.user = ko.observableDefault(job.user);
       self.cluster = ko.observableDefault(job.cluster);
-      self.progress = ko.observableDefault(job.progress);
       self.duration = ko.observableDefault(job.duration);
       self.submitted = ko.observableDefault(job.submitted);
 
@@ -831,9 +835,15 @@ ${ commonheader("Job Browser", "jobbrowser", user, request) | n,unicode }
 
       self.loadingJob = ko.observable(false);
 
+
       self.fetchJob = function () {
         self.loadingJob(true);
 
+        if (self.checkStatusTimeout != null) {
+          clearTimeout(self.checkStatusTimeout);
+          self.checkStatusTimeout = null;
+        }
+
         var interface = vm.interface();
         if (/oozie-oozi-W/.test(self.id())) { interface = 'workflows'; };
 
@@ -849,6 +859,8 @@ ${ commonheader("Job Browser", "jobbrowser", user, request) | n,unicode }
             vm.breadcrumbs.push({'id': vm.job().id(), 'name': vm.job().name(), 'type': vm.job().type()});
 
             vm.job().fetchLogs();
+            vm.job().fetchStatus();
+
             if (self.mainType() == 'schedules') {
               //vm.job().coordVM.setActions(data.app.actions);
             }
@@ -892,6 +904,29 @@ ${ commonheader("Job Browser", "jobbrowser", user, request) | n,unicode }
         });
       };
 
+      self.fetchStatus = function () {
+        if (self.apiStatus() != 'RUNNING') {
+          return;
+        }
+
+        $.post("/jobbrowser/api/job", {
+          app_id: ko.mapping.toJSON(self.id),
+          interface: ko.mapping.toJSON(self.mainType)
+        }, function (data) {
+          if (data.status == 0) {
+            self.status(data.app.status);
+            self.apiStatus(data.app.apiStatus);
+            self.progress(data.app.progress);
+
+            if (self.apiStatus() == 'RUNNING') {
+              self.checkStatusTimeout = setTimeout(self.fetchStatus, 2000);
+            }
+          } else {
+            $(document).trigger("error", data.message);
+          }
+        });
+      };
+
       self.control = function (action) {
         $.post("/jobbrowser/api/job/action", {
           app_id: ko.mapping.toJSON(self.id),
@@ -927,7 +962,7 @@ ${ commonheader("Job Browser", "jobbrowser", user, request) | n,unicode }
           if (data.status == 0) {
             var apps = [];
             if (data && data.apps) {
-              data.apps.forEach(function (job) { // TODO: update and merge
+              data.apps.forEach(function (job) { // TODO: update and merge with status and progress
                 apps.push(new Job(vm, job));
               });
             }

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

@@ -191,6 +191,8 @@ class Job(object):
 
     self._fixup()
 
+    self.progress = None
+
     # Set MAPS/REDUCES completion percentage
     if hasattr(self, 'mapsTotal'):
       self.desiredMaps = self.mapsTotal
@@ -198,6 +200,7 @@ class Job(object):
         self.maps_percent_complete = 0
       else:
         self.maps_percent_complete = int(round(float(self.finishedMaps) / self.desiredMaps * 100))
+      self.progress = self.maps_percent_complete
 
     if hasattr(self, 'reducesTotal'):
       self.desiredReduces = self.reducesTotal
@@ -205,6 +208,10 @@ class Job(object):
         self.reduces_percent_complete = 0
       else:
         self.reduces_percent_complete = int(round(float(self.finishedReduces) / self.desiredReduces * 100))
+      if self.progress is not None:
+        self.progress = int((self.progress + self.reduces_percent_complete) / 2)
+      else:
+        self.progress = self.reduces_percent_complete
 
 
   def _fixup(self):
@@ -287,6 +294,7 @@ class KilledJob(Job):
       setattr(self, 'finishTime', self.finishedTime)
     if not hasattr(self, 'startTime'):
       setattr(self, 'startTime', self.startedTime)
+    self.progress = 100
 
     super(KilledJob, self)._fixup()