소스 검색

HUE-4526 [oozie] Refactor to have only one way to save a workflow

Romain Rigaux 9 년 전
부모
커밋
ac53520

+ 59 - 43
apps/oozie/src/oozie/models2.py

@@ -26,6 +26,7 @@ from dateutil.parser import parse
 from string import Template
 from string import Template
 
 
 from django.core.urlresolvers import reverse
 from django.core.urlresolvers import reverse
+from django.db.models import Q
 from django.utils.encoding import force_unicode
 from django.utils.encoding import force_unicode
 from django.utils.translation import ugettext as _
 from django.utils.translation import ugettext as _
 from django.contrib.auth.models import User
 from django.contrib.auth.models import User
@@ -48,7 +49,6 @@ from notebook.models import Notebook
 from oozie.conf import REMOTE_SAMPLE_DIR
 from oozie.conf import REMOTE_SAMPLE_DIR
 from oozie.utils import utc_datetime_format, UTC_TIME_FORMAT, convert_to_server_timezone
 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.importlib.workflows import generate_v2_graph_nodes, MalformedWfDefException, InvalidTagWithNamespaceException
-# from oozie.views.editor2 import _save_workflow
 
 
 
 
 LOG = logging.getLogger(__name__)
 LOG = logging.getLogger(__name__)
@@ -433,7 +433,7 @@ class Workflow(Job):
     data = self.get_data()
     data = self.get_data()
     nodes = [node for node in self.nodes if node.name != 'End'] + [node for node in self.nodes if
     nodes = [node for node in self.nodes if node.name != 'End'] + [node for node in self.nodes if
                                                                    node.name == 'End']  # End at the end
                                                                    node.name == 'End']  # End at the end
-    node_mapping = dict([(node.id, node) for node in nodes]) 
+    node_mapping = dict([(node.id, node) for node in nodes])
     sub_wfs_ids = [node.data['properties']['workflow'] for node in nodes if node.data['type'] == 'subworkflow']
     sub_wfs_ids = [node.data['properties']['workflow'] for node in nodes if node.data['type'] == 'subworkflow']
     workflow_mapping = dict(
     workflow_mapping = dict(
       [(workflow.uuid, Workflow(document=workflow, user=self.user)) for workflow in Document2.objects.filter(uuid__in=sub_wfs_ids)])
       [(workflow.uuid, Workflow(document=workflow, user=self.user)) for workflow in Document2.objects.filter(uuid__in=sub_wfs_ids)])
@@ -2915,6 +2915,48 @@ class History(object):
       pass
       pass
 
 
 
 
+def _import_workspace(fs, user, job):
+  source_workspace_dir = job.deployment_dir
+
+  job.set_workspace(user)
+  job.check_workspace(fs, user)
+  job.import_workspace(fs, source_workspace_dir, user)
+
+
+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=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.
+  workflow_doc.dependencies = workflow_doc.dependencies.exclude(Q(is_history=False) & Q(type__in=['oozie-workflow2', 'query-hive', 'query-java']))
+
+  dependencies = \
+      [node['properties']['workflow'] for node in workflow['nodes'] if node['type'] == 'subworkflow-widget'] + \
+      [node['properties']['uuid'] for node in workflow['nodes'] if 'document-widget' in node['type']]
+  if dependencies:
+    dependency_docs = Document2.objects.filter(uuid__in=dependencies)
+    workflow_doc.dependencies.add(*dependency_docs)
+
+  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=user)
+    _import_workspace(fs, user, workflow_instance)
+    workflow['properties']['deployment_dir'] = workflow_instance.deployment_dir
+
+  workflow_doc.update_data({'workflow': workflow})
+  workflow_doc.update_data({'layout': layout})
+  workflow_doc1 = workflow_doc.doc.get()
+  workflow_doc.name = workflow_doc1.name = workflow['name']
+  workflow_doc.description = workflow_doc1.description = workflow['properties']['description']
+  workflow_doc.save()
+  workflow_doc1.save()
+
+  return workflow_doc
+
+
 class WorkflowBuilder():
 class WorkflowBuilder():
   """
   """
   Building a workflow that has saved Documents for nodes (e.g Saved Hive query, saved Pig script...).
   Building a workflow that has saved Documents for nodes (e.g Saved Hive query, saved Pig script...).
@@ -2928,16 +2970,16 @@ class WorkflowBuilder():
     if name is None:
     if name is None:
       name = _('Schedule of ') + ','.join([document.name or document.type for document in documents])
       name = _('Schedule of ') + ','.join([document.name or document.type for document in documents])
 
 
-    for document in documents:     
+    for document in documents:
       if document.type == 'query-java':
       if document.type == 'query-java':
         node = self.get_java_document_node(document, name)
         node = self.get_java_document_node(document, name)
       else:
       else:
         node = self.get_hive_document_node(document, name, user)
         node = self.get_hive_document_node(document, name, user)
-  
+
       nodes.append(node)
       nodes.append(node)
 
 
     workflow_doc = self.get_workflow(nodes, name, document.uuid, user, managed=managed)
     workflow_doc = self.get_workflow(nodes, name, document.uuid, user, managed=managed)
-    
+
     for document in documents:
     for document in documents:
       workflow_doc.dependencies.add(document)
       workflow_doc.dependencies.add(document)
 
 
@@ -2956,8 +2998,7 @@ class WorkflowBuilder():
     return {
     return {
         u'name': u'doc-hive-%s' % document.uuid[:4],
         u'name': u'doc-hive-%s' % document.uuid[:4],
         u'id': str(uuid.uuid4()),
         u'id': str(uuid.uuid4()),
-        u'type': u'hive-document-widget',        
-        u'actionParametersUI': [],
+        u'type': u'hive-document-widget',
         u'properties': {
         u'properties': {
             u'files': [],
             u'files': [],
             u'job_xml': u'',
             u'job_xml': u'',
@@ -2982,8 +3023,7 @@ class WorkflowBuilder():
             u'credentials': credentials,
             u'credentials': credentials,
             u'password': u'',
             u'password': u'',
             u'jdbc_url': u'',
             u'jdbc_url': u'',
-            },
-        u'actionParametersFetched': False,
+        },
         u'children': [
         u'children': [
             {u'to': u'33430f0f-ebfa-c3ec-f237-3e77efa03d0a'},
             {u'to': u'33430f0f-ebfa-c3ec-f237-3e77efa03d0a'},
             {u'error': u'17c9c895-5a16-7443-bb81-f34b30b21548'
             {u'error': u'17c9c895-5a16-7443-bb81-f34b30b21548'
@@ -3028,26 +3068,6 @@ class WorkflowBuilder():
     data = {
     data = {
       'workflow': {
       'workflow': {
       u'name': name,
       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'nodes': [{
       u'nodes': [{
           u'name': u'Start',
           u'name': u'Start',
           u'properties': {},
           u'properties': {},
@@ -3071,8 +3091,7 @@ class WorkflowBuilder():
               u'cc': u'',
               u'cc': u'',
               u'to': u'',
               u'to': u'',
               u'enableMail': False,
               u'enableMail': False,
-              u'message': u'Action failed, error message[${wf:errorMessage(wf:lastErrorNode())}]'
-                  ,
+              u'message': u'Action failed, error message[${wf:errorMessage(wf:lastErrorNode())}]',
               u'subject': u'',
               u'subject': u'',
               },
               },
           u'actionParametersFetched': False,
           u'actionParametersFetched': False,
@@ -3104,12 +3123,6 @@ class WorkflowBuilder():
           u'parameters': parameters,
           u'parameters': parameters,
           u'properties': [],
           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': str(uuid.uuid4()),
       u'uuid': str(uuid.uuid4()),
       }
       }
     }
     }
@@ -3122,12 +3135,15 @@ class WorkflowBuilder():
       _prev_node['children'][0]['to'] = node['id'] # We link nodes
       _prev_node['children'][0]['to'] = node['id'] # We link nodes
       _prev_node = node
       _prev_node = node
 
 
-    data = json.dumps(data)
-    
-#     workflow_doc = _save_workflow(workflow, layout, user, fs)
+#     data = json.dumps(data)
+
+
+    workflow_doc = _save_workflow(data['workflow'], {}, user,) # # from oozie.views.editor2 import _save_workflow
+    workflow_doc.is_managed = managed
+    workflow_doc.save()
 # is_managed
 # 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')
+
+#     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')
 
 
     return workflow_doc
     return workflow_doc

+ 4 - 1
apps/oozie/src/oozie/tests2.py

@@ -985,7 +985,10 @@ class TestModelAPI(OozieMockBase):
     notebook = make_notebook(name='Browse', editor_type='hive', statement='SHOW TABLES', status='ready')
     notebook = make_notebook(name='Browse', editor_type='hive', statement='SHOW TABLES', status='ready')
     notebook_doc = _save_notebook(notebook, self.user)
     notebook_doc = _save_notebook(notebook, self.user)
 
 
-    workflow_doc = WorkflowBuilder().create_workflow(documents=[notebook_doc, notebook_doc], user=self.user, managed=True)
+    notebook2 = make_notebook(name='Browse', editor_type='hive', statement='SHOW TABLES', status='ready')
+    notebook2_doc = _save_notebook(notebook2, self.user)
+
+    workflow_doc = WorkflowBuilder().create_workflow(documents=[notebook_doc, notebook_doc2], user=self.user, managed=True)
     
     
     workflow = Workflow(document=workflow_doc, user=self.user)
     workflow = Workflow(document=workflow_doc, user=self.user)
 
 

+ 2 - 42
apps/oozie/src/oozie/views/editor2.py

@@ -19,7 +19,6 @@ import json
 import logging
 import logging
 
 
 from django.core.urlresolvers import reverse
 from django.core.urlresolvers import reverse
-from django.db.models import Q
 from django.forms.formsets import formset_factory
 from django.forms.formsets import formset_factory
 from django.shortcuts import redirect
 from django.shortcuts import redirect
 from django.utils.translation import ugettext as _
 from django.utils.translation import ugettext as _
@@ -43,7 +42,8 @@ from oozie.decorators import check_document_access_permission, check_document_mo
 from oozie.forms import ParameterForm
 from oozie.forms import ParameterForm
 from oozie.models import Workflow as OldWorklow, Coordinator as OldCoordinator, Bundle as OldBundle, Job
 from oozie.models import Workflow as OldWorklow, Coordinator as OldCoordinator, Bundle as OldBundle, Job
 from oozie.models2 import Node, Workflow, Coordinator, Bundle, NODES, WORKFLOW_NODE_PROPERTIES, import_workflow_from_hue_3_7,\
 from oozie.models2 import Node, Workflow, Coordinator, Bundle, NODES, WORKFLOW_NODE_PROPERTIES, import_workflow_from_hue_3_7,\
-    find_dollar_variables, find_dollar_braced_variables, WorkflowBuilder
+    find_dollar_variables, find_dollar_braced_variables, WorkflowBuilder,\
+  _import_workspace, _save_workflow
 from oozie.utils import convert_to_server_timezone
 from oozie.utils import convert_to_server_timezone
 from oozie.views.editor import edit_workflow as old_edit_workflow, edit_coordinator as old_edit_coordinator, edit_bundle as old_edit_bundle
 from oozie.views.editor import edit_workflow as old_edit_workflow, edit_coordinator as old_edit_coordinator, edit_bundle as old_edit_bundle
 
 
@@ -196,46 +196,6 @@ def copy_workflow(request):
   return JsonResponse(response)
   return JsonResponse(response)
 
 
 
 
-def _import_workspace(fs, user, job):
-  source_workspace_dir = job.deployment_dir
-
-  job.set_workspace(user)
-  job.check_workspace(fs, user)
-  job.import_workspace(fs, source_workspace_dir, user)
-
-
-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=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.
-  workflow_doc.dependencies = workflow_doc.dependencies.exclude(Q(is_history=False) & Q(type__in=['oozie-workflow2', 'query-hive', 'query-java']))
-
-  dependencies = \
-      [node['properties']['workflow'] for node in workflow['nodes'] if node['type'] == 'subworkflow-widget'] + \
-      [node['properties']['uuid'] for node in workflow['nodes'] if 'document-widget' in node['type']]
-  if dependencies:
-    dependency_docs = Document2.objects.filter(uuid__in=dependencies)
-    workflow_doc.dependencies.add(*dependency_docs)
-
-  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=user)
-    _import_workspace(fs, user, workflow_instance)
-    workflow['properties']['deployment_dir'] = workflow_instance.deployment_dir
-
-  workflow_doc.update_data({'workflow': workflow})
-  workflow_doc.update_data({'layout': layout})
-  workflow_doc1 = workflow_doc.doc.get()
-  workflow_doc.name = workflow_doc1.name = workflow['name']
-  workflow_doc.description = workflow_doc1.description = workflow['properties']['description']
-  workflow_doc.save()
-  workflow_doc1.save()
-
-
 @check_editor_access_permission
 @check_editor_access_permission
 @check_document_modify_permission()
 @check_document_modify_permission()
 def save_workflow(request):
 def save_workflow(request):

+ 1 - 1
desktop/libs/notebook/src/notebook/connectors/oozie_batch.py

@@ -63,7 +63,7 @@ class OozieApi(Api):
     notebook_doc = Document2.objects.get_by_uuid(user=self.user, uuid=notebook['uuid'], perm_type='read')
     notebook_doc = Document2.objects.get_by_uuid(user=self.user, uuid=notebook['uuid'], perm_type='read')
 
 
     # Create a managed workflow from the notebook doc
     # Create a managed workflow from the notebook doc
-    workflow_doc = WorkflowBuilder().create_workflow(document=notebook_doc, user=self.user, managed=True, name=_("Batch job for %s") % notebook_doc.name or notebook_doc.type)
+    workflow_doc = WorkflowBuilder().create_workflow(document=notebook_doc, user=self.user, managed=True, name=_("Batch job for %s") % (notebook_doc.name or notebook_doc.type))
     workflow = Workflow(document=workflow_doc, user=self.user)
     workflow = Workflow(document=workflow_doc, user=self.user)
 
 
     # Submit workflow
     # Submit workflow