浏览代码

HUE-1604 [oozie] Jar path in mapreduce action and java action should use File arg

Abraham Elmahrek 12 年之前
父节点
当前提交
e1f292f

+ 34 - 5
apps/oozie/src/oozie/import_workflow.py

@@ -32,10 +32,7 @@ Action extensions are also versioned.
 Every action extension will have its own version via /xslt/<workflow version>/extensions/<name of extensions>.<version>.xslt
 """
 
-try:
-  import json
-except ImportError:
-  import simplejson as json
+import json
 
 import logging
 from lxml import etree
@@ -501,7 +498,6 @@ def _resolve_subworkflow_from_deployment_dir(fs, workflow, app_path):
   """
   Resolves subworkflow in a subworkflow node
   Looks at path and interrogates all workflows until the proper deployment path is found.
-  If the proper deployment path is never found, then
   """
   if not fs:
     raise RuntimeError(_("No hadoop file system to operate on."))
@@ -547,7 +543,39 @@ def _save_nodes(workflow, nodes):
       node.save()
 
 
+def _resolve_jar_paths(workflow):
+  """
+  Make first file in "files" field the "jar path".
+  """
+  for node in workflow.node_list:
+    if hasattr(node, 'jar_path') and hasattr(node, 'files'):
+      files = json.loads(node.files)
+      if files:
+        node.jar_path = files.pop(0)
+        node.files = json.dumps(files)
+        node.save()
+
+
+def _postprocess_workflow(workflow):
+  """
+  Post processing step.
+  """
+  _resolve_jar_paths(workflow)
+
+
 def import_workflow(workflow, workflow_definition, fs=None):
+  """
+  Import workflow takes 7 steps:
+  1. Perform XSLT.
+  2. Verify schema version.
+  3. Prepare nodes for importing.
+  4. Preprocess nodes before importing.
+  5. Save nodes after they've been processed.
+  6. Save links after the nodes have been added.
+  7. Post process the workflow.
+
+  Most logic is in steps 3, 4, and 6, 7.
+  """
   xslt_definition_fh = open("%(xslt_dir)s/workflow.xslt" % {
     'xslt_dir': DEFINITION_XSLT_DIR.get()
   })
@@ -581,6 +609,7 @@ def import_workflow(workflow, workflow_definition, fs=None):
   _preprocess_nodes(workflow, transformed_root, workflow_definition_root, nodes, fs)
   _save_nodes(workflow, nodes)
   _save_links(workflow, workflow_definition_root)
+  _postprocess_workflow(workflow)
 
   # Update schema_version
   workflow.schema_version = schema_version

+ 8 - 2
apps/oozie/src/oozie/models.py

@@ -707,7 +707,10 @@ class Mapreduce(Action):
     return json.loads(self.job_properties)
 
   def get_files(self):
-    return json.loads(self.files)
+    files = json.loads(self.files)
+    if self.jar_path:
+      files.insert(0, self.jar_path)
+    return files
 
   def get_archives(self):
     return json.loads(self.archives)
@@ -780,7 +783,10 @@ class Java(Action):
     return json.loads(self.job_properties)
 
   def get_files(self):
-    return json.loads(self.files)
+    files = json.loads(self.files)
+    if self.jar_path:
+      files.insert(0, self.jar_path)
+    return files
 
   def get_archives(self):
     return json.loads(self.archives)

+ 1 - 20
desktop/libs/liboozie/src/liboozie/submittion.py

@@ -149,7 +149,7 @@ class Submission(object):
       raise PopupException(message=msg, detail=str(ex))
 
     oozie_xml = self.job.to_xml(self.properties)
-    self._do_as(self.user.username , self._copy_files, deployment_dir, oozie_xml)
+    self._do_as(self.user.username, self._copy_files, deployment_dir, oozie_xml)
 
     if hasattr(self.job, 'actions'):
       for action in self.job.actions:
@@ -222,25 +222,6 @@ class Submission(object):
     self.fs.create(xml_path, overwrite=True, permission=0644, data=oozie_xml)
     LOG.debug("Created %s" % (xml_path,))
 
-    # Copy the files over
-    files = []
-    if hasattr(self.job, 'node_list'):
-      for node in self.job.node_list:
-        if hasattr(node, 'jar_path') and node.jar_path.startswith('/'):
-          files.append(node.jar_path)
-
-    if files:
-      lib_path = self.fs.join(deployment_dir, 'lib')
-      if self.fs.exists(lib_path):
-        LOG.debug("Cleaning up old %s" % (lib_path,))
-        self.fs.rmtree(lib_path)
-
-      self.fs.mkdir(lib_path, 0755)
-      LOG.debug("Created %s" % (lib_path,))
-
-      for file in files:
-        self.fs.copyfile(file, self.fs.join(lib_path, self.fs.basename(file)))
-
   def _do_as(self, username, fn, *args, **kwargs):
     prev_user = self.fs.user
     try: