Browse Source

HUE-3103 [oozie] Add missing data to external workflow graph

krish 9 năm trước cách đây
mục cha
commit
84a793c

+ 9 - 1
apps/oozie/src/oozie/importlib/xslt2/workflows/0.5/action.xslt

@@ -28,7 +28,15 @@
     "ok_to": "<xsl:value-of select="*[name()=$ok]/@to"/>",
     "error_to": "<xsl:value-of select="*[name()=$error]/@to"/>"
 
-    <xsl:apply-templates select="*"/>
+    <xsl:variable name="name" select="@name"/>
+    <xsl:choose>
+      <xsl:when test="contains($name, 'streaming')">
+        <xsl:call-template name="streaming"/>
+      </xsl:when>
+      <xsl:otherwise>
+        <xsl:apply-templates select="*"/>
+      </xsl:otherwise>
+    </xsl:choose>
   },
 </xsl:template>
 

+ 11 - 7
apps/oozie/src/oozie/importlib/xslt2/workflows/0.5/extensions/distcp.0.1.xslt

@@ -4,14 +4,18 @@
 
 <xsl:template match="distcp:distcp">
 
-  ,"distcp": {
-        <xsl:for-each select="arg">
-          "path<xsl:value-of select='position()'/>": "<xsl:value-of select="arg"/>"
-          <xsl:if  test="position() &lt; last()">
-            ,
-          </xsl:if>
+  ,"params": [
+        <xsl:for-each select="*[local-name()='arg']">
+          <xsl:choose>
+            <xsl:when test="position() &lt; last()">
+              {"type":"arg","value":"<xsl:value-of select="."/>"},
+            </xsl:when>
+            <xsl:otherwise>
+              {"type":"arg","value":"<xsl:value-of select="."/>"}
+            </xsl:otherwise>
+          </xsl:choose>
         </xsl:for-each>
-    }
+    ]
 
 </xsl:template>
 

+ 1 - 1
apps/oozie/src/oozie/importlib/xslt2/workflows/0.5/extensions/sqoop.0.1.xslt

@@ -6,7 +6,7 @@
 
 <xsl:template match="sqoop:sqoop">
 
-  ,"sqoop": {<xsl:call-template name="script_path"/>}
+  ,"sqoop": { "command": "<xsl:value-of select="*[local-name()='command']"/>"}
 
 </xsl:template>
 

+ 1 - 1
apps/oozie/src/oozie/importlib/xslt2/workflows/0.5/extensions/sqoop.0.2.xslt

@@ -4,7 +4,7 @@
 
 <xsl:template match="sqoop:sqoop">
 
-  ,"sqoop": {"script_path": "<xsl:value-of select="*[local-name()='command']"/>"}
+  ,"sqoop": {"command": "<xsl:value-of select="*[local-name()='command']"/>"}
 
 </xsl:template>
 

+ 4 - 5
apps/oozie/src/oozie/importlib/xslt2/workflows/0.5/nodes/fields/job_properties.xslt

@@ -3,19 +3,18 @@
 <xsl:stylesheet version="1.0" xmlns:xsl="http://www.w3.org/1999/XSL/Transform" xmlns:workflow="uri:oozie:workflow:0.5" exclude-result-prefixes="workflow">
 
 <xsl:template name="job_properties">
-  "job_properties":
-    <xsl:text>[</xsl:text>
+  ,"job_properties": [
     <xsl:for-each select="*[local-name()='configuration']/*[local-name()='property']">
       <xsl:choose>
         <xsl:when test="position() &lt; last()">
-          <xsl:text><![CDATA[{"name":"]]></xsl:text><xsl:value-of select="*[local-name()='name']" /><xsl:text><![CDATA[","value":"]]></xsl:text><xsl:value-of select="*[local-name()='value']" /><xsl:text><![CDATA["},]]></xsl:text>
+          {"name": "<xsl:value-of select="*[local-name()='name']" />", "value": "<xsl:value-of select="*[local-name()='value']" />"},
         </xsl:when>
         <xsl:otherwise>
-          <xsl:text><![CDATA[{"name":"]]></xsl:text><xsl:value-of select="*[local-name()='name']" /><xsl:text><![CDATA[","value":"]]></xsl:text><xsl:value-of select="*[local-name() ='value']" /><xsl:text><![CDATA["}]]></xsl:text>
+          {"name": "<xsl:value-of select="*[local-name()='name']" />", "value": "<xsl:value-of select="*[local-name()='value']" />"}
         </xsl:otherwise>
       </xsl:choose>
     </xsl:for-each>
-    <xsl:text>]</xsl:text>
+    ]
 </xsl:template>
 
 </xsl:stylesheet>

