Browse Source

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

This reverts commit e1f292f98fbfe269a78fb4dbd00225f183250312.
Abraham Elmahrek 12 years ago
parent
commit
52e25f5cf1

+ 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
 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
 import logging
 from lxml import etree
 from lxml import etree
@@ -499,6 +502,7 @@ def _resolve_subworkflow_from_deployment_dir(fs, workflow, app_path):
   """
   """
   Resolves subworkflow in a subworkflow node
   Resolves subworkflow in a subworkflow node
   Looks at path and interrogates all workflows until the proper deployment path is found.
   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:
   if not fs:
     raise RuntimeError(_("No hadoop file system to operate on."))
     raise RuntimeError(_("No hadoop file system to operate on."))
@@ -544,39 +548,7 @@ def _save_nodes(workflow, nodes):
       node.save()
       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):
 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_definition_fh = open("%(xslt_dir)s/workflow.xslt" % {
     'xslt_dir': DEFINITION_XSLT_DIR.get()
     '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)
   _preprocess_nodes(workflow, transformed_root, workflow_definition_root, nodes, fs)
   _save_nodes(workflow, nodes)
   _save_nodes(workflow, nodes)
   _save_links(workflow, workflow_definition_root)
   _save_links(workflow, workflow_definition_root)
-  _postprocess_workflow(workflow)
 
 
   # Update schema_version
   # Update schema_version
   workflow.schema_version = 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)
     return json.loads(self.job_properties)
 
 
   def get_files(self):
   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):
   def get_archives(self):
     return json.loads(self.archives)
     return json.loads(self.archives)
@@ -788,10 +785,7 @@ class Java(Action):
     return json.loads(self.job_properties)
     return json.loads(self.job_properties)
 
 
   def get_files(self):
   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):
   def get_archives(self):
     return json.loads(self.archives)
     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))
       raise PopupException(message=msg, detail=str(ex))
 
 
     oozie_xml = self.job.to_xml(self.properties)
     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'):
     if hasattr(self.job, 'actions'):
       for action in 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)
     self.fs.create(xml_path, overwrite=True, permission=0644, data=oozie_xml)
     LOG.debug("Created %s" % (xml_path,))
     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):
   def _do_as(self, username, fn, *args, **kwargs):
     prev_user = self.fs.user
     prev_user = self.fs.user
     try:
     try: