Browse Source

HUE-8330 [oozie] Submit a mocked scheduled job to de-cluster

Romain Rigaux 7 years ago
parent
commit
166e4722c2

+ 2 - 1
apps/oozie/src/oozie/static/oozie/js/workflow-editor.ko.js

@@ -1241,7 +1241,8 @@ var WorkflowEditorViewModel = function (layout_json, workflow_json, credentials_
   self.showSubmitPopup = function () {
   self.showSubmitPopup = function () {
     $(".jHueNotify").remove();
     $(".jHueNotify").remove();
     $.get("/oozie/editor/workflow/submit/" + self.workflow.id(), {
     $.get("/oozie/editor/workflow/submit/" + self.workflow.id(), {
-      format: IS_HUE_4 ? 'json' : 'html'
+      format: IS_HUE_4 ? 'json' : 'html',
+      cluster: self.compute() ? ko.mapping.toJSON(self.compute()) : '{}'
     }, function (data) {
     }, function (data) {
       $(document).trigger("showSubmitPopup", data);
       $(document).trigger("showSubmitPopup", data);
     }).fail(function (xhr, textStatus, errorThrown) {
     }).fail(function (xhr, textStatus, errorThrown) {

+ 3 - 0
apps/oozie/src/oozie/templates/editor2/submit_job_popup.mako

@@ -92,6 +92,9 @@
         % else:
         % else:
           ${_('Email not set in ')}<a href="/useradmin/users/edit/${user.username}#step2"> ${_('profile.')} </a>
           ${_('Email not set in ')}<a href="/useradmin/users/edit/${user.username}#step2"> ${_('profile.')} </a>
         % endif
         % endif
+        % if cluster_json:
+          <input type="hidden" name="cluster" value="${ cluster_json }"></input>
+        % endif
         </label>
         </label>
         %endif       
         %endif       
       % if return_json:
       % if return_json:

+ 11 - 1
apps/oozie/src/oozie/views/editor2.py

@@ -46,6 +46,7 @@ from oozie.models2 import Node, Workflow, Coordinator, Bundle, NODES, WORKFLOW_N
   _import_workspace, _save_workflow
   _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
+from notebook.connectors.dataeng import DataEngApi
 
 
 
 
 LOG = logging.getLogger(__name__)
 LOG = logging.getLogger(__name__)
@@ -390,6 +391,7 @@ def _submit_workflow_helper(request, workflow, submit_action):
   ParametersFormSet = formset_factory(ParameterForm, extra=0)
   ParametersFormSet = formset_factory(ParameterForm, extra=0)
 
 
   if request.method == 'POST':
   if request.method == 'POST':
+    cluster = json.loads(request.POST.get('cluster', '{}'))
     params_form = ParametersFormSet(request.POST)
     params_form = ParametersFormSet(request.POST)
 
 
     if params_form.is_valid():
     if params_form.is_valid():
@@ -399,6 +401,12 @@ def _submit_workflow_helper(request, workflow, submit_action):
       if '/submit_single_action/' in submit_action:
       if '/submit_single_action/' in submit_action:
         mapping['submit_single_action'] = True
         mapping['submit_single_action'] = True
 
 
+      if cluster.get('type') == 'altus-de':
+        notebook = {}
+        snippet = {'statement': 'SELECT 1'}
+        # TODO: open in Job Browser, Jobs, compute context
+        print DataEngApi(user=request.user, request=request, cluster_name=cluster.get('name')).execute(notebook, snippet)
+
       try:
       try:
         job_id = _submit_workflow(request.user, request.fs, request.jt, workflow, mapping)
         job_id = _submit_workflow(request.user, request.fs, request.jt, workflow, mapping)
       except Exception, e:
       except Exception, e:
@@ -412,6 +420,7 @@ def _submit_workflow_helper(request, workflow, submit_action):
     else:
     else:
       request.error(_('Invalid submission form: %s' % params_form.errors))
       request.error(_('Invalid submission form: %s' % params_form.errors))
   else:
   else:
+    cluster_json = request.GET.get('cluster', '{}')
     parameters = workflow and workflow.find_all_parameters() or []
     parameters = workflow and workflow.find_all_parameters() or []
     initial_params = ParameterForm.get_initial_params(dict([(param['name'], param['value']) for param in parameters]))
     initial_params = ParameterForm.get_initial_params(dict([(param['name'], param['value']) for param in parameters]))
     params_form = ParametersFormSet(initial=initial_params)
     params_form = ParametersFormSet(initial=initial_params)
@@ -424,7 +433,8 @@ def _submit_workflow_helper(request, workflow, submit_action):
                      'show_dryrun': True,
                      'show_dryrun': True,
                      'email_id': request.user.email,
                      'email_id': request.user.email,
                      'is_oozie_mail_enabled': _is_oozie_mail_enabled(request.user),
                      'is_oozie_mail_enabled': _is_oozie_mail_enabled(request.user),
-                     'return_json': request.GET.get('format') == 'json'
+                     'return_json': request.GET.get('format') == 'json',
+                     'cluster_json': cluster_json
                    }, force_template=True).content
                    }, force_template=True).content
     return JsonResponse(popup, safe=False)
     return JsonResponse(popup, safe=False)
 
 

+ 7 - 2
desktop/core/src/desktop/api2.py

@@ -114,14 +114,17 @@ def get_context_computes(request, interface):
 
 
   clusters = get_clusters(request.user).values()
   clusters = get_clusters(request.user).values()
 
 
-  if interface == 'hive':
+  if interface == 'hive' or interface == 'oozie':
     computes.extend([{
     computes.extend([{
         'id': cluster['id'],
         'id': cluster['id'],
         'name': cluster['name'],
         'name': cluster['name'],
-        'namespace': cluster['id'] # Dummy
+        'namespace': cluster['id'], # Dummy
+        'interface': interface,
+        'type': 'default'
       } for cluster in clusters
       } for cluster in clusters
     ])
     ])
 
 
+  if interface == 'hive':
     if [cluster for cluster in clusters if cluster['type'] == 'altus']:
     if [cluster for cluster in clusters if cluster['type'] == 'altus']:
       computes.extend([{
       computes.extend([{
           'id': cluster.get('crn', 'None'),
           'id': cluster.get('crn', 'None'),
@@ -131,6 +134,7 @@ def get_context_computes(request, interface):
           # environmentType
           # environmentType
           # secured
           # secured
           # cdhVersion
           # cdhVersion
+          'type': 'altus-adb'
         } for cluster in AnalyticDbApi(request.user).list_clusters()['clusters']]
         } for cluster in AnalyticDbApi(request.user).list_clusters()['clusters']]
       )
       )
 
 
@@ -142,6 +146,7 @@ def get_context_computes(request, interface):
           'status': cluster.get('status'),
           'status': cluster.get('status'),
           'environmentType': cluster.get('environmentType'),
           'environmentType': cluster.get('environmentType'),
           'serviceType': cluster.get('serviceType'),
           'serviceType': cluster.get('serviceType'),
+          'type': 'altus-de'
         } for cluster in DataEngApi(request.user).list_clusters()['clusters']]
         } for cluster in DataEngApi(request.user).list_clusters()['clusters']]
       )
       )
 
 

+ 2 - 1
desktop/libs/notebook/src/notebook/connectors/base.py

@@ -281,6 +281,7 @@ def get_api(request, snippet):
 
 
   # Multi cluster
   # Multi cluster
   cluster = json.loads(request.POST.get('cluster', '""'))
   cluster = json.loads(request.POST.get('cluster', '""'))
+  print cluster
 
 
   if interface == 'hiveserver2':
   if interface == 'hiveserver2':
     from notebook.connectors.hiveserver2 import HS2Api
     from notebook.connectors.hiveserver2 import HS2Api
@@ -301,7 +302,7 @@ def get_api(request, snippet):
     return RdbmsApi(request.user, interpreter=snippet['type'])
     return RdbmsApi(request.user, interpreter=snippet['type'])
   elif interface == 'dataeng':
   elif interface == 'dataeng':
     from notebook.connectors.dataeng import DataEngApi
     from notebook.connectors.dataeng import DataEngApi
-    return DataEngApi(user=request.user, request=request, cluster_name=cluster.get_interface())
+    return DataEngApi(user=request.user, request=request, cluster_name=cluster.get('name'))
   elif interface == 'jdbc' or interface == 'teradata':
   elif interface == 'jdbc' or interface == 'teradata':
     from notebook.connectors.jdbc import JdbcApi
     from notebook.connectors.jdbc import JdbcApi
     return JdbcApi(request.user, interpreter=interpreter)
     return JdbcApi(request.user, interpreter=interpreter)