ソースを参照

HUE-4530 [notebook] Submit indexer as a batch job

Romain Rigaux 9 年 前
コミット
4d5dd79fd3

+ 2 - 3
desktop/libs/indexer/src/indexer/api3.py

@@ -141,6 +141,5 @@ def index_file(request):
   else:
   else:
     input_path = file_format["path"]
     input_path = file_format["path"]
 
 
-  job_id = indexer.run_morphline(collection_name, morphline, input_path) #TODO if query generate insert
-
-  return JsonResponse({"jobId": job_id})
+  job_handle = indexer.run_morphline(request, collection_name, morphline, input_path) #TODO if query generate insert
+  return JsonResponse(job_handle)

+ 29 - 31
desktop/libs/indexer/src/indexer/smart_indexer.py

@@ -23,10 +23,9 @@ from mako.lookup import TemplateLookup
 from mako.template import Template
 from mako.template import Template
 
 
 from collections import deque
 from collections import deque
-from notebook.api import _save_notebook
-from notebook.models import make_notebook
-from oozie.views.editor2 import _submit_workflow
-from oozie.models2 import Job, WorkflowBuilder, Workflow
+from notebook.api import _save_notebook, _execute_notebook
+from notebook.models import make_notebook, make_notebook2
+from oozie.models2 import Job
 
 
 from indexer.fields import get_field_type
 from indexer.fields import get_field_type
 from indexer.operations import get_checked_args
 from indexer.operations import get_checked_args
@@ -62,42 +61,41 @@ class Indexer(object):
 
 
     return hdfs_workspace_path
     return hdfs_workspace_path
 
 
-  def run_morphline(self, collection_name, morphline, input_path):
+  def run_morphline(self, request, collection_name, morphline, input_path):
     workspace_path = self._upload_workspace(morphline)
     workspace_path = self._upload_workspace(morphline)
 
 
     snippet_properties =  {
     snippet_properties =  {
-      u'files': [
-          {u'path': u'%s/log4j.properties' % workspace_path, u'type': u'file'},
-          {u'path': u'%s/morphline.conf' % workspace_path, u'type': u'file'}
-      ],
-      u'class': u'org.apache.solr.hadoop.MapReduceIndexerTool',
-      u'app_jar': CONFIG_INDEXER_LIBS_PATH.get(),
-      u'arguments': [
-          u'--morphline-file',
-          u'morphline.conf',
-          u'--output-dir',
-          u'${nameNode}/user/%s/indexer' % self.username,
-          u'--log4j',
-          u'log4j.properties',
-          u'--go-live',
-          u'--zk-host',
-          zkensemble(),
-          u'--collection',
-          collection_name,
-          input_path,
-      ],
-      u'archives': [],
+       u'files': [
+           {u'path': u'%s/log4j.properties' % workspace_path, u'type': u'file'},
+           {u'path': u'%s/morphline.conf' % workspace_path, u'type': u'file'}
+       ],
+       u'class': u'org.apache.solr.hadoop.MapReduceIndexerTool',
+       u'app_jar': CONFIG_INDEXER_LIBS_PATH.get(),
+       u'arguments': [
+           u'--morphline-file',
+           u'morphline.conf',
+           u'--output-dir',
+           u'${nameNode}/user/%s/indexer' % self.username,
+           u'--log4j',
+           u'log4j.properties',
+           u'--go-live',
+           u'--zk-host',
+           zkensemble(),
+           u'--collection',
+           collection_name,
+           input_path,
+       ],
+       u'archives': [],
     }
     }
 
 
-    notebook = make_notebook(name='Indexer', editor_type='java', snippet_properties=snippet_properties).get_data()
+    notebook = make_notebook(name='Indexer', editor_type='java', snippet_properties=snippet_properties, status='running').get_data()
     notebook_doc, created = _save_notebook(notebook, self.user)
     notebook_doc, created = _save_notebook(notebook, self.user)
 
 
-    workflow_doc = WorkflowBuilder().create_workflow(document=notebook_doc, user=self.user, managed=True, name=_("Batch job for %s") % notebook_doc.name)
-    workflow = Workflow(document=workflow_doc, user=self.user)
+    snippet = {'wasBatchExecuted': True, 'id': notebook['snippets'][0]['id'], 'statement': ''}
 
 
-    job_id = _submit_workflow(user=self.user, fs=self.fs, jt=self.jt, workflow=workflow, mapping=None)
+    job_handle = _execute_notebook(request, notebook, snippet)
 
 
-    return job_id
+    return job_handle
 
 
   def guess_format(self, data):
   def guess_format(self, data):
     """
     """

+ 11 - 9
desktop/libs/indexer/src/indexer/templates/indexer.mako

@@ -301,7 +301,6 @@ ${ assist.assistPanel() }
       </div>
       </div>
     <!-- /ko -->
     <!-- /ko -->
 
 
-
     <!-- ko if: previousStepVisible -->
     <!-- ko if: previousStepVisible -->
       <a class="btn" data-bind="click: previousStep">${ _('Previous') }</a>
       <a class="btn" data-bind="click: previousStep">${ _('Previous') }</a>
     <!-- /ko -->
     <!-- /ko -->
@@ -318,14 +317,15 @@ ${ assist.assistPanel() }
 
 
     <div data-bind="visible: createWizard.jobId">
     <div data-bind="visible: createWizard.jobId">
       <a href="javascript:void(0)" class="btn btn-success" data-bind="attr: {href: '/oozie/list_oozie_workflow/' + createWizard.jobId() }" target="_blank" title="${ _('Open') }">
       <a href="javascript:void(0)" class="btn btn-success" data-bind="attr: {href: '/oozie/list_oozie_workflow/' + createWizard.jobId() }" target="_blank" title="${ _('Open') }">
-        ${_('View Indexing Status')}
+        ${_('Oozie Status')}
+      </a>
+      <a href="javascript:void(0)" class="btn btn-success" data-bind="attr: {href: '${ url('notebook:editor') }?editor=' + createWizard.editorId() }" target="_blank" title="${ _('Open') }">
+        ${_('View indexing status')}
       </a>
       </a>
 
 
       ${ _('View collection') } <a href="javascript:void(0)" data-bind="attr: {href: '${ url("indexer:collections") }' +'#edit/' + createWizard.fileFormat().name() }, text: createWizard.fileFormat().name" target="_blank"></a>
       ${ _('View collection') } <a href="javascript:void(0)" data-bind="attr: {href: '${ url("indexer:collections") }' +'#edit/' + createWizard.fileFormat().name() }, text: createWizard.fileFormat().name" target="_blank"></a>
     </div>
     </div>
-
   </div>
   </div>
-
 </script>
 </script>
 
 
 <script type="text/html" id="format-settings">
 <script type="text/html" id="format-settings">
@@ -336,7 +336,7 @@ ${ assist.assistPanel() }
 
 
 <script type="text/html" id="field-template">
 <script type="text/html" id="field-template">
   <label>${ _('Name') }
   <label>${ _('Name') }
-    <input type="text" class="input-small" placeholder="${ _('Field name') }" data-bind="value: name">
+    <input type="text" class="input-large" placeholder="${ _('Field name') }" data-bind="value: name">
   </label>
   </label>
   <label>${ _('Type') }
   <label>${ _('Type') }
     <select data-bind="options: $root.createWizard.fieldTypes, value: type"></select>
     <select data-bind="options: $root.createWizard.fieldTypes, value: type"></select>
@@ -679,7 +679,8 @@ ${ assist.assistPanel() }
       self.fileFormat = ko.observable(new IndexerFormat(vm));
       self.fileFormat = ko.observable(new IndexerFormat(vm));
       self.sample = ko.observableArray();
       self.sample = ko.observableArray();
 
 
-      self.jobId = ko.observable(null);
+      self.jobId = ko.observable();
+      self.editorId = ko.observable();
 
 
       self.indexingStarted = ko.observable(false);
       self.indexingStarted = ko.observable(false);
 
 
@@ -756,20 +757,21 @@ ${ assist.assistPanel() }
         if (!self.readyToIndex()) return;
         if (!self.readyToIndex()) return;
 
 
         self.indexingStarted(true);
         self.indexingStarted(true);
-
         viewModel.isLoading(true);
         viewModel.isLoading(true);
-
         self.isIndexing(true);
         self.isIndexing(true);
 
 
         $.post("${ url('indexer:index_file') }", {
         $.post("${ url('indexer:index_file') }", {
           "fileFormat": ko.mapping.toJSON(self.fileFormat)
           "fileFormat": ko.mapping.toJSON(self.fileFormat)
         }, function (resp) {
         }, function (resp) {
           self.showCreate(true);
           self.showCreate(true);
-          self.jobId(resp.jobId);
+          self.jobId(resp.handle.id);
+          self.editorId(resp.history_id);
           viewModel.isLoading(false);
           viewModel.isLoading(false);
         }).fail(function (xhr, textStatus, errorThrown) {
         }).fail(function (xhr, textStatus, errorThrown) {
           $(document).trigger("error", xhr.responseText);
           $(document).trigger("error", xhr.responseText);
           viewModel.isLoading(false);
           viewModel.isLoading(false);
+          self.indexingStarted(false);
+          self.isIndexing(false);
         });
         });
       }
       }
 
 

+ 12 - 6
desktop/libs/notebook/src/notebook/api.py

@@ -96,16 +96,11 @@ def close_session(request):
   return JsonResponse(response)
   return JsonResponse(response)
 
 
 
 
-@require_POST
-@check_document_access_permission()
-@api_error_handler
-def execute(request):
+def _execute_notebook(request, notebook, snippet):
   response = {'status': -1}
   response = {'status': -1}
   result = None
   result = None
   history = None
   history = None
 
 
-  notebook = json.loads(request.POST.get('notebook', '{}'))
-  snippet = json.loads(request.POST.get('snippet', '{}'))
   is_query = notebook['type'].startswith('query-') or snippet['type'] == 'java'
   is_query = notebook['type'].startswith('query-') or snippet['type'] == 'java'
 
 
   try:
   try:
@@ -153,6 +148,17 @@ def execute(request):
 
 
   response['status'] = 0
   response['status'] = 0
 
 
+  return response
+
+@require_POST
+@check_document_access_permission()
+@api_error_handler
+def execute(request):
+  notebook = json.loads(request.POST.get('notebook', '{}'))
+  snippet = json.loads(request.POST.get('snippet', '{}'))
+
+  response = _execute_notebook(request, notebook, snippet)
+
   return JsonResponse(response)
   return JsonResponse(response)
 
 
 
 

+ 4 - 0
desktop/libs/notebook/src/notebook/models.py

@@ -121,6 +121,10 @@ def make_notebook(name='Browse', description='', editor_type='hive', statement='
   return editor
   return editor
 
 
 
 
+def make_notebook2(name='Browse', description='', is_saved=False, snippets=None):
+  pass
+
+
 def import_saved_beeswax_query(bquery):
 def import_saved_beeswax_query(bquery):
   design = bquery.get_design()
   design = bquery.get_design()