浏览代码

HUE-822 [oozie] Prepare statements should follow fs action convention

Works when not using a variable as the path
Romain Rigaux 13 年之前
父节点
当前提交
dc7f032
共有 2 个文件被更改,包括 39 次插入31 次删除
  1. 13 2
      apps/oozie/src/oozie/templates/editor/gen/workflow-common.xml.mako
  2. 26 29
      apps/oozie/src/oozie/tests.py

+ 13 - 2
apps/oozie/src/oozie/templates/editor/gen/workflow-common.xml.mako

@@ -26,7 +26,18 @@ import posixpath
         % if prepares:
         % if prepares:
             <prepare>
             <prepare>
                 % for p in prepares:
                 % for p in prepares:
-                <${ p['type'] } path="${'${'}nameNode}${ p['value'] }"/>
+                  <%
+                    # Same as Fs action convention. No change if path is just a parameter.
+                    operation = p['type']
+                    path = p['value']
+
+                    if not path.startswith('/') and not path.startswith('$') and not path.startswith('hdfs://'):
+                      path = '/user/%(username)s/%(path)s' % {'username': '${wf:user()}', 'path': path}
+
+                    path = '%(nameNode)s%(path)s' % {'nameNode': '${nameNode}', 'path': path}
+                  %>
+
+                  <${ operation } path="${ path }"/>
                 % endfor
                 % endfor
             </prepare>
             </prepare>
         % endif
         % endif
@@ -50,7 +61,7 @@ import posixpath
 <%def name="distributed_cache(files, archives)">
 <%def name="distributed_cache(files, archives)">
     % for f in files:
     % for f in files:
         % if len(f) != 0:
         % if len(f) != 0:
-            <file>${ f + '#' + posixpath.basename(f) }</file>
+            <file>${ filelink(f) }</file>
         % endif
         % endif
     % endfor
     % endfor
     % for a in archives:
     % for a in archives:

+ 26 - 29
apps/oozie/src/oozie/tests.py

@@ -114,7 +114,7 @@ class MockOozieApi:
     return '<xml></xml>'
     return '<xml></xml>'
 
 
 
 
-class OozieMockBase:
+class OozieMockBase(object):
 
 
   def setUp(self):
   def setUp(self):
     # Beware: Monkey patch Oozie/LibOozie with Mock API
     # Beware: Monkey patch Oozie/LibOozie with Mock API
@@ -143,6 +143,7 @@ class OozieMockBase:
     """ Creates a linear workflow """
     """ Creates a linear workflow """
     Link.objects.filter(parent__workflow=self.wf).delete()
     Link.objects.filter(parent__workflow=self.wf).delete()
     Link(parent=self.wf.start, child=self.wf.end, name="related").save()
     Link(parent=self.wf.start, child=self.wf.end, name="related").save()
+
     action1 = add_node(self.wf, 'action-name-1', 'mapreduce', [self.wf.start], {
     action1 = add_node(self.wf, 'action-name-1', 'mapreduce', [self.wf.start], {
       'description': '',
       'description': '',
       'files': '[]',
       'files': '[]',
@@ -167,6 +168,7 @@ class OozieMockBase:
       'prepares': '[{"value":"${output}","type":"delete"},{"value":"/test","type":"mkdir"}]',
       'prepares': '[{"value":"${output}","type":"delete"},{"value":"/test","type":"mkdir"}]',
       'archives': '[]',
       'archives': '[]',
     })
     })
+
     Link(parent=action3, child=self.wf.end, name="ok").save()
     Link(parent=action3, child=self.wf.end, name="ok").save()
 
 
 
 
@@ -335,7 +337,7 @@ class TestAPI(OozieMockBase):
     assert_equal(0, test_response_json_object['status'])
     assert_equal(0, test_response_json_object['status'])
 
 
 
 
-class TestAPIPermissionsWithOozie(OozieBase):
+class TestApiPermissionsWithOozie(OozieBase):
 
 
   def setUp(self):
   def setUp(self):
     OozieBase.setUp(self)
     OozieBase.setUp(self)
@@ -417,9 +419,12 @@ class TestAPIPermissionsWithOozie(OozieBase):
 
 
 class TestEditor(OozieMockBase):
 class TestEditor(OozieMockBase):
 
 
-  def test_find_parameters(self):
+  def setUp(self):
+    super(TestEditor, self).setUp()
     self.setup_simple_workflow()
     self.setup_simple_workflow()
 
 
+
+  def test_find_parameters(self):
     jobs = [Job(name="$a"),
     jobs = [Job(name="$a"),
             Job(name="foo ${b} $$"),
             Job(name="foo ${b} $$"),
             Job(name="${foo}", description="xxx ${foo}")]
             Job(name="${foo}", description="xxx ${foo}")]
@@ -429,15 +434,11 @@ class TestEditor(OozieMockBase):
 
 
 
 
   def test_find_all_parameters(self):
   def test_find_all_parameters(self):
-    self.setup_simple_workflow()
-
     assert_equal([{'name': u'output', 'value': u''}, {'name': u'SLEEP', 'value': ''}, {'name': u'market', 'value': u'US'}],
     assert_equal([{'name': u'output', 'value': u''}, {'name': u'SLEEP', 'value': ''}, {'name': u'market', 'value': u'US'}],
                  self.wf.find_all_parameters())
                  self.wf.find_all_parameters())
 
 
 
 
   def test_workflow_has_cycle(self):
   def test_workflow_has_cycle(self):
-    self.setup_simple_workflow()
-
     action1 = Node.objects.get(name='action-name-1')
     action1 = Node.objects.get(name='action-name-1')
     action3 = Node.objects.get(name='action-name-3')
     action3 = Node.objects.get(name='action-name-3')
 
 
@@ -451,8 +452,6 @@ class TestEditor(OozieMockBase):
 
 
 
 
   def test_workflow_gen_xml(self):
   def test_workflow_gen_xml(self):
-    self.setup_simple_workflow()
-
     assert_equal(
     assert_equal(
         '<workflow-app name="wf-name-1" xmlns="uri:oozie:workflow:0.2">\n'
         '<workflow-app name="wf-name-1" xmlns="uri:oozie:workflow:0.2">\n'
         '    <global>\n'
         '    <global>\n'
@@ -527,8 +526,6 @@ class TestEditor(OozieMockBase):
 
 
 
 
   def test_workflow_flatten_list(self):
   def test_workflow_flatten_list(self):
-    self.setup_simple_workflow()
-
     assert_equal('[<Start: start>, <Mapreduce: action-name-1>, <Mapreduce: action-name-2>, <Mapreduce: action-name-3>, '
     assert_equal('[<Start: start>, <Mapreduce: action-name-1>, <Mapreduce: action-name-2>, <Mapreduce: action-name-3>, '
                  '<Kill: kill>, <End: end>]',
                  '<Kill: kill>, <End: end>]',
                  str(self.wf.node_list))
                  str(self.wf.node_list))
@@ -543,14 +540,10 @@ class TestEditor(OozieMockBase):
 
 
 
 
   def test_create_coordinator(self):
   def test_create_coordinator(self):
-    self.setup_simple_workflow()
-
     create_coordinator(self.wf, self.c)
     create_coordinator(self.wf, self.c)
 
 
 
 
   def test_clone_coordinator(self):
   def test_clone_coordinator(self):
-    self.setup_simple_workflow()
-
     coord = create_coordinator(self.wf, self.c)
     coord = create_coordinator(self.wf, self.c)
     coordinator_count = Coordinator.objects.count()
     coordinator_count = Coordinator.objects.count()
 
 
@@ -581,8 +574,6 @@ class TestEditor(OozieMockBase):
 
 
 
 
   def test_coordinator_workflow_access_permissions(self):
   def test_coordinator_workflow_access_permissions(self):
-    self.setup_simple_workflow()
-
     self.wf.is_shared = True
     self.wf.is_shared = True
     self.wf.save()
     self.wf.save()
 
 
@@ -629,8 +620,6 @@ class TestEditor(OozieMockBase):
 
 
 
 
   def test_coordinator_gen_xml(self):
   def test_coordinator_gen_xml(self):
-    self.setup_simple_workflow()
-
     coord = create_coordinator(self.wf, self.c)
     coord = create_coordinator(self.wf, self.c)
 
 
     assert_equal(
     assert_equal(
@@ -653,8 +642,6 @@ class TestEditor(OozieMockBase):
 
 
 
 
   def test_coordinator_with_data_input_gen_xml(self):
   def test_coordinator_with_data_input_gen_xml(self):
-    self.setup_simple_workflow()
-
     coord = create_coordinator(self.wf, self.c)
     coord = create_coordinator(self.wf, self.c)
     create_dataset(coord, self.c)
     create_dataset(coord, self.c)
     create_coordinator_data(coord, self.c)
     create_coordinator_data(coord, self.c)
@@ -695,15 +682,11 @@ class TestEditor(OozieMockBase):
 
 
 
 
   def test_create_coordinator_dataset(self):
   def test_create_coordinator_dataset(self):
-    self.setup_simple_workflow()
-
     coord = create_coordinator(self.wf, self.c)
     coord = create_coordinator(self.wf, self.c)
     create_dataset(coord, self.c)
     create_dataset(coord, self.c)
 
 
 
 
   def test_create_coordinator_input_data(self):
   def test_create_coordinator_input_data(self):
-    self.setup_simple_workflow()
-
     coord = create_coordinator(self.wf, self.c)
     coord = create_coordinator(self.wf, self.c)
     create_dataset(coord, self.c)
     create_dataset(coord, self.c)
 
 
@@ -714,16 +697,30 @@ class TestEditor(OozieMockBase):
     self.c.post(reverse('oozie:setup_app'))
     self.c.post(reverse('oozie:setup_app'))
 
 
 
 
-  def test_get_workflow_parameters(self):
-    self.setup_simple_workflow()
+  def test_workflow_prepare(self):
+    action1 = Node.objects.get(name='action-name-1').get_full_node()
+
+    action1.prepares = json.dumps([
+                           {"type": "delete","value": "${output}"},
+                           {"type": "delete","value": "out"},
+                           {"type": "delete","value": "/user/test/out"},
+                           {"type": "delete","value": "hdfs://localhost:8020/user/test/out"}])
+    action1.save()
+
+    xml = self.wf.to_xml()
+
+    assert_true('<delete path="${nameNode}${output}"/>' in xml, xml)
+    assert_true(re.search(re.escape('<delete path="${nameNode}/user/${wf:user()}/out"/>'), xml, re.IGNORECASE), xml)
+    assert_true(re.search(re.escape('<delete path="${nameNode}/user/test/out"/>'), xml, re.IGNORECASE), xml)
+    assert_true(re.search(re.escape('<delete path="hdfs://localhost:8020/user/test/out"/>'), xml, re.IGNORECASE), xml)
 
 
+
+  def test_get_workflow_parameters(self):
     assert_equal([{'name': u'output', 'value': ''}, {'name': u'SLEEP', 'value': ''}, {'name': u'market', 'value': u'US'}],
     assert_equal([{'name': u'output', 'value': ''}, {'name': u'SLEEP', 'value': ''}, {'name': u'market', 'value': u'US'}],
                  self.wf.find_all_parameters())
                  self.wf.find_all_parameters())
 
 
 
 
   def test_get_coordinator_parameters(self):
   def test_get_coordinator_parameters(self):
-    self.setup_simple_workflow()
-
     coord = create_coordinator(self.wf, self.c)
     coord = create_coordinator(self.wf, self.c)
 
 
     create_dataset(coord, self.c)
     create_dataset(coord, self.c)