فهرست منبع

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

This reverts commit e1f292f98fbfe269a78fb4dbd00225f183250312.
Abraham Elmahrek 12 سال پیش
والد
کامیت
52e25f5
3فایلهای تغییر یافته به همراه27 افزوده شده و 43 حذف شده
  1. 5 34
      apps/oozie/src/oozie/import_workflow.py
  2. 2 8
      apps/oozie/src/oozie/models.py
  3. 20 1
      desktop/libs/liboozie/src/liboozie/submittion.py

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

@@ -32,7 +32,10 @@ Action extensions are also versioned.
 Every action extension will have its own version via /xslt/<workflow version>/extensions/<name of extensions>.<version>.xslt
 """
 
-import json
+try:
+  import json
+except ImportError:
+  import simplejson as json
 
 import logging
 from lxml import etree
@@ -499,6 +502,7 @@ 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."))
@@ -544,39 +548,7 @@ 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()
   })
@@ -610,7 +582,6 @@ 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

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

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

+ 20 - 1
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,6 +222,25 @@ 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: