فهرست منبع

HUE-1176 [jb] Adding bulk actions to jobs and workflows

Romain Rigaux 8 سال پیش
والد
کامیت
76c8d3a

+ 3 - 3
apps/jobbrowser/src/jobbrowser/api2.py

@@ -76,14 +76,14 @@ def job(request):
 
 @api_error_handler
 def action(request):
-  response = {'status': -1}
+  response = {'status': -1, 'message': ''}
 
   interface = json.loads(request.POST.get('interface'))
-  app_id = json.loads(request.POST.get('app_id'))
+  app_ids = json.loads(request.POST.get('app_ids'))
   operation = json.loads(request.POST.get('operation'))
 
   response['operation'] = operation
-  response.update(get_api(request.user, interface).action(app_id, operation))
+  response.update(get_api(request.user, interface).action(app_ids, operation))
 
   return JsonResponse(response)
 

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

@@ -57,7 +57,7 @@ class Api(object):
 
   def app(self, appid): return {} # Also contains progress (0-100) and status [RUNNING, FINISHED, PAUSED]
 
-  def action(self, appid, operation): return {}
+  def action(self, app_ids, operation): return {}
 
   def logs(self, appid, app_type, log_name): return {'progress': 0, 'logs': ''}
 

+ 21 - 12
apps/jobbrowser/src/jobbrowser/apis/job_api.py

@@ -21,6 +21,7 @@ import logging
 from django.utils.translation import ugettext as _
 from hadoop.yarn import resource_manager_api
 
+from desktop.lib.exceptions import MessageException
 from desktop.lib.exceptions_renderable import PopupException
 
 
@@ -53,8 +54,8 @@ class JobApi(Api):
   def app(self, appid):
     return self._get_api(appid).app(appid)
 
-  def action(self, appid, operation):
-    return self._get_api(appid).action(operation, appid)
+  def action(self, app_ids, operation):
+    return self._get_api(app_ids).action(operation, app_ids)
 
   def logs(self, appid, app_type, log_name):
     return self._get_api(appid).logs(appid, app_type, log_name)
@@ -63,7 +64,9 @@ class JobApi(Api):
     return self._get_api(appid).profile(appid, app_type, app_property)
 
   def _get_api(self, appid):
-    if appid.startswith('task_'):
+    if type(appid) == list:
+      return self.yarn_api
+    elif appid.startswith('task_'):
       return YarnMapReduceTaskApi(self.user, appid)
     elif appid.startswith('attempt_'):
       return YarnMapReduceTaskAttemptApi(self.user, appid)
@@ -166,9 +169,15 @@ class YarnApi(Api):
     return common
 
 
-  def action(self, operation, appid):
+  def action(self, operation, app_ids):
     if operation['action'] == 'kill':
-      return kill_job(MockDjangoRequest(self.user), job=appid)
+      kills = []
+      for app_id in app_ids:
+        try:
+          kill_job(MockDjangoRequest(self.user), job=app_id)
+        except MessageException:
+          kills.append(app_id)
+      return {'kills': kills, 'status': len(app_ids) - len(kills), 'message': _('Stop signal sent to %s') % kills}
     else:
       return {}
 
@@ -223,22 +232,22 @@ class YarnMapReduceTaskApi(Api):
       'jobid': self.app_id,
       'pagenum': 1
     }
- 
+
 #     filter_params.update(_extract_query_params(filters)
-#     
+#
 #     #filter_params['text']
-#  
+#
 #     if filters.get('states'):
 #       filter_params['states'] = filters['states']
-#  
+#
 #     if 'time' in filters:
 #       filter_params['time_value'] = int(filters['time']['time_value'])
 #       filter_params['time_unit'] = filters['time']['time_unit']
- 
+
 #     jobs = NativeYarnApi(self.user).get_jobs(**filter_params)
-#  
+#
 #     apps = [massage_job_for_json(job, user=self.user) for job in jobs]
-#  
+#
 #     return {
 #       'apps': [{
 #         'id': app['id'],

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

@@ -109,15 +109,18 @@ class WorkflowApi(Api):
     return common
 
 
-  def action(self, appid, action):
-    if action == 'change' or action == 'ignore' or ',' not in appid:
+  def action(self, app_ids, action):
+    if action == 'change' or action == 'ignore' or len(app_ids) == 1:
       request = MockDjangoRequest(self.user)
-      response = manage_oozie_jobs(request, appid, action['action'])
+      response = manage_oozie_jobs(request, app_ids[0], action['action'])
     else:
-      request = MockDjangoRequest(self.user, post={'job_ids': appid, 'action': action['action']})
+      request = MockDjangoRequest(self.user, post={'job_ids': ' '.join(app_ids), 'action': action['action']})
       response = bulk_manage_oozie_jobs(request)
 
-    return json.loads(response.content)
+    result = json.loads(response.content)
+    result['status'] = result.get('totalErrors', 0)
+    result['message'] = _('%s action sent to %s jobs') % (action['action'], result.get('totalRequests', 1))
+    return result
 
 
   def logs(self, appid, app_type, log_name=None):

+ 22 - 16
apps/jobbrowser/src/jobbrowser/templates/job_browser.mako

@@ -1012,7 +1012,7 @@ ${ commonheader("Job Browser", "jobbrowser", user, request) | n,unicode }
       self.apiStatus = ko.observableDefault(job.apiStatus);
       self.progress = ko.observableDefault(job.progress);
       self.isRunning = ko.computed(function() {
-        return self.apiStatus() != 'SUCCEEDED' && self.apiStatus() != 'FAILED';
+        return self.apiStatus() == 'RUNNING' || self.apiStatus() == 'PAUSED';
       });
 
       self.user = ko.observableDefault(job.user);
@@ -1197,20 +1197,11 @@ ${ commonheader("Job Browser", "jobbrowser", user, request) | n,unicode }
       };
 
       self.control = function (action) {
-        $.post("/jobbrowser/api/job/action", {
-          app_id: ko.mapping.toJSON(self.id),
-          interface: ko.mapping.toJSON(vm.interface),
-          app_type: ko.mapping.toJSON(self.type),
-          operation: ko.mapping.toJSON({action: action})
-        }, function (data) {
-          if (data.status == 0) {
-             $(document).trigger("info", data.message);
-             self.fetchStatus();
-          } else {
-            $(document).trigger("error", data.message);
-          }
+        vm.jobs._control([self.id()], action, function(data) {
+            $(document).trigger("info", data.message);
+            self.fetchStatus();
         });
-      };
+      }
 
       self.updateWorkflowGraph = function() {
         var lastPosition = {
@@ -1445,13 +1436,28 @@ ${ commonheader("Job Browser", "jobbrowser", user, request) | n,unicode }
       };
 
       self.control = function (action) {
+        self._control(
+          $.map(self.selectedJobs(), function(job) {
+            return job.id();
+          }),
+          action,
+          function(data) {
+            $(document).trigger("info", data.message);
+            self.updateJobs();
+          }
+        )
+      }
+
+      self._control = function (app_ids, action, callback) {
         $.post("/jobbrowser/api/job/action", {
-          app_id: ko.mapping.toJSON(self.id), // CSV list
+          app_ids: ko.mapping.toJSON(app_ids),
           interface: ko.mapping.toJSON(vm.interface),
           operation: ko.mapping.toJSON({action: action})
         }, function (data) {
           if (data.status == 0) {
-            $(document).trigger("info", data.message);
+            if (callback) {
+              callback(data);
+            }
           } else {
             $(document).trigger("error", data.message);
           }