Browse Source

HUE-1160 [oozie] Import workflow doesn't create link from start to end

Abraham Elmahrek 12 years ago
parent
commit
c6de770
2 changed files with 28 additions and 13 deletions
  1. 15 0
      apps/oozie/src/oozie/import_workflow.py
  2. 13 13
      apps/oozie/src/oozie/tests.py

+ 15 - 0
apps/oozie/src/oozie/import_workflow.py

@@ -103,9 +103,11 @@ def _save_links(workflow, root):
   workflow.end = End.objects.get(workflow=workflow).get_full_node()
   workflow.save()
 
+  _resolve_start_relationships(workflow)
   _resolve_fork_relationships(workflow)
   _resolve_decision_relationships(workflow)
 
+
 def _start_relationships(workflow, parent, child_el):
   """
   Resolve start node links.
@@ -125,6 +127,7 @@ def _start_relationships(workflow, parent, child_el):
   obj = Link.objects.create(name='to', parent=parent, child=child)
   obj.save()
 
+
 def _join_relationships(workflow, parent, child_el):
   """
   Resolves join node links.
@@ -143,6 +146,7 @@ def _join_relationships(workflow, parent, child_el):
   obj = Link.objects.create(name='to', parent=parent, child=child)
   obj.save()
 
+
 def _decision_relationships(workflow, parent, child_el):
   """
   Resolves the switch statement like nature of decision nodes.
@@ -178,6 +182,7 @@ def _decision_relationships(workflow, parent, child_el):
 
       obj.save()
 
+
 def _node_relationships(workflow, parent, child_el):
   """
   Resolves node links.
@@ -210,6 +215,16 @@ def _node_relationships(workflow, parent, child_el):
       obj.save()
 
 
+def _resolve_start_relationships(workflow):
+  if not workflow.start:
+    raise RuntimeError(_("Workflow start has not been created."))
+
+  if not workflow.end:
+    raise RuntimeError(_("Workflow end has not been created."))
+
+  obj = Link.objects.get_or_create(name='related', parent=workflow.start, child=workflow.end)
+
+
 def _resolve_fork_relationships(workflow):
   """
   Requires proper workflow structure.

+ 13 - 13
apps/oozie/src/oozie/tests.py

@@ -1449,7 +1449,7 @@ class TestImportWorkflow04(OozieMockBase):
     f.close()
     workflow.save()
     assert_equal(2, len(Node.objects.filter(workflow=workflow)))
-    assert_equal(1, len(Link.objects.filter(parent__workflow=workflow)))
+    assert_equal(2, 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)
     workflow.delete()
@@ -1466,7 +1466,7 @@ class TestImportWorkflow04(OozieMockBase):
     f.close()
     workflow.save()
     assert_equal(12, len(Node.objects.filter(workflow=workflow)))
-    assert_equal(20, len(Link.objects.filter(parent__workflow=workflow)))
+    assert_equal(21, len(Link.objects.filter(parent__workflow=workflow)))
     assert_equal(1, len(Link.objects.filter(parent__workflow=workflow, parent__node_type='decision', comment='${1 gt 2}', name='start')))
     assert_equal(1, len(Link.objects.filter(parent__workflow=workflow, parent__node_type='decision', comment='', name='start')))
     assert_equal(1, len(Link.objects.filter(parent__workflow=workflow, parent__node_type='decision', name='default')))
@@ -1482,7 +1482,7 @@ class TestImportWorkflow04(OozieMockBase):
     f.close()
     workflow.save()
     assert_equal(14, len(Node.objects.filter(workflow=workflow)))
-    assert_equal(26, len(Link.objects.filter(parent__workflow=workflow)))
+    assert_equal(27, len(Link.objects.filter(parent__workflow=workflow)))
     assert_equal(3, len(Link.objects.filter(parent__workflow=workflow, parent__node_type='decision', comment='${ 1 gt 2 }', name='start')))
     assert_equal(0, len(Link.objects.filter(parent__workflow=workflow, parent__node_type='decision', comment='', name='start')))
     assert_equal(3, len(Link.objects.filter(parent__workflow=workflow, parent__node_type='decision', name='default')))
@@ -1501,7 +1501,7 @@ class TestImportWorkflow04(OozieMockBase):
     f.close()
     workflow.save()
     assert_equal(4, len(Node.objects.filter(workflow=workflow)))
-    assert_equal(3, len(Link.objects.filter(parent__workflow=workflow)))
+    assert_equal(4, len(Link.objects.filter(parent__workflow=workflow)))
     assert_equal('[{"type":"arg","value":"-overwrite"},{"type":"arg","value":"-m"},{"type":"arg","value":"${MAP_NUMBER}"},{"type":"arg","value":"/user/hue/oozie/workspaces/data"},{"type":"arg","value":"${OUTPUT}"}]', Node.objects.get(workflow=workflow, node_type='distcp').get_full_node().params)
     workflow.delete()
 
@@ -1514,7 +1514,7 @@ class TestImportWorkflow04(OozieMockBase):
     f.close()
     workflow.save()
     assert_equal(12, len(Node.objects.filter(workflow=workflow)))
-    assert_equal(19, len(Link.objects.filter(parent__workflow=workflow)))
+    assert_equal(20, len(Link.objects.filter(parent__workflow=workflow)))
     assert_equal(6, len(Link.objects.filter(parent__workflow=workflow, parent__node_type='fork')))
     assert_equal(4, len(Link.objects.filter(parent__workflow=workflow, parent__node_type='fork', name='start')))
     assert_equal(2, len(Link.objects.filter(parent__workflow=workflow, parent__node_type='fork', child__node_type='join', name='related')))
@@ -1532,7 +1532,7 @@ class TestImportWorkflow04(OozieMockBase):
     f.close()
     workflow.save()
     assert_equal(4, len(Node.objects.filter(workflow=workflow)))
-    assert_equal(3, len(Link.objects.filter(parent__workflow=workflow)))
+    assert_equal(4, len(Link.objects.filter(parent__workflow=workflow)))
     assert_equal('[{"name":"mapred.reduce.tasks","value":"1"},{"name":"mapred.mapper.class","value":"org.apache.hadoop.examples.SleepJob"},{"name":"mapred.reducer.class","value":"org.apache.hadoop.examples.SleepJob"},{"name":"mapred.mapoutput.key.class","value":"org.apache.hadoop.io.IntWritable"},{"name":"mapred.mapoutput.value.class","value":"org.apache.hadoop.io.NullWritable"},{"name":"mapred.output.format.class","value":"org.apache.hadoop.mapred.lib.NullOutputFormat"},{"name":"mapred.input.format.class","value":"org.apache.hadoop.examples.SleepJob$SleepInputFormat"},{"name":"mapred.partitioner.class","value":"org.apache.hadoop.examples.SleepJob"},{"name":"mapred.speculative.execution","value":"false"},{"name":"sleep.job.map.sleep.time","value":"0"},{"name":"sleep.job.reduce.sleep.time","value":"1"}]', Node.objects.get(workflow=workflow, node_type='mapreduce').get_full_node().job_properties)
     workflow.delete()
 
@@ -1549,7 +1549,7 @@ class TestImportWorkflow04(OozieMockBase):
     workflow.save()
     node = Node.objects.get(workflow=workflow, node_type='pig').get_full_node()
     assert_equal(4, len(Node.objects.filter(workflow=workflow)))
-    assert_equal(3, len(Link.objects.filter(parent__workflow=workflow)))
+    assert_equal(4, len(Link.objects.filter(parent__workflow=workflow)))
     assert_equal('aggregate.pig', node.script_path)
     assert_equal('[{"type":"argument","value":"-param"},{"type":"argument","value":"INPUT=/user/hue/oozie/workspaces/data"},{"type":"argument","value":"-param"},{"type":"argument","value":"OUTPUT=${output}"}]', node.params)
     workflow.delete()
@@ -1566,7 +1566,7 @@ class TestImportWorkflow04(OozieMockBase):
     f.close()
     workflow.save()
     assert_equal(4, len(Node.objects.filter(workflow=workflow)))
-    assert_equal(3, len(Link.objects.filter(parent__workflow=workflow)))
+    assert_equal(4, len(Link.objects.filter(parent__workflow=workflow)))
     node = Node.objects.get(workflow=workflow, node_type='sqoop').get_full_node()
     assert_equal('["db.hsqldb.properties#db.hsqldb.properties","db.hsqldb.script#db.hsqldb.script"]', node.files)
     assert_equal('import --connect jdbc:hsqldb:file:db.hsqldb --table TT --target-dir ${output} -m 1', node.script_path)
@@ -1584,7 +1584,7 @@ class TestImportWorkflow04(OozieMockBase):
     f.close()
     workflow.save()
     assert_equal(5, len(Node.objects.filter(workflow=workflow)))
-    assert_equal(5, len(Link.objects.filter(parent__workflow=workflow)))
+    assert_equal(6, len(Link.objects.filter(parent__workflow=workflow)))
     nodes = [Node.objects.filter(workflow=workflow, node_type='java')[0].get_full_node(),
              Node.objects.filter(workflow=workflow, node_type='java')[1].get_full_node()]
     assert_equal('org.apache.hadoop.examples.terasort.TeraGen', nodes[0].main_class)
@@ -1604,7 +1604,7 @@ class TestImportWorkflow04(OozieMockBase):
     f.close()
     workflow.save()
     assert_equal(5, len(Node.objects.filter(workflow=workflow)))
-    assert_equal(5, len(Link.objects.filter(parent__workflow=workflow)))
+    assert_equal(6, len(Link.objects.filter(parent__workflow=workflow)))
     nodes = [Node.objects.filter(workflow=workflow, node_type='shell')[0].get_full_node(),
              Node.objects.filter(workflow=workflow, node_type='shell')[1].get_full_node()]
     assert_equal('shell-1', nodes[0].name)
@@ -1627,7 +1627,7 @@ class TestImportWorkflow04(OozieMockBase):
     f.close()
     workflow.save()
     assert_equal(4, len(Node.objects.filter(workflow=workflow)))
-    assert_equal(3, len(Link.objects.filter(parent__workflow=workflow)))
+    assert_equal(4, len(Link.objects.filter(parent__workflow=workflow)))
     node = Node.objects.get(workflow=workflow, node_type='fs').get_full_node()
     assert_equal('[{"path":"${nameNode}${output}/testfs/renamed","permissions":"700","recursive":"false"}]', node.chmods)
     assert_equal('[{"name":"${nameNode}${output}/testfs"}]', node.deletes)
@@ -1648,7 +1648,7 @@ class TestImportWorkflow04(OozieMockBase):
     f.close()
     workflow.save()
     assert_equal(4, len(Node.objects.filter(workflow=workflow)))
-    assert_equal(3, len(Link.objects.filter(parent__workflow=workflow)))
+    assert_equal(4, len(Link.objects.filter(parent__workflow=workflow)))
     node = Node.objects.get(workflow=workflow, node_type='email').get_full_node()
     assert_equal('example@example.org', node.to)
     assert_equal('', node.cc)
@@ -1668,7 +1668,7 @@ class TestImportWorkflow04(OozieMockBase):
     f.close()
     workflow.save()
     assert_equal(4, len(Node.objects.filter(workflow=workflow)))
-    assert_equal(3, len(Link.objects.filter(parent__workflow=workflow)))
+    assert_equal(4, len(Link.objects.filter(parent__workflow=workflow)))
     node = Node.objects.get(workflow=workflow, node_type='generic').get_full_node()
     assert_equal("<bleh test=\"test\">\n              <test>test</test>\n        </bleh>", node.xml)
     workflow.delete()