瀏覽代碼

HUE-4526 [oozie] API to generate sequential workflow with sequential document actions

Romain Rigaux 9 年之前
父節點
當前提交
a72a419
共有 3 個文件被更改,包括 148 次插入78 次删除
  1. 87 62
      apps/oozie/src/oozie/models2.py
  2. 39 3
      apps/oozie/src/oozie/tests2.py
  3. 22 13
      apps/oozie/src/oozie/views/editor2.py

+ 87 - 62
apps/oozie/src/oozie/models2.py

@@ -48,6 +48,7 @@ from notebook.models import Notebook
 from oozie.conf import REMOTE_SAMPLE_DIR
 from oozie.utils import utc_datetime_format, UTC_TIME_FORMAT, convert_to_server_timezone
 from oozie.importlib.workflows import generate_v2_graph_nodes, MalformedWfDefException, InvalidTagWithNamespaceException
+# from oozie.views.editor2 import _save_workflow
 
 
 LOG = logging.getLogger(__name__)
@@ -2916,22 +2917,29 @@ class History(object):
 
 class WorkflowBuilder():
   """
-  Focus on building nodes, not the UI layout (should be graphed automatically in dashboard).
-  Only support Hive document currently, but then will have Pig, PySpark, MapReduce...
+  Building a workflow that has saved Documents for nodes (e.g Saved Hive query, saved Pig script...).
   """
 
-  def create_workflow(self, document, user, name=None, managed=False):
+  def create_workflow(self, user, document=None, documents=None, name=None, managed=False):
+    nodes = []
+    if documents is None:
+      documents = [document]
 
     if name is None:
-      name = _('Schedule of ') + document.name or document.type
+      name = _('Schedule of ') + ','.join([document.name or document.type for document in documents])
 
-    if document.type == 'query-java':
-      node = self.get_java_document_node(document, name)
-    else:
-      node = self.get_hive_document_node(document, name, user)
+    for document in documents:     
+      if document.type == 'query-java':
+        node = self.get_java_document_node(document, name)
+      else:
+        node = self.get_hive_document_node(document, name, user)
+  
+      nodes.append(node)
 
-    workflow_doc = self.get_workflow(node, name, document.uuid, user, managed=managed)
-    workflow_doc.dependencies.add(document)
+    workflow_doc = self.get_workflow(nodes, name, document.uuid, user, managed=managed)
+    
+    for document in documents:
+      workflow_doc.dependencies.add(document)
 
     return workflow_doc
 
@@ -2946,7 +2954,9 @@ class WorkflowBuilder():
     parameters = [{u'value': u'%s=${%s}' % (p, p)} for p in parameters]
 
     return {
-        u'name': u'doc-1',
+        u'name': u'doc-hive-%s' % document.uuid[:4],
+        u'id': str(uuid.uuid4()),
+        u'type': u'hive-document-widget',        
         u'actionParametersUI': [],
         u'properties': {
             u'files': [],
@@ -2966,7 +2976,7 @@ class WorkflowBuilder():
                 {u'key': u'alert-contact', u'value': u''},
                 {u'key': u'notification-msg', u'value': u''},
                 {u'key': u'upstream-apps', u'value': u''},
-                ],
+            ],
             u'archives': [],
             u'prepares': [],
             u'credentials': credentials,
@@ -2974,11 +2984,10 @@ class WorkflowBuilder():
             u'jdbc_url': u'',
             },
         u'actionParametersFetched': False,
-        u'id': u'0aec471d-2b7c-d93d-b22c-2110fd17ea2c',
-        u'type': u'hive-document-widget',
-        u'children': [{u'to': u'33430f0f-ebfa-c3ec-f237-3e77efa03d0a'},
-                      {u'error': u'17c9c895-5a16-7443-bb81-f34b30b21548'
-                      }],
+        u'children': [
+            {u'to': u'33430f0f-ebfa-c3ec-f237-3e77efa03d0a'},
+            {u'error': u'17c9c895-5a16-7443-bb81-f34b30b21548'
+        }],
         u'actionParameters': [],
     }
 
@@ -2986,11 +2995,11 @@ class WorkflowBuilder():
     credentials = []
 
     return {
-         "id":"0aec471d-2b7c-d93d-b22c-2110fd17ea2c",
-         "name":"doc-1",
-         "type":"java-document-widget",
-         "properties":{
-              u'uuid': document.uuid, # files, main_class, arguments comes from there
+        "id": str(uuid.uuid4()),
+        'name': u'doc-hive-%s' % document.uuid[:4],
+        "type":"java-document-widget",
+        "properties":{
+              u'uuid': document.uuid, # Files, main_class, arguments comes from there
               "job_xml":[],
               "jar_path": "",
               "java_opts":[],
@@ -3002,48 +3011,52 @@ class WorkflowBuilder():
               "credentials": credentials,
               "sla":[{"value":False, "key":"enabled"}, {"value":"${nominal_time}", "key":"nominal-time"}, {"value":"", "key":"should-start"}, {"value":"${30 * MINUTES}", "key":"should-end"}, {"value":"", "key":"max-duration"}, {"value":"", "key":"alert-events"}, {"value":"", "key":"alert-contact"}, {"value":"", "key":"notification-msg"}, {"value":"", "key":"upstream-apps"}],
               "archives":[]
-          },
-          "children":[
-              {"to":"33430f0f-ebfa-c3ec-f237-3e77efa03d0a"},
-              {"error":"17c9c895-5a16-7443-bb81-f34b30b21548"}],
-          "actionParameters":[],
-          "actionParametersFetched": False
+        },
+        "children":[
+            {"to":"33430f0f-ebfa-c3ec-f237-3e77efa03d0a"},
+            {"error":"17c9c895-5a16-7443-bb81-f34b30b21548"}
+        ],
+        "actionParameters":[],
+        "actionParametersFetched": False
     }
 
 
 
-  def get_workflow(self, node, name, doc_uuid, user, managed=False):
+  def get_workflow(self, nodes, name, doc_uuid, user, managed=False):
     parameters = []
 
-    data = json.dumps({'workflow': {
+    data = {
+      'workflow': {
       u'name': name,
-      u'versions': [u'uri:oozie:workflow:0.4', u'uri:oozie:workflow:0.4.5'
-                    , u'uri:oozie:workflow:0.5'],
-      u'isDirty': False,
-      u'movedNode': None,
-      u'linkMapping': {
-          u'33430f0f-ebfa-c3ec-f237-3e77efa03d0a': [],
-          u'3f107997-04cc-8733-60a9-a4bb62cebffc': [u'0aec471d-2b7c-d93d-b22c-2110fd17ea2c'
-                  ],
-          u'0aec471d-2b7c-d93d-b22c-2110fd17ea2c': [u'33430f0f-ebfa-c3ec-f237-3e77efa03d0a'
-                  ],
-          u'17c9c895-5a16-7443-bb81-f34b30b21548': [],
-          },
-      u'nodeIds': [u'3f107997-04cc-8733-60a9-a4bb62cebffc',
-                   u'33430f0f-ebfa-c3ec-f237-3e77efa03d0a',
-                   u'17c9c895-5a16-7443-bb81-f34b30b21548',
-                   u'0aec471d-2b7c-d93d-b22c-2110fd17ea2c'],
-      u'id': 47,
+#       u'versions': [u'uri:oozie:workflow:0.4', u'uri:oozie:workflow:0.4.5', u'uri:oozie:workflow:0.5'],
+#       u'isDirty': False,
+#       u'movedNode': None,
+#       u'linkMapping': {
+#           u'33430f0f-ebfa-c3ec-f237-3e77efa03d0a': [],
+#           u'3f107997-04cc-8733-60a9-a4bb62cebffc': [
+#               u'0aec471d-2b7c-d93d-b22c-2110fd17ea2c'
+#           ],
+#           u'0aec471d-2b7c-d93d-b22c-2110fd17ea2c': [
+#               u'33430f0f-ebfa-c3ec-f237-3e77efa03d0a'
+#           ],
+#           u'17c9c895-5a16-7443-bb81-f34b30b21548': [],
+#       },
+#       u'nodeIds': [
+#           u'3f107997-04cc-8733-60a9-a4bb62cebffc',
+#           u'33430f0f-ebfa-c3ec-f237-3e77efa03d0a',
+#           u'17c9c895-5a16-7443-bb81-f34b30b21548',
+#           u'0aec471d-2b7c-d93d-b22c-2110fd17ea2c'
+#       ],
+#       u'id': 47,
       u'nodes': [{
           u'name': u'Start',
           u'properties': {},
           u'actionParametersFetched': False,
           u'id': u'3f107997-04cc-8733-60a9-a4bb62cebffc',
           u'type': u'start-widget',
-          u'children': [{u'to': u'0aec471d-2b7c-d93d-b22c-2110fd17ea2c'
-                        }],
+          u'children': [{u'to': u'33430f0f-ebfa-c3ec-f237-3e77efa03d0a'}],
           u'actionParameters': [],
-          }, {
+        }, {
           u'name': u'End',
           u'properties': {},
           u'actionParametersFetched': False,
@@ -3051,7 +3064,7 @@ class WorkflowBuilder():
           u'type': u'end-widget',
           u'children': [],
           u'actionParameters': [],
-          }, {
+        }, {
           u'name': u'Kill',
           u'properties': {
               u'body': u'',
@@ -3067,9 +3080,8 @@ class WorkflowBuilder():
           u'type': u'kill-widget',
           u'children': [],
           u'actionParameters': [],
-          },
-            node
-          ],
+        }
+      ],
       u'properties': {
           u'job_xml': u'',
           u'description': u'',
@@ -3092,16 +3104,29 @@ class WorkflowBuilder():
           u'parameters': parameters,
           u'properties': [],
           },
-      u'nodeNamesMapping': {
-          u'33430f0f-ebfa-c3ec-f237-3e77efa03d0a': u'End',
-          u'3f107997-04cc-8733-60a9-a4bb62cebffc': u'Start',
-          u'0aec471d-2b7c-d93d-b22c-2110fd17ea2c': u'doc-1',
-          u'17c9c895-5a16-7443-bb81-f34b30b21548': u'Kill',
-          },
-      u'uuid': u'433922e5-e616-dfe0-1cba-7fe744c9305c',
+#       u'nodeNamesMapping': {
+#           u'33430f0f-ebfa-c3ec-f237-3e77efa03d0a': u'End',
+#           u'3f107997-04cc-8733-60a9-a4bb62cebffc': u'Start',
+#           u'0aec471d-2b7c-d93d-b22c-2110fd17ea2c': u'doc-1',
+#           u'17c9c895-5a16-7443-bb81-f34b30b21548': u'Kill',
+#           },
+      u'uuid': str(uuid.uuid4()),
       }
-    })
+    }
+
+    _prev_node = data['workflow']['nodes'][0]
+
+    for node in nodes:
+      data['workflow']['nodes'].append(node)
+
+      _prev_node['children'][0]['to'] = node['id'] # We link nodes
+      _prev_node = node
 
+    data = json.dumps(data)
+    
+#     workflow_doc = _save_workflow(workflow, layout, user, fs)
+# is_managed
+    
     workflow_doc = Document2.objects.create(name=name, type='oozie-workflow2', owner=user, data=data, is_managed=managed)
     Document.objects.link(workflow_doc, owner=workflow_doc.owner, name=workflow_doc.name, description=workflow_doc.description, extra='workflow2')
 

+ 39 - 3
apps/oozie/src/oozie/tests2.py

@@ -33,8 +33,10 @@ from desktop.models import DefaultConfiguration, Document, Document2
 from oozie.conf import ENABLE_V2
 from oozie.importlib.workflows import generate_v2_graph_nodes
 from oozie.models2 import Node, Workflow, WorkflowConfiguration, find_dollar_variables, find_dollar_braced_variables, \
-    _create_graph_adjaceny_list, _get_hierarchy_from_adj_list
+    _create_graph_adjaceny_list, _get_hierarchy_from_adj_list, WorkflowBuilder
 from oozie.tests import OozieMockBase, save_temp_workflow, MockOozieApi
+from notebook.models import make_notebook
+from notebook.api import _save_notebook
 
 
 LOG = logging.getLogger(__name__)
@@ -432,8 +434,7 @@ LIMIT $limit"""))
 
       # other user can access document
       response = self.client_not_me.get(reverse('oozie:edit_workflow'), {'workflow': wf_doc.uuid})
-      assert_false('Document does not exist or you don't have the permission to access it.' in response.content,
-                   response.content)
+      assert_false('Document does not exist or you don't have the permission to access it.' in response.content, response.content)
     finally:
       wf_doc.delete()
 
@@ -955,3 +956,38 @@ class TestExternalWorkflowGraph(object):
     assert_true(len(workflow_data['workflow']['nodes']) == 4)
     assert_equal(workflow_data['layout'][0]['rows'][1]['widgets'][0]['widgetType'], 'spark-widget')
     assert_true(len(workflow_data['workflow']['nodes'][1]['children']) == 2)
+
+
+class TestModelAPI(OozieMockBase):
+
+  def setUp(self):
+    super(TestModelAPI, self).setUp()
+    self.wf = Workflow()
+
+    self.client_not_me = make_logged_in_client(username="not_perm_user", groupname="default", recreate=True,
+                                               is_superuser=False)
+    self.user_not_me = User.objects.get(username="not_perm_user")
+
+
+  def test_gen_workflow_from_document(self):
+    notebook = make_notebook(name='Browse', editor_type='hive', statement='SHOW TABLES', status='ready')
+    notebook_doc = _save_notebook(notebook, self.user)
+
+    workflow_doc = WorkflowBuilder().create_workflow(document=notebook_doc, user=self.user, managed=True)
+    
+    workflow = Workflow(document=workflow_doc, user=self.user)
+
+    _data = workflow.get_data()
+    assert_equal(len(_data['nodes']), 4)
+
+
+  def test_gen_workflow_from_documents(self):
+    notebook = make_notebook(name='Browse', editor_type='hive', statement='SHOW TABLES', status='ready')
+    notebook_doc = _save_notebook(notebook, self.user)
+
+    workflow_doc = WorkflowBuilder().create_workflow(documents=[notebook_doc, notebook_doc], user=self.user, managed=True)
+    
+    workflow = Workflow(document=workflow_doc, user=self.user)
+
+    _data = workflow.get_data()
+    assert_equal(len(_data['nodes']), 5)

+ 22 - 13
apps/oozie/src/oozie/views/editor2.py

@@ -204,18 +204,11 @@ def _import_workspace(fs, user, job):
   job.import_workspace(fs, source_workspace_dir, user)
 
 
-@check_editor_access_permission
-@check_document_modify_permission()
-def save_workflow(request):
-  response = {'status': -1}
-
-  workflow = json.loads(request.POST.get('workflow', '{}'))
-  layout = json.loads(request.POST.get('layout', '{}'))
-
+def _save_workflow(workflow, layout, user, fs=None):
   if workflow.get('id'):
     workflow_doc = Document2.objects.get(id=workflow['id'])
   else:
-    workflow_doc = Document2.objects.create(name=workflow['name'], uuid=workflow['uuid'], type='oozie-workflow2', owner=request.user, description=workflow['properties']['description'])
+    workflow_doc = Document2.objects.create(name=workflow['name'], uuid=workflow['uuid'], type='oozie-workflow2', owner=user, description=workflow['properties']['description'])
     Document.objects.link(workflow_doc, owner=workflow_doc.owner, name=workflow_doc.name, description=workflow_doc.description, extra='workflow2')
 
   # Excludes all the sub-workflow and Hive dependencies. Contains list of history and coordinator dependencies.
@@ -228,12 +221,11 @@ def save_workflow(request):
     dependency_docs = Document2.objects.filter(uuid__in=dependencies)
     workflow_doc.dependencies.add(*dependency_docs)
 
-  if workflow['properties'].get('imported'): # We save and old format workflow to the latest
+  if workflow['properties'].get('imported'): # We convert from and old workflow format (3.8 <) to the latest
     workflow['properties']['imported'] = False
-    workflow_instance = Workflow(workflow=workflow, user=request.user)
-    _import_workspace(request.fs, request.user, workflow_instance)
+    workflow_instance = Workflow(workflow=workflow, user=user)
+    _import_workspace(fs, user, workflow_instance)
     workflow['properties']['deployment_dir'] = workflow_instance.deployment_dir
-    response['url'] = reverse('oozie:edit_workflow') + '?workflow=' + str(workflow_doc.id)
 
   workflow_doc.update_data({'workflow': workflow})
   workflow_doc.update_data({'layout': layout})
@@ -243,6 +235,23 @@ def save_workflow(request):
   workflow_doc.save()
   workflow_doc1.save()
 
+
+@check_editor_access_permission
+@check_document_modify_permission()
+def save_workflow(request):
+  response = {'status': -1}
+
+  workflow = json.loads(request.POST.get('workflow', '{}'))
+  layout = json.loads(request.POST.get('layout', '{}'))
+
+  is_imported = workflow['properties'].get('imported')
+
+  workflow_doc = _save_workflow(workflow, layout, request.user)
+
+  # For old workflow import
+  if is_imported:
+    response['url'] = reverse('oozie:edit_workflow') + '?workflow=' + str(workflow_doc.id)
+
   response['status'] = 0
   response['id'] = workflow_doc.id
   response['doc_uuid'] = workflow_doc.uuid