Pārlūkot izejas kodu

HUE-2291 [oozie] Remove check status from workflow dashboard display

Romain Rigaux 11 gadi atpakaļ
vecāks
revīzija
3b4c5d2f5a

+ 20 - 16
apps/oozie/src/oozie/views/dashboard.py

@@ -37,6 +37,7 @@ from desktop.log.access import access_warn
 
 from liboozie.oozie_api import get_oozie
 from liboozie.submittion import Submission
+from liboozie.types import Workflow as OozieWorkflow
 
 from oozie.conf import OOZIE_JOBS_COUNT, ENABLE_CRON_SCHEDULING
 from oozie.forms import RerunForm, ParameterForm, RerunCoordForm,\
@@ -100,20 +101,23 @@ def list_oozie_workflows(request):
   if not has_dashboard_jobs_access(request.user):
     kwargs['user'] = request.user.username
 
-  workflows = get_oozie(request.user).get_workflows(**kwargs)
-
-  if request.GET.get('format') == 'json':
-    json_jobs = workflows.jobs
+  if request.GET.get('format') == 'json':    
     just_sla = request.GET.get('justsla') == 'true'
     if request.GET.get('type') == 'running':
-      json_jobs = split_oozie_jobs(request.user, workflows.jobs)['running_jobs']
-    if request.GET.get('type') == 'completed':
-      json_jobs = split_oozie_jobs(request.user, workflows.jobs)['completed_jobs']
+      kwargs['filters'] = [('status', status) for status in OozieWorkflow.RUNNING_STATUSES]
+      json_jobs = get_oozie(request.user).get_workflows(**kwargs).jobs
+    elif request.GET.get('type') == 'completed':
+      kwargs['filters'] = [('status', status) for status in OozieWorkflow.FINISHED_STATUSES]
+      json_jobs = get_oozie(request.user).get_workflows(**kwargs).jobs
+    elif request.GET.get('type') == 'progress':
+      kwargs['filters'] = [('status', status) for status in OozieWorkflow.RUNNING_STATUSES]
+      json_jobs = get_oozie(request.user).get_workflows(**kwargs).jobs
+      json_jobs = [get_oozie(request.user).get_job(job.id) for job in json_jobs] 
     return HttpResponse(encode_json_for_js(massaged_oozie_jobs_for_json(json_jobs, request.user, just_sla)), mimetype="application/json")
 
   return render('dashboard/list_oozie_workflows.mako', request, {
     'user': request.user,
-    'jobs': split_oozie_jobs(request.user, workflows.jobs),
+    'jobs': [],
     'has_job_edition_permission':  has_job_edition_permission,
   })
 
@@ -750,14 +754,14 @@ def massaged_oozie_jobs_for_json(oozie_jobs, user, just_sla=False):
   jobs = []
 
   for job in oozie_jobs:
-    if job.is_running():
-      if job.type == 'Workflow':
-        job = get_oozie(user).get_job(job.id)
-      elif job.type == 'Coordinator':
-        job = get_oozie(user).get_coordinator(job.id)
-      else:
-        job = get_oozie(user).get_bundle(job.id)
-    if not just_sla or (just_sla and job.has_sla):
+#    if job.is_running():
+#      if job.type == 'Workflow':
+#        job = get_oozie(user).get_job(job.id)
+#      elif job.type == 'Coordinator':
+#        job = get_oozie(user).get_coordinator(job.id)
+#      else:
+#        job = get_oozie(user).get_bundle(job.id)
+    if not just_sla or (just_sla and job.has_sla) and job.appName != 'pig-app-hue-script':
       massaged_job = {
         'id': job.id,
         'lastModTime': hasattr(job, 'lastModTime') and job.lastModTime and format_time(job.lastModTime) or None,

+ 14 - 13
desktop/libs/liboozie/src/liboozie/oozie_api.py

@@ -105,11 +105,10 @@ class OozieApi(object):
 
   VALID_JOB_FILTERS = ('name', 'user', 'group', 'status')
 
-  def get_jobs(self, jobtype, offset=None, cnt=None, **kwargs):
+  def get_jobs(self, jobtype, offset=None, cnt=None, filters=None):
     """
     Get a list of Oozie jobs.
 
-    jobtype is 'wf', 'coord'
     Note that offset is 1-based.
     kwargs is used for filtering and may be one of VALID_FILTERS: name, user, group, status
     """
@@ -118,10 +117,12 @@ class OozieApi(object):
       params['offset'] = str(offset)
     if cnt is not None:
       params['len'] = str(cnt)
+    if filters is None:
+      filters = []
     params['jobtype'] = jobtype
 
-    filter_list = [ ]
-    for key, val in kwargs.iteritems():
+    filter_list = []
+    for key, val in filters:
       if key not in OozieApi.VALID_JOB_FILTERS:
         raise ValueError('"%s" is not a valid filter for selecting jobs' % (key,))
       filter_list.append('%s=%s' % (key, val))
@@ -130,21 +131,21 @@ class OozieApi(object):
     # Send the request
     resp = self._root.get('jobs', params)
     if jobtype == 'wf':
-      wf_list = WorkflowList(self, resp, filters=kwargs)
+      wf_list = WorkflowList(self, resp, filters=filters)
     elif jobtype == 'coord':
-      wf_list = CoordinatorList(self, resp, filters=kwargs)
+      wf_list = CoordinatorList(self, resp, filters=filters)
     else:
-      wf_list = BundleList(self, resp, filters=kwargs)
+      wf_list = BundleList(self, resp, filters=filters)
     return wf_list
 
-  def get_workflows(self, offset=None, cnt=None, **kwargs):
-    return self.get_jobs('wf', offset, cnt, **kwargs)
+  def get_workflows(self, offset=None, cnt=None, filters=None):
+    return self.get_jobs('wf', offset, cnt, filters)
 
-  def get_coordinators(self, offset=None, cnt=None, **kwargs):
-    return self.get_jobs('coord', offset, cnt, **kwargs)
+  def get_coordinators(self, offset=None, cnt=None, filters=None):
+    return self.get_jobs('coord', offset, cnt, filters)
 
-  def get_bundles(self, offset=None, cnt=None, **kwargs):
-    return self.get_jobs('bundle', offset, cnt, **kwargs)
+  def get_bundles(self, offset=None, cnt=None, filters=None):
+    return self.get_jobs('bundle', offset, cnt, filters)
 
   # TODO: make get_job accept any jobid
   def get_job(self, jobid):

+ 10 - 9
desktop/libs/liboozie/src/liboozie/types.py

@@ -275,18 +275,13 @@ class BundleAction(Action):
     return progress
 
 
-class Job(object):
-  RUNNING_STATUSES = set([
-     'PREP', 'RUNNING', 'SUSPENDED', # Workflow
-     'RUNNING', 'PREPSUSPENDED', 'SUSPENDED', 'PREPPAUSED', 'PAUSED' # Coordinator
-    ]
-  )
+class Job(object):  
   MAX_LOG_SIZE = 3500 * 20 # 20 pages
 
   """
   Accessing log and definition will trigger Oozie API calls.
   """
-  def __init__(self, api, json_dict):
+  def __init__(self, api, json_dict):    
     for attr in self._ATTRS:
       setattr(self, attr, json_dict.get(attr))
     self._fixup()
@@ -374,7 +369,7 @@ class Job(object):
     return [action for action in self.actions if not ControlFlowAction.is_control_flow(action.type)]
 
   def is_running(self):
-    return self.status in Job.RUNNING_STATUSES
+    return self.status in (Workflow.RUNNING_STATUSES, Coordinator.RUNNING_STATUSES, Bundle.RUNNING_STATUSES)
 
   def __str__(self):
     return '%s - %s' % (self.id, self.status)
@@ -405,6 +400,8 @@ class Workflow(Job):
     'parentId'
   ]
   ACTION = WorkflowAction
+  RUNNING_STATUSES = set(['PREP', 'RUNNING', 'SUSPENDED'])
+  FINISHED_STATUSES = set(['SUCCEEDED' , 'KILLED', 'FAILED'])
 
   def _fixup(self):
     super(Workflow, self)._fixup()
@@ -476,6 +473,8 @@ class Coordinator(Job):
     'bundleId'
   ]
   ACTION = CoordinatorAction
+  RUNNING_STATUSES = set(['PREP', 'RUNNING', 'RUNNINGWITHERROR', 'PREPSUSPENDED', 'SUSPENDED', 'SUSPENDEDWITHERROR', 'PREPPAUSED', 'PAUSED', 'PAUSEDWITHERROR'])
+  FINISHED_STATUSES = set(['SUCCEEDED', 'DONEWITHERROR', 'KILLED', 'FAILED'])
 
   def _fixup(self):
     super(Coordinator, self)._fixup()
@@ -571,6 +570,8 @@ class Bundle(Job):
   ]
 
   ACTION = BundleAction
+  RUNNING_STATUSES = set(['PREP', 'RUNNING', 'RUNNINGWITHERROR', 'SUSPENDED', 'PREPSUSPENDED', 'SUSPENDEDWITHERROR', 'PAUSED', 'PAUSEDWITHERROR', 'PREPPAUSED'])
+  FINISHED_STATUSES = set(['SUCCEEDED', 'DONEWITHERROR', 'KILLED', 'FAILED'])
 
   def _fixup(self):
     self.actions = self.bundleCoordJobs
@@ -616,7 +617,7 @@ class JobList(object):
   def __init__(self, klass, jobs_key, api, json_dict, filters=None):
     """
     json_dict is the oozie json.
-    filters is (optionally) the dictionary of filters used to select this list
+    filters is (optionally) the list of filters used to select this list
     """
     self._api = api
     self.offset = int(json_dict['offset'])