+ 3 - 1
apps/oozie/src/oozie/importlib/xslt2/workflows/0.5/nodes/java.xslt

@@ -6,7 +6,9 @@
 
 <xsl:template match="workflow:java" xmlns:workflow="uri:oozie:workflow:0.5">
 
-  ,"java": {<xsl:call-template name="jar_path"/>}
+  ,"java": {
+        "main-class": "<xsl:value-of select="*[local-name()='main-class']"/>"
+    }
 
 
 </xsl:template>

+ 2 - 2
apps/oozie/src/oozie/importlib/xslt2/workflows/0.5/nodes/mapreduce.xslt

@@ -2,11 +2,11 @@
 
 <xsl:stylesheet version="1.0" xmlns:xsl="http://www.w3.org/1999/XSL/Transform" xmlns:workflow="uri:oozie:workflow:0.5" exclude-result-prefixes="workflow">
 
-<xsl:import href="fields/jar_path.xslt"/>
+<xsl:import href="fields/job_properties.xslt"/>
 
 <xsl:template match="workflow:map-reduce" xmlns:workflow="uri:oozie:workflow:0.5">
 
-  ,"mapreduce": {<xsl:call-template name="jar_path"/>}
+    <xsl:call-template name="job_properties"/>
 
 </xsl:template>
 

+ 3 - 3
apps/oozie/src/oozie/importlib/xslt2/workflows/0.5/nodes/streaming.xslt

@@ -5,11 +5,11 @@
 <xsl:import href="fields/mapper.xslt"/>
 <xsl:import href="fields/reducer.xslt"/>
 
-<xsl:template match="workflow:streaming" xmlns:workflow="uri:oozie:workflow:0.5">
+<xsl:template name="streaming">
 
   ,"streaming": {
-        <xsl:call-template name="mapper"/>,
-        <xsl:call-template name="reducer"/>
+        "mapper": "<xsl:value-of select="*[local-name()='map-reduce']/*[local-name()='streaming']/*[local-name()='mapper']"/>",
+        "reducer": "<xsl:value-of select="*[local-name()='map-reduce']/*[local-name()='streaming']/*[local-name()='reducer']"/>"
     }
 
 </xsl:template>

+ 1 - 1
apps/oozie/src/oozie/importlib/xslt2/workflows/0.5/nodes/subworkflow.xslt

@@ -3,7 +3,7 @@
 <xsl:stylesheet version="1.0" xmlns:xsl="http://www.w3.org/1999/XSL/Transform" xmlns:workflow="uri:oozie:workflow:0.5" exclude-result-prefixes="workflow">
 
 <xsl:template match="workflow:sub-workflow" xmlns:workflow="uri:oozie:workflow:0.5">
-  , "sub-workflow": { "app-path":"<xsl:value-of select="*[local-name()='app-path']"/>" }
+  , "subworkflow": { "app-path":"<xsl:value-of select="*[local-name()='app-path']"/>" }
 </xsl:template>
 
 </xsl:stylesheet>

+ 20 - 14
apps/oozie/src/oozie/models2.py

@@ -366,7 +366,7 @@ class Workflow(Job):
     node_hierarchy = ['start']
     _get_hierarchy_from_adj_list(adj_list, adj_list['start']['ok_to'], node_hierarchy)
 
-    _add_uuids_to_adj_list(adj_list)
+    _update_adj_list(adj_list)
 
     wf_rows = _create_workflow_layout(node_hierarchy, adj_list)
     data = {'layout': [{}], 'workflow': {}}
@@ -434,12 +434,21 @@ class Workflow(Job):
 
     return data
 
-def _add_uuids_to_adj_list(adj_list):
+def _update_adj_list(adj_list):
   uuids = {}
   id = 1
   for node in adj_list.keys():
     adj_list[node]['id'] = id
 
+    # Oozie uses same action for streaming and mapreduce but Hue manages them differently
+    if adj_list[node]['node_type'] == 'map-reduce':
+      if 'streaming' in adj_list[node]['name']:
+        adj_list[node]['node_type'] = 'streaming'
+      else:
+        adj_list[node]['node_type'] = 'mapreduce'
+    elif adj_list[node]['node_type'] == 'sub-workflow':
+      adj_list[node]['node_type'] = 'subworkflow'
+
     if adj_list[node]['node_type'] == 'kill':
       adj_list[node]['uuid'] = '17c9c895-5a16-7443-bb81-f34b30b21548'
     elif adj_list[node]['node_type'] == 'start':
@@ -462,19 +471,18 @@ def _dig_nodes(nodes, adj_list, user, wf_nodes):
         properties = dict(NODES['%s-widget' % node['node_type']].get_fields())
 
       if node['node_type'] == 'pig':
-        properties['script_path'] = node.get('script_path')
+        properties['script_path'] = node.get('pig').get('script_path')
       elif node['node_type'] == 'spark':
         properties['class'] = node.get('spark').get('class')
-        #TBD: jar_path
-        properties['jar_path'] = node.get('spark').get('jar')
+        properties['jars'] = node.get('spark').get('jar')
       elif node['node_type'] == 'hive' or node['node_type'] == 'hive2':
         properties['script_path'] = node.get('hive2').get('script')
       elif node['node_type'] == 'java':
-        properties['jar_path'] = node.get('jar_path')
+        properties['main_class'] = node.get('java').get('main-class')
       elif node['node_type'] == 'sqoop':
-        properties['command'] = node.get('script_path')
+        properties['command'] = node.get('sqoop').get('command')
       elif node['node_type'] == 'mapreduce':
-        properties['jar_path'] = node.get('jar_path')
+        properties['job_properties'] = node.get('job_properties')
       elif node['node_type'] == 'shell':
         properties['shell_command'] = node.get('shell').get('command')
       elif node['node_type'] == 'ssh':
@@ -495,11 +503,9 @@ def _dig_nodes(nodes, adj_list, user, wf_nodes):
         properties['mapper'] = node.get('streaming').get('mapper')
         properties['reducer'] = node.get('streaming').get('reducer')
       elif node['node_type'] == 'distcp':
-        #TBD: both
-        properties['source'] = node.get('distcp').get('source')
-        properties['destination'] = node.get('distcp').get('source')
-      elif node['node_type'] == 'sub-workflow':
-        properties['app-path'] = node.get('sub-workflow').get('app-path')
+        properties['distcp_parameters'] = node.get('params')
+      elif node['node_type'] == 'subworkflow':
+        properties['app-path'] = node.get('subworkflow').get('app-path')
         properties['workflow'] = node.get('uuid')
         properties['job_properties'] = []
         properties['sla'] = ''
@@ -531,7 +537,7 @@ def _create_workflow_layout(nodes, adj_list, size=12):
     if type(node) == list and len(node) == 1:
       node = node[0]
     if type(node) != list:
-      wf_rows.append({"widgets":[{"size":size, "name": adj_list[node]['node_type'], "id":  adj_list[node]['uuid'], "widgetType": "%s-widget" % adj_list[node]['node_type'] if adj_list[node]['node_type'] != 'sub-workflow' else 'subworkflow-widget', "properties":{}, "offset":0, "isLoading":False, "klass":"card card-widget span%s" % size, "columns":[]}]})
+      wf_rows.append({"widgets":[{"size":size, "name": adj_list[node]['node_type'], "id":  adj_list[node]['uuid'], "widgetType": "%s-widget" % adj_list[node]['node_type'], "properties":{}, "offset":0, "isLoading":False, "klass":"card card-widget span%s" % size, "columns":[]}]})
     else:
       if adj_list[node[0]]['node_type'] == 'fork':
         wf_rows.append({"widgets":[{"size":size, "name": 'Fork', "id":  adj_list[node[0]]['uuid'], "widgetType": "%s-widget" % adj_list[node[0]]['node_type'], "properties":{}, "offset":0, "isLoading":False, "klass":"card card-widget span%s" % size, "columns":[]}]})

+ 2 - 2
apps/oozie/src/oozie/tests2.py

@@ -25,7 +25,7 @@ from nose.tools import assert_true, assert_false, assert_equal, assert_not_equal
 
 from oozie.conf import ENABLE_V2
 from oozie.importlib.workflows import generate_v2_graph_nodes
-from oozie.models2 import Workflow, find_dollar_variables, find_dollar_braced_variables, Node, _create_graph_adjaceny_list, _get_hierarchy_from_adj_list, _create_workflow_layout, _add_uuids_to_adj_list
+from oozie.models2 import Workflow, find_dollar_variables, find_dollar_braced_variables, Node, _create_graph_adjaceny_list, _get_hierarchy_from_adj_list, _create_workflow_layout
 from oozie.tests import OozieMockBase, save_temp_workflow, MockOozieApi
 
 
@@ -303,7 +303,7 @@ class TestExternalWorkflowGraph():
   def test_graph_generation_from_xml(self):
     f = open('apps/oozie/src/oozie/test_data/xslt2/test-workflow.xml')
     self.wf.definition = f.read()
-    self.node_list = [{u'node_type': u'start', u'ok_to': u'fork-68d4', u'name': u''}, {u'node_type': u'kill', u'ok_to': u'', u'name': u'Kill'}, {u'path2': u'shell-0f44', u'node_type': u'fork', u'ok_to': u'', u'name': u'fork-68d4', u'path1': u'subworkflow-a13f'}, {u'node_type': u'join', u'ok_to': u'End', u'name': u'join-775e'}, {u'node_type': u'end', u'ok_to': u'', u'name': u'End'}, {u'node_type': u'sub-workflow', u'ok_to': u'join-775e', u'sub-workflow': {u'app-path': u'${nameNode}/user/hue/oozie/deployments/_admin_-oozie-50001-1427488969.48'}, u'name': u'subworkflow-a13f', u'error_to': u'Kill'}, {u'shell': {u'command': u'ls'}, u'node_type': u'shell', u'ok_to': u'join-775e', u'name': u'shell-0f44', u'error_to': u'Kill'}, {}]
+    self.node_list = [{u'node_type': u'start', u'ok_to': u'fork-68d4', u'name': u''}, {u'node_type': u'kill', u'ok_to': u'', u'name': u'Kill'}, {u'path2': u'shell-0f44', u'node_type': u'fork', u'ok_to': u'', u'name': u'fork-68d4', u'path1': u'subworkflow-a13f'}, {u'node_type': u'join', u'ok_to': u'End', u'name': u'join-775e'}, {u'node_type': u'end', u'ok_to': u'', u'name': u'End'}, {u'subworkflow': {u'app-path': u'${nameNode}/user/hue/oozie/deployments/_admin_-oozie-50001-1427488969.48'}, u'node_type': u'sub-workflow', u'ok_to': u'join-775e', u'name': u'subworkflow-a13f', u'error_to': u'Kill'}, {u'shell': {u'command': u'ls'}, u'node_type': u'shell', u'ok_to': u'join-775e', u'name': u'shell-0f44', u'error_to': u'Kill'}, {}]
     assert_equal(self.node_list, generate_v2_graph_nodes(self.wf.definition))
 
   def test_get_graph_adjacency_list(self):