瀏覽代碼

HUE-5304 [indexer] Port morphline job to work with task framework

Romain Rigaux 8 年之前
父節點
當前提交
8de909d

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

@@ -169,7 +169,7 @@ def importer_submit(request):
     _convert_format(source["format"], inverse=True)
     collection_name = destination["name"]
     source['columns'] = destination['columns']
-    job_handle = _index(request, source, collection_name)
+    job_handle = _index(request, source, collection_name, start_time=start_time)
   elif destination['ouputFormat'] == 'database':
     job_handle = create_database(request, source, destination, start_time)
   else:
@@ -390,7 +390,7 @@ def _create_table_from_a_file(request, source, destination, start_time=-1):
   )
 
 
-def _index(request, file_format, collection_name, query=None):
+def _index(request, file_format, collection_name, query=None, start_time=None):
   indexer = Indexer(request.user, request.fs)
 
   unique_field = indexer.get_unique_field(file_format)
@@ -425,4 +425,4 @@ def _index(request, file_format, collection_name, query=None):
 
   morphline = indexer.generate_morphline_config(collection_name, file_format, unique_field)
 
-  return indexer.run_morphline(request, collection_name, morphline, input_path, query)
+  return indexer.run_morphline(request, collection_name, morphline, input_path, query, start_time=start_time)

+ 13 - 15
desktop/libs/indexer/src/indexer/smart_indexer.py

@@ -25,12 +25,12 @@ from django.utils.translation import ugettext as _
 from mako.lookup import TemplateLookup
 
 from desktop.models import Document2
+from libzookeeper.conf import ENSEMBLE
 from notebook.connectors.base import get_api
 from notebook.models import Notebook, make_notebook
 
 from indexer.conf import CONFIG_INDEXING_TEMPLATES_PATH
 from indexer.conf import CONFIG_INDEXER_LIBS_PATH
-from indexer.conf import zkensemble
 from indexer.fields import get_field_type
 from indexer.file_format import get_file_format_instance, get_file_format_class
 from indexer.operations import get_checked_args
@@ -64,19 +64,17 @@ class Indexer(object):
 
     return hdfs_workspace_path
 
-  def run_morphline(self, request, collection_name, morphline, input_path, query=None):
+  def run_morphline(self, request, collection_name, morphline, input_path, query=None, start_time=None):
     workspace_path = self._upload_workspace(morphline)
-# 
-    task = Notebook(
-        name='Indexer job for %s' % collection_name,
-        isManaged=True
+
+    task = make_notebook(
+      name=_('Indexer job for %s') % collection_name,
+      editor_type='notebook',
+      on_success_url=reverse('search:browse', kwargs={'name': collection_name}),
+      is_task=True,
+      is_notebook=True,
+      last_executed=start_time
     )
-#     task = make_notebook(
-#       name=_('Indexer job for %s') % collection_name,
-#       editor_type='notebook',
-#       on_success_url=reverse('search:browse', kwargs={'name': collection_name}),
-#       is_task=True
-#     )
 
     if query:
       q = Notebook(document=Document2.objects.get_by_uuid(user=self.user, uuid=query))
@@ -87,7 +85,7 @@ class Indexer(object):
 
       destination = '__hue_%s' % notebook_data['uuid'][:4]
       location = '/user/%s/__hue-%s' % (request.user,  notebook_data['uuid'][:4])
-      sql, success_url = api.export_data_as_table(notebook_data, snippet, destination, is_temporary=True, location=location)
+      sql, _success_url = api.export_data_as_table(notebook_data, snippet, destination, is_temporary=True, location=location)
       input_path = '${nameNode}%s' % location
 
       task.add_hive_snippet(snippet['database'], sql)
@@ -104,7 +102,7 @@ class Indexer(object):
           u'log4j.properties',
           u'--go-live',
           u'--zk-host',
-          zkensemble(),
+          ENSEMBLE.get(),
           u'--collection',
           collection_name,
           input_path,
@@ -185,7 +183,7 @@ class Indexer(object):
       "get_kept_args": get_checked_args,
       "grok_dictionaries_location" : grok_dicts_loc if self.fs and self.fs.exists(grok_dicts_loc) else None,
       "geolite_db_location" : geolite_loc if self.fs and self.fs.exists(geolite_loc) else None,
-      "zk_host": zkensemble()
+      "zk_host": ENSEMBLE.get()
     }
 
     oozie_workspace = CONFIG_INDEXING_TEMPLATES_PATH.get()

+ 1 - 1
desktop/libs/libsolr/src/libsolr/api.py

@@ -472,7 +472,7 @@ class SolrApi(object):
       )
 
       data = self._root.post('admin/collections', params=params, contenttype='application/json')
-      if 'success' in data['responseHeader']['status'] == 0:
+      if data['responseHeader']['status'] == 0:
         response['status'] = 0
       else:
         response['message'] = "Could not remove collection: %s" % data

+ 6 - 6
desktop/libs/notebook/src/notebook/models.py

@@ -57,7 +57,7 @@ def escape_rows(rows, nulls_only=False):
 
 def make_notebook(name='Browse', description='', editor_type='hive', statement='', status='ready',
                   files=None, functions=None, settings=None, is_saved=False, database='default', snippet_properties=None, batch_submit=False,
-                  on_success_url=None, skip_historify=False, is_task=False, last_executed=-1):
+                  on_success_url=None, skip_historify=False, is_task=False, last_executed=-1, is_notebook=False):
   '''
   skip_historify: do not add the task to the query history. e.g. SQL Dashboard
   isManaged: true when being a managed by Hue operation (include_managed=True in document), e.g. exporting query result, dropping some tables
@@ -99,7 +99,7 @@ def make_notebook(name='Browse', description='', editor_type='hive', statement='
       }
     ],
     'selectedSnippet': editor_type,
-    'type': 'query-%s' % editor_type,
+    'type': 'notebook' if is_notebook else 'query-%s' % editor_type,
     'showHistory': True,
     'isSaved': is_saved,
     'onSuccessUrl': on_success_url,
@@ -124,7 +124,7 @@ def make_notebook(name='Browse', description='', editor_type='hive', statement='
          'result': {'handle':{}},
          'variables': []
       }
-    ]
+    ] if not is_notebook else []
   }
 
   if snippet_properties:
@@ -162,9 +162,9 @@ def make_notebook2(name='Browse', description='', is_saved=False, snippets=None)
     'description': description,
     'sessions': [
       {
-         'type': _snippet['type'],
-         'properties': HS2Api.get_properties(snippet['type']),
-         'id': None
+        'type': _snippet['type'],
+        'properties': HS2Api.get_properties(snippet['type']),
+        'id': None
       } for _snippet in _snippets # Non unique types currently
     ],
     'selectedSnippet': _snippets[0]['type'],