瀏覽代碼

HUE-1231 [oozie] Support global configuration when importing

Section for assigning workflow properties.
Skip job-tracker and name-node for now.
Abraham Elmahrek 11 年之前
父節點
當前提交
5300fba862

+ 47 - 1
apps/oozie/src/oozie/importlib/workflows.py

@@ -31,6 +31,7 @@ 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
 import logging
 from lxml import etree
 import os
@@ -53,6 +54,50 @@ OOZIE_NAMESPACES = ['uri:oozie:workflow:0.1', 'uri:oozie:workflow:0.2', 'uri:ooz
 LINKS = ('ok', 'error', 'path')
 
 
+def _set_properties(workflow, root, namespace):
+  # root should be config element.
+  properties = []
+  seen = {}
+  namespaces = {
+    'n': namespace
+  }
+
+  for prop in root.xpath('n:property', namespaces=namespaces):
+    name = prop.xpath('n:name', namespaces=namespaces)[0].text
+    value = prop.xpath('n:value', namespaces=namespaces)[0].text
+    if name not in seen:
+      properties.append({'name': name, 'value': value})
+      seen[name] = True
+
+  workflow.job_properties = json.dumps(properties)
+
+
+def _global_configuration(workflow, root, namespace):
+  # root should be global config element.
+  namespaces = {
+    'n': namespace
+  }
+
+  job_xml = root.xpath('n:job-xml', namespaces=namespaces)
+  configuration = root.xpath('n:configuration', namespaces=namespaces)
+  if job_xml:
+    workflow.job_xml = job_xml[0].text
+  if configuration:
+    _set_properties(workflow, configuration[0], namespace)
+
+
+def _assign_workflow_properties(workflow, root, namespace):
+  namespaces = {
+    'n': namespace
+  }
+
+  global_config = root.xpath('n:global', namespaces=namespaces)
+  if global_config:
+    _global_configuration(workflow, global_config[0], namespace)
+
+  LOG.debug("Finished assigning properties to workflow %s" % smart_str(workflow.name))
+
+
 def _save_links(workflow, root):
   """
   Iterates over all links in the passed XML doc and creates links.
@@ -77,7 +122,7 @@ def _save_links(workflow, root):
 
   Note: The nodes that these links point to should exist already.
   Note: Nodes are looked up by workflow and name.
-  Note: Skip global configuration explicitly. Unknown knows should throw an error.
+  Note: Unknown elements should throw an error.
   """
   # Iterate over nodes
   for child_el in root:
@@ -608,6 +653,7 @@ def import_workflow_root(workflow, workflow_definition_root, metadata=None, fs=N
     _preprocess_nodes(workflow, transformed_root, workflow_definition_root, nodes, fs)
     _save_nodes(workflow, nodes)
     _save_links(workflow, workflow_definition_root)
+    _assign_workflow_properties(workflow, workflow_definition_root, schema_version)
     if metadata:
       _process_metadata(workflow, metadata)
 

+ 64 - 2
apps/oozie/src/oozie/test_data/workflows/0.4/test-basic-global-config.xml

@@ -1,5 +1,67 @@
 <workflow-app name="test-workflow" xmlns="uri:oozie:workflow:0.4">
-  <global></global>
-  <start to="done"/>
+  <global>
+    <job-tracker>${job-tracker}</job-tracker>
+    <name-node>${namd-node}</name-node>
+    <job-xml>job1.xml</job-xml>
+    <configuration>
+      <property>
+        <name>mapred.job.queue.name</name>
+        <value>${queueName}</value>
+      </property>
+    </configuration>
+  </global>
+  <start to="Sleep-1"/>
+  <action name="Sleep-1">
+    <map-reduce>
+      <configuration>
+        <property>
+          <name>mapred.reduce.tasks</name>
+          <value>1</value>
+        </property>
+        <property>
+          <name>mapred.mapper.class</name>
+          <value>org.apache.hadoop.examples.SleepJob</value>
+        </property>
+        <property>
+          <name>mapred.reducer.class</name>
+          <value>org.apache.hadoop.examples.SleepJob</value>
+        </property>
+        <property>
+          <name>mapred.mapoutput.key.class</name>
+          <value>org.apache.hadoop.io.IntWritable</value>
+        </property>
+        <property>
+          <name>mapred.mapoutput.value.class</name>
+          <value>org.apache.hadoop.io.NullWritable</value>
+        </property>
+        <property>
+          <name>mapred.output.format.class</name>
+          <value>org.apache.hadoop.mapred.lib.NullOutputFormat</value>
+        </property>
+        <property>
+          <name>mapred.input.format.class</name>
+          <value>org.apache.hadoop.examples.SleepJob$SleepInputFormat</value>
+        </property>
+        <property>
+          <name>mapred.partitioner.class</name>
+          <value>org.apache.hadoop.examples.SleepJob</value>
+        </property>
+        <property>
+          <name>mapred.speculative.execution</name>
+          <value>false</value>
+        </property>
+        <property>
+          <name>sleep.job.map.sleep.time</name>
+          <value>0</value>
+        </property>
+        <property>
+          <name>sleep.job.reduce.sleep.time</name>
+          <value>1</value>
+        </property>
+      </configuration>
+    </map-reduce>
+    <ok to="done"/>
+    <error to="kill"/>
+  </action>
   <end name="done"/>
 </workflow-app>

+ 4 - 2
apps/oozie/src/oozie/tests.py

@@ -2115,10 +2115,12 @@ class TestImportWorkflow04(OozieMockBase):
     import_workflow(workflow, f.read())
     f.close()
     workflow.save()
-    assert_equal(2, len(Node.objects.filter(workflow=workflow)))
-    assert_equal(2, len(Link.objects.filter(parent__workflow=workflow)))
+    assert_equal(4, len(Node.objects.filter(workflow=workflow)))
+    assert_equal(4, len(Link.objects.filter(parent__workflow=workflow)))
     assert_equal('done', Node.objects.get(workflow=workflow, node_type='end').name)
     assert_equal('uri:oozie:workflow:0.4', workflow.schema_version)
+    assert_equal('job1.xml', workflow.job_xml)
+    assert_equal('[{"name": "mapred.job.queue.name", "value": "${queueName}"}]', workflow.job_properties)
     workflow.delete(skip_trash=True)