Explorar o código

[notebook] Display output of Pig scripts as result

Very tricky to get the logs correctly because of YARN behavior.
Remove jt, fs from API and just reply with request.

Not perfect yet:
 - links are escaped, we should allow them
 - output logs are appended and grow
 - little lag after execution is ready and displaying the rows
 - maybe we should not display the grid but the logs only, like Jar snippet
Romain Rigaux %!s(int64=9) %!d(string=hai) anos
pai
achega
bc0167d

+ 4 - 2
apps/jobbrowser/src/jobbrowser/views.py

@@ -74,7 +74,9 @@ def check_job_permission(view_func):
     except JobExpired, e:
       raise PopupException(_('Job %s has expired.') % jobid, detail=_('Cannot be found on the History Server.'))
     except Exception, e:
-      raise PopupException(_('Could not find job %s.') % jobid, detail=e)
+      msg = 'Could not find job %s.'
+      LOGGER.exception(msg % jobid)
+      raise PopupException(_(msg) % jobid, detail=e)
 
     if not SHARE_JOBS.get() and not request.user.is_superuser \
         and job.user != request.user.username and not can_view_job(request.user.username, job):
@@ -308,7 +310,6 @@ def job_attempt_logs_json(request, job, attempt_index=0, name='syslog', offset=0
   return JsonResponse(response)
 
 
-
 @check_job_permission
 def job_single_logs(request, job):
   """
@@ -339,6 +340,7 @@ def job_single_logs(request, job):
 
   return single_task_attempt_logs(request, **{'job': job.jobId, 'taskid': task.taskId, 'attemptid': task.taskAttemptIds[-1]})
 
+
 @check_job_permission
 def tasks(request, job):
   """

+ 11 - 11
desktop/libs/notebook/src/notebook/api.py

@@ -52,7 +52,7 @@ def create_session(request):
     if any(old_session) and 'properties' in old_session[0]:
       properties = old_session[0]['properties']
 
-  response['session'] = get_api(request.user, session, request.fs, request.jt).create_session(lang=session['type'], properties=properties)
+  response['session'] = get_api(request, session).create_session(lang=session['type'], properties=properties)
   response['status'] = 0
 
   return JsonResponse(response)
@@ -66,7 +66,7 @@ def close_session(request):
 
   session = json.loads(request.POST.get('session', '{}'))
 
-  response['session'] = get_api(request.user, {'type': session['type']}, request.fs, request.jt).close_session(session=session)
+  response['session'] = get_api(request, {'type': session['type']}).close_session(session=session)
   response['status'] = 0
 
   return JsonResponse(response)
@@ -81,7 +81,7 @@ def execute(request):
   notebook = json.loads(request.POST.get('notebook', '{}'))
   snippet = json.loads(request.POST.get('snippet', '{}'))
 
-  response['handle'] = get_api(request.user, snippet, request.fs, request.jt).execute(notebook, snippet)
+  response['handle'] = get_api(request, snippet).execute(notebook, snippet)
 
   # Materialize and HTML escape results
   if response['handle'].get('sync') and response['handle']['result'].get('data'):
@@ -101,7 +101,7 @@ def check_status(request):
   notebook = json.loads(request.POST.get('notebook', '{}'))
   snippet = json.loads(request.POST.get('snippet', '{}'))
 
-  response['query_status'] = get_api(request.user, snippet, request.fs, request.jt).check_status(notebook, snippet)
+  response['query_status'] = get_api(request, snippet).check_status(notebook, snippet)
   response['status'] = 0
 
   return JsonResponse(response)
@@ -118,7 +118,7 @@ def fetch_result_data(request):
   rows = json.loads(request.POST.get('rows', 100))
   start_over = json.loads(request.POST.get('startOver', False))
 
-  response['result'] = get_api(request.user, snippet, request.fs, request.jt).fetch_result(notebook, snippet, rows, start_over)
+  response['result'] = get_api(request, snippet).fetch_result(notebook, snippet, rows, start_over)
 
   # Materialize and HTML escape results
   if response['result'].get('data') and response['result'].get('type') == 'table':
@@ -138,7 +138,7 @@ def fetch_result_metadata(request):
   notebook = json.loads(request.POST.get('notebook', '{}'))
   snippet = json.loads(request.POST.get('snippet', '{}'))
 
-  response['result'] = get_api(request.user, snippet, request.fs, request.jt).fetch_result_metadata(notebook, snippet)
+  response['result'] = get_api(request, snippet).fetch_result_metadata(notebook, snippet)
   response['status'] = 0
 
   return JsonResponse(response)
@@ -153,7 +153,7 @@ def cancel_statement(request):
   notebook = json.loads(request.POST.get('notebook', '{}'))
   snippet = json.loads(request.POST.get('snippet', '{}'))
 
-  response['result'] = get_api(request.user, snippet, request.fs, request.jt).cancel(notebook, snippet)
+  response['result'] = get_api(request, snippet).cancel(notebook, snippet)
   response['status'] = 0
 
   return JsonResponse(response)
@@ -174,7 +174,7 @@ def get_logs(request):
   size = request.POST.get('size')
   size = int(size) if size else None
 
-  db = get_api(request.user, snippet, request.fs, request.jt)
+  db = get_api(request, snippet)
 
   logs = db.get_log(notebook, snippet, startFrom=startFrom, size=size)
 
@@ -314,7 +314,7 @@ def close_notebook(request):
 
   for session in notebook['sessions']:
     try:
-      response['result'].append(get_api(request.user, session, request.fs, request.jt).close_session(session))
+      response['result'].append(get_api(request, session).close_session(session))
     except QueryExpired:
       pass
     except Exception, e:
@@ -336,7 +336,7 @@ def close_statement(request):
   snippet = json.loads(request.POST.get('snippet', '{}'))
 
   try:
-    response['result'] = get_api(request.user, snippet, request.fs, request.jt).close_statement(snippet)
+    response['result'] = get_api(request, snippet).close_statement(snippet)
   except QueryExpired:
     pass
 
@@ -357,7 +357,7 @@ def autocomplete(request, server=None, database=None, table=None, column=None, n
   snippet = json.loads(request.POST.get('snippet', '{}'))
 
   try:
-    autocomplete_data = get_api(request.user, snippet, request.fs, request.jt).autocomplete(snippet, database, table, column, nested)
+    autocomplete_data = get_api(request, snippet).autocomplete(snippet, database, table, column, nested)
     response.update(autocomplete_data)
   except QueryExpired:
     pass

+ 11 - 12
desktop/libs/notebook/src/notebook/connectors/base.py

@@ -78,7 +78,7 @@ class Notebook(object):
     return '\n\n'.join([snippet['statement_raw'] for snippet in self.get_data()['snippets']])
 
 
-def get_api(user, snippet, fs, jt):
+def get_api(request, snippet):
   from notebook.connectors.hiveserver2 import HS2Api
   from notebook.connectors.jdbc import JdbcApi
   from notebook.connectors.rdbms import RdbmsApi
@@ -87,26 +87,26 @@ def get_api(user, snippet, fs, jt):
   from notebook.connectors.spark_batch import SparkBatchApi
   from notebook.connectors.text import TextApi
 
-  interpreter = [interpreter for interpreter in get_interpreters(user) if interpreter['type'] == snippet['type']]
+  interpreter = [interpreter for interpreter in get_interpreters(request.user) if interpreter['type'] == snippet['type']]
   if not interpreter:
     raise PopupException(_('Snippet type %(type)s is not configured in hue.ini') % snippet)
   interpreter = interpreter[0]
   interface = interpreter['interface']
 
   if interface == 'hiveserver2':
-    return HS2Api(user)
+    return HS2Api(user=request.user)
   elif interface == 'livy':
-    return SparkApi(user)
+    return SparkApi(request.user)
   elif interface == 'livy-batch':
-    return SparkBatchApi(user)
+    return SparkBatchApi(request.user)
   elif interface == 'text' or interface == 'markdown':
-    return TextApi(user)
+    return TextApi(request.user)
   elif interface == 'rdbms':
-    return RdbmsApi(user, interpreter=snippet['type'])
+    return RdbmsApi(request.user, interpreter=snippet['type'])
   elif interface == 'jdbc':
-    return JdbcApi(user, interpreter=interpreter)
+    return JdbcApi(request.user, interpreter=interpreter)
   elif interface == 'pig':
-    return PigApi(user, fs=fs, jt=jt)
+    return PigApi(user=request.user, request=request)
   else:
     raise PopupException(_('Notebook connector interface not recognized: %s') % interface)
 
@@ -123,11 +123,10 @@ def _get_snippet_session(notebook, snippet):
 
 class Api(object):
 
-  def __init__(self, user, fs=None, jt=None, interpreter=None):
+  def __init__(self, user, interpreter=None, request=None):
     self.user = user
-    self.fs = fs
-    self.jt = jt
     self.interpreter = interpreter
+    self.request = request
 
   def create_session(self, lang, properties=None):
     return {

+ 36 - 23
desktop/libs/notebook/src/notebook/connectors/pig_batch.py

@@ -18,8 +18,9 @@
 import logging
 import json
 
-from django.utils.translation import ugettext as _
 from django.core.urlresolvers import reverse
+from django.http import QueryDict
+from django.utils.translation import ugettext as _
 
 from notebook.connectors.base import Api, QueryError
 
@@ -37,6 +38,12 @@ except Exception, e:
 
 class PigApi(Api):
 
+  def __init__(self, *args, **kwargs):
+    Api.__init__(self, *args, **kwargs)
+
+    self.fs = self.request.fs
+    self.jt = self.request.jt
+
   def execute(self, notebook, snippet):
 
     attrs = {
@@ -54,20 +61,23 @@ class PigApi(Api):
 
     return {
       'id': oozie_id,
-      'watchUrl': reverse('pig:watch', kwargs={'job_id': oozie_id}) + '?format=python'
+      'watchUrl': reverse('pig:watch', kwargs={'job_id': oozie_id}) + '?format=python',
+      'has_result_set': True,
     }
 
   def check_status(self, notebook, snippet):
     job_id = snippet['result']['handle']['id']
-    request = MockRequest(self.user, self.fs, self.jt)
 
-    oozie_workflow = check_job_access_permission(request, job_id)
-    logs, workflow_actions, is_really_done = api.get(self.jt, self.jt, self.user).get_log(request, oozie_workflow)
+    oozie_workflow = check_job_access_permission(self.request, job_id)
+    logs, workflow_actions, is_really_done = self._get_output(oozie_workflow)
 
     if is_really_done and not oozie_workflow.is_running():
       if oozie_workflow.status in ('KILLED', 'FAILED'):
         raise QueryError(_('The script failed to run and was stopped'))
-      status = 'available'
+      if logs:
+        status = 'available'
+      else:
+        status = 'running' # Tricky case when the logs are being moved by YARN at job completion
     elif oozie_workflow.is_running():
       status = 'running'
     else:
@@ -77,16 +87,28 @@ class PigApi(Api):
         'status': status
     }
 
+  def _get_output(self, oozie_workflow):
+    q = QueryDict(self.request.GET, mutable=True)
+    q['format'] = 'python' # Hack for triggering the good section in single_task_attempt_logs
+    self.request.GET = q
+
+    logs, workflow_actions, is_really_done = api.get(self.fs, self.jt, self.user).get_log(self.request, oozie_workflow)
+
+    return logs, workflow_actions, is_really_done
+
   def fetch_result(self, notebook, snippet, rows, start_over):
     job_id = snippet['result']['handle']['id']
 
-    oozie_workflow = check_job_access_permission(MockRequest(self.user, self.fs, self.jt), job_id)
-    output = get_workflow_output(oozie_workflow, self.fs)
+    oozie_workflow = check_job_access_permission(self.request, job_id)
+    logs, workflow_actions, is_really_done = self._get_output(oozie_workflow)
+
+    output = logs.get('pig', _('No result'))
 
     return {
-        'data':  [hdfs_link(output)],
+        'data':  [[line] for line in output.split('\n')], # hdfs_link()
         'meta': [{'name': 'Header', 'type': 'STRING_TYPE', 'comment': ''}],
-        'type': 'text'
+        'type': 'table',
+        'has_more': False,
     }
 
   def cancel(self, notebook, snippet):
@@ -101,17 +123,16 @@ class PigApi(Api):
 
   def get_log(self, notebook, snippet, startFrom=0, size=None):
     job_id = snippet['result']['handle']['id']
-    request = MockRequest(self.user, self.fs, self.jt)
 
-    oozie_workflow = check_job_access_permission(MockRequest(self.user, self.fs, self.jt), job_id)
-    logs, workflow_actions, is_really_done = api.get(self.jt, self.jt, self.user).get_log(request, oozie_workflow)
+    oozie_workflow = check_job_access_permission(self.request, job_id)
+    logs, workflow_actions, is_really_done = self._get_output(oozie_workflow)
 
-    return logs
+    return logs.get('pig', _('No result'))
 
   def progress(self, snippet, logs):
     job_id = snippet['result']['handle']['id']
 
-    oozie_workflow = check_job_access_permission(MockRequest(self.user, self.fs, self.jt), job_id)
+    oozie_workflow = check_job_access_permission(self.request, job_id)
     return oozie_workflow.get_progress(),
 
   def close_statement(self, snippet):
@@ -119,11 +140,3 @@ class PigApi(Api):
 
   def close_session(self, session):
     pass
-
-
-class MockRequest():
-
-  def __init__(self, user, fs, jt):
-    self.user = user
-    self.fs = fs
-    self.js = jt