Sfoglia il codice sorgente

[core] Latest Oozie 3.2 breaks the app

Updating the liboozie
Romain Rigaux 13 anni fa
parent
commit
d7a6d69ac5

+ 2 - 2
apps/oozie/src/oozie/templates/dashboard/list_oozie_coordinator.mako

@@ -95,7 +95,7 @@ ${ layout.menubar(section='dashboard') }
             </tr>
           </thead>
           <tbody>
-            % for i, action in enumerate(reversed(oozie_coordinator.actions)):
+            % for i, action in enumerate(reversed(oozie_coordinator.get_working_actions())):
               <tr>
                 <td>
                   % if action.externalId:
@@ -133,7 +133,7 @@ ${ layout.menubar(section='dashboard') }
             </tr>
           </thead>
           <tbody>
-            % for i, action in enumerate(oozie_coordinator.actions):
+            % for i, action in enumerate(oozie_coordinator.get_working_actions()):
               <tr>
 
                 <td>${ action.actionNumber }</td>

+ 2 - 2
apps/oozie/src/oozie/templates/dashboard/list_oozie_workflow.mako

@@ -135,7 +135,7 @@ ${ layout.menubar(section='dashboard') }
            forms = WorkflowFormSet(instance=hue_workflow.get_full_node()).forms
          %>
 
-           ${ hue_workflow.get_full_node().gen_status_graph(forms, oozie_workflow.actions) }
+           ${ hue_workflow.get_full_node().gen_status_graph(forms, oozie_workflow.get_working_actions()) }
          % endif
        </div>
      % endif
@@ -162,7 +162,7 @@ ${ layout.menubar(section='dashboard') }
             </tr>
           </thead>
           <tbody>
-            % for i, action in enumerate(oozie_workflow.actions):
+            % for i, action in enumerate(oozie_workflow.get_working_actions()):
               <tr>
                 <td>
                   <a href="${ url('oozie:list_oozie_workflow_action', action=action.id) }" data-row-selector='true'>${ action.id }</a>

+ 14 - 0
apps/oozie/src/oozie/views/dashboard.py

@@ -46,6 +46,15 @@ Permissions checking happens by calling check_access_and_get_oozie_job().
 """
 
 
+def show_oozie_error(view_func):
+  def decorate(request, *args, **kwargs):
+    try:
+      return view_func(request, *args, **kwargs)
+    except RestException, ex:
+      raise PopupException(_('Sorry, an error with Oozie happened.'), detail=ex._headers.get('oozie-error-message', ex))
+  return wraps(view_func)(decorate)
+
+
 def manage_oozie_jobs(request, job_id, action):
   if request.method != 'POST':
     raise PopupException(_('Please use a POST request to manage an Oozie job.'))
@@ -64,6 +73,7 @@ def manage_oozie_jobs(request, job_id, action):
   return HttpResponse(json.dumps(response), mimetype="application/json")
 
 
+@show_oozie_error
 def list_oozie_workflows(request):
   kwargs = {'cnt': 50,}
   if not request.user.is_superuser:
@@ -77,6 +87,7 @@ def list_oozie_workflows(request):
   })
 
 
+@show_oozie_error
 def list_oozie_coordinators(request):
   kwargs = {'cnt': 50,}
   if not request.user.is_superuser:
@@ -89,6 +100,7 @@ def list_oozie_coordinators(request):
   })
 
 
+@show_oozie_error
 def list_oozie_workflow(request, job_id, coordinator_job_id=None):
   oozie_workflow = check_access_and_get_oozie_job(request, job_id)
 
@@ -124,6 +136,7 @@ def list_oozie_workflow(request, job_id, coordinator_job_id=None):
   })
 
 
+@show_oozie_error
 def list_oozie_coordinator(request, job_id):
   oozie_coordinator = check_access_and_get_oozie_job(request, job_id)
 
@@ -140,6 +153,7 @@ def list_oozie_coordinator(request, job_id):
   })
 
 
+@show_oozie_error
 def list_oozie_workflow_action(request, action):
   try:
     action = get_oozie().get_action(action)

+ 2 - 1
apps/oozie/src/oozie/views/editor.py

@@ -36,12 +36,13 @@ from desktop.log.access import access_warn
 from hadoop.fs.exceptions import WebHdfsException
 from liboozie.submittion import Submission
 
+from oozie.conf import SHARE_JOBS
 from oozie.models import Workflow, Node, Link, History, Coordinator,\
   Dataset, DataInput, DataOutput, Job, _STD_PROPERTIES_JSON
 from oozie.forms import NodeForm, WorkflowForm, CoordinatorForm, DatasetForm,\
   DataInputForm, DataInputSetForm, DataOutputForm, DataOutputSetForm, LinkForm,\
   DefaultLinkForm, design_form_by_type
-from oozie.conf import SHARE_JOBS
+
 
 LOG = logging.getLogger(__name__)
 

+ 73 - 7
desktop/libs/liboozie/src/liboozie/types.py

@@ -16,9 +16,10 @@
 # limitations under the License.
 
 """
-Oozie objects.
+Oozie API classes.
 
-This is mostly just codifying the oozie json.
+This is mostly just codifying the datastructure of the Oozie REST API.
+http://incubator.apache.org/oozie/docs/3.2.0-incubating/docs/WebServicesAPI.html
 """
 
 from cStringIO import StringIO
@@ -41,9 +42,63 @@ class Action(object):
       setattr(self, attr, json_dict.get(attr))
     self._fixup()
 
+  def _fixup(self): pass
+
   def is_finished(self):
     return self.status == 'OK'
 
+  @classmethod
+  def create(self, action_class, action_dict):
+    if ControlFlowAction.is_control_flow(action_dict['type']):
+      return ControlFlowAction(action_dict)
+    else:
+      return action_class(action_dict)
+
+
+class ControlFlowAction(Action):
+  _ATTRS = [
+    'errorMessage',
+    'status',
+    'stats',
+    'data',
+    'transition',
+    'externalStatus',
+    'cred',
+    'conf',
+    'type',
+    'endTime',
+    'externalId',
+    'id',
+    'startTime',
+    'externalChildIDs',
+    'name',
+    'errorCode',
+    'trackerUri',
+    'retries',
+    'toString',
+    'consoleUrl'
+  ]
+
+  @classmethod
+  def is_control_flow(self, action_type):
+    return action_type is not None and ':' in action_type
+
+  def _fixup(self):
+    """
+    Fixup:
+      - time fields as struct_time
+      - config dict
+    """
+    super(ControlFlowAction, self)._fixup()
+
+    if self.startTime:
+      self.startTime = parse_timestamp(self.startTime)
+    if self.endTime:
+      self.endTime = parse_timestamp(self.endTime)
+    if self.retries:
+      self.retries = int(self.retries)
+
+    self.conf_dict = {}
 
 class CoordinatorAction(Action):
   _ATTRS = [
@@ -72,6 +127,8 @@ class CoordinatorAction(Action):
       - time fields as struct_time
       - config dict
     """
+    super(CoordinatorAction, self)._fixup()
+
     if self.createdTime:
       self.createdTime = parse_timestamp(self.createdTime)
     if self.nominalTime:
@@ -79,9 +136,11 @@ class CoordinatorAction(Action):
     if self.lastModifiedTime:
       self.lastModifiedTime = parse_timestamp(self.lastModifiedTime)
 
-    xml = StringIO(i18n.smart_str(self.runConf))
-    self.conf_dict = hadoop.confparse.ConfParse(xml)
-
+    if self.runConf:
+      xml = StringIO(i18n.smart_str(self.runConf))
+      self.conf_dict = hadoop.confparse.ConfParse(xml)
+    else:
+      self.conf_dict = {}
 
 class WorkflowAction(Action):
   _ATTRS = [
@@ -103,13 +162,14 @@ class WorkflowAction(Action):
     'type',
   ]
 
-
   def _fixup(self):
     """
     Fixup:
       - time fields as struct_time
       - config dict
     """
+    super(WorkflowAction, self)._fixup()
+
     if self.startTime:
       self.startTime = parse_timestamp(self.startTime)
     if self.endTime:
@@ -153,7 +213,7 @@ class Job(object):
     if self.endTime:
       self.endTime = parse_timestamp(self.endTime)
 
-    self.actions = [ self.ACTION(act_dict) for act_dict in self.actions ]
+    self.actions = [Action.create(self.ACTION, act_dict) for act_dict in self.actions]
     if self.conf is not None:
       xml = StringIO(i18n.smart_str(self.conf))
       self.conf_dict = hadoop.confparse.ConfParse(xml)
@@ -213,6 +273,12 @@ class Job(object):
       raise PopupException(_("Permission denied. User %(username)s cannot modify user %(user)s's job.") %
                            dict(username=request.user.username, user=self.user))
 
+  def get_control_flow_actions(self):
+    return [action for action in self.actions if ControlFlowAction.is_control_flow(action.type)]
+
+  def get_working_actions(self):
+    return [action for action in self.actions if not ControlFlowAction.is_control_flow(action.type)]
+
 class Coordinator(Job):
   _ATTRS = [
     'acl',