浏览代码

[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 9 年之前
父节点
当前提交
bc0167d

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

@@ -74,7 +74,9 @@ def check_job_permission(view_func):
     except JobExpired, e:
     except JobExpired, e:
       raise PopupException(_('Job %s has expired.') % jobid, detail=_('Cannot be found on the History Server.'))
       raise PopupException(_('Job %s has expired.') % jobid, detail=_('Cannot be found on the History Server.'))
     except Exception, e:
     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 \
     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):
         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)
   return JsonResponse(response)
 
 
 
 
-
 @check_job_permission
 @check_job_permission
 def job_single_logs(request, job):
 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]})
   return single_task_attempt_logs(request, **{'job': job.jobId, 'taskid': task.taskId, 'attemptid': task.taskAttemptIds[-1]})
 
 
+
 @check_job_permission
 @check_job_permission
 def tasks(request, job):
 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]:
     if any(old_session) and 'properties' in old_session[0]:
       properties = old_session[0]['properties']
       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
   response['status'] = 0
 
 
   return JsonResponse(response)
   return JsonResponse(response)
@@ -66,7 +66,7 @@ def close_session(request):
 
 
   session = json.loads(request.POST.get('session', '{}'))
   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
   response['status'] = 0
 
 
   return JsonResponse(response)
   return JsonResponse(response)
@@ -81,7 +81,7 @@ def execute(request):
   notebook = json.loads(request.POST.get('notebook', '{}'))
   notebook = json.loads(request.POST.get('notebook', '{}'))
   snippet = json.loads(request.POST.get('snippet', '{}'))
   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
   # Materialize and HTML escape results
   if response['handle'].get('sync') and response['handle']['result'].get('data'):
   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', '{}'))
   notebook = json.loads(request.POST.get('notebook', '{}'))
   snippet = json.loads(request.POST.get('snippet', '{}'))
   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
   response['status'] = 0
 
 
   return JsonResponse(response)
   return JsonResponse(response)
@@ -118,7 +118,7 @@ def fetch_result_data(request):
   rows = json.loads(request.POST.get('rows', 100))
   rows = json.loads(request.POST.get('rows', 100))
   start_over = json.loads(request.POST.get('startOver', False))
   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
   # Materialize and HTML escape results
   if response['result'].get('data') and response['result'].get('type') == 'table':
   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', '{}'))
   notebook = json.loads(request.POST.get('notebook', '{}'))
   snippet = json.loads(request.POST.get('snippet', '{}'))
   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
   response['status'] = 0
 
 
   return JsonResponse(response)
   return JsonResponse(response)
@@ -153,7 +153,7 @@ def cancel_statement(request):
   notebook = json.loads(request.POST.get('notebook', '{}'))
   notebook = json.loads(request.POST.get('notebook', '{}'))
   snippet = json.loads(request.POST.get('snippet', '{}'))
   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
   response['status'] = 0
 
 
   return JsonResponse(response)
   return JsonResponse(response)
@@ -174,7 +174,7 @@ def get_logs(request):
   size = request.POST.get('size')
   size = request.POST.get('size')
   size = int(size) if size else None
   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)
   logs = db.get_log(notebook, snippet, startFrom=startFrom, size=size)
 
 
@@ -314,7 +314,7 @@ def close_notebook(request):
 
 
   for session in notebook['sessions']:
   for session in notebook['sessions']:
     try:
     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:
     except QueryExpired:
       pass
       pass
     except Exception, e:
     except Exception, e:
@@ -336,7 +336,7 @@ def close_statement(request):
   snippet = json.loads(request.POST.get('snippet', '{}'))
   snippet = json.loads(request.POST.get('snippet', '{}'))
 
 
   try:
   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:
   except QueryExpired:
     pass
     pass
 
 
@@ -357,7 +357,7 @@ def autocomplete(request, server=None, database=None, table=None, column=None, n
   snippet = json.loads(request.POST.get('snippet', '{}'))
   snippet = json.loads(request.POST.get('snippet', '{}'))
 
 
   try:
   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)
     response.update(autocomplete_data)
   except QueryExpired:
   except QueryExpired:
     pass
     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']])
     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.hiveserver2 import HS2Api
   from notebook.connectors.jdbc import JdbcApi
   from notebook.connectors.jdbc import JdbcApi
   from notebook.connectors.rdbms import RdbmsApi
   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.spark_batch import SparkBatchApi
   from notebook.connectors.text import TextApi
   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:
   if not interpreter:
     raise PopupException(_('Snippet type %(type)s is not configured in hue.ini') % snippet)
     raise PopupException(_('Snippet type %(type)s is not configured in hue.ini') % snippet)
   interpreter = interpreter[0]
   interpreter = interpreter[0]
   interface = interpreter['interface']
   interface = interpreter['interface']
 
 
   if interface == 'hiveserver2':
   if interface == 'hiveserver2':
-    return HS2Api(user)
+    return HS2Api(user=request.user)
   elif interface == 'livy':
   elif interface == 'livy':
-    return SparkApi(user)
+    return SparkApi(request.user)
   elif interface == 'livy-batch':
   elif interface == 'livy-batch':
-    return SparkBatchApi(user)
+    return SparkBatchApi(request.user)
   elif interface == 'text' or interface == 'markdown':
   elif interface == 'text' or interface == 'markdown':
-    return TextApi(user)
+    return TextApi(request.user)
   elif interface == 'rdbms':
   elif interface == 'rdbms':
-    return RdbmsApi(user, interpreter=snippet['type'])
+    return RdbmsApi(request.user, interpreter=snippet['type'])
   elif interface == 'jdbc':
   elif interface == 'jdbc':
-    return JdbcApi(user, interpreter=interpreter)
+    return JdbcApi(request.user, interpreter=interpreter)
   elif interface == 'pig':
   elif interface == 'pig':
-    return PigApi(user, fs=fs, jt=jt)
+    return PigApi(user=request.user, request=request)
   else:
   else:
     raise PopupException(_('Notebook connector interface not recognized: %s') % interface)
     raise PopupException(_('Notebook connector interface not recognized: %s') % interface)
 
 
@@ -123,11 +123,10 @@ def _get_snippet_session(notebook, snippet):
 
 
 class Api(object):
 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.user = user
-    self.fs = fs
-    self.jt = jt
     self.interpreter = interpreter
     self.interpreter = interpreter
+    self.request = request
 
 
   def create_session(self, lang, properties=None):
   def create_session(self, lang, properties=None):
     return {
     return {

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

@@ -18,8 +18,9 @@
 import logging
 import logging
 import json
 import json
 
 
-from django.utils.translation import ugettext as _
 from django.core.urlresolvers import reverse
 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
 from notebook.connectors.base import Api, QueryError
 
 
@@ -37,6 +38,12 @@ except Exception, e:
 
 
 class PigApi(Api):
 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):
   def execute(self, notebook, snippet):
 
 
     attrs = {
     attrs = {
@@ -54,20 +61,23 @@ class PigApi(Api):
 
 
     return {
     return {
       'id': oozie_id,
       '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):
   def check_status(self, notebook, snippet):
     job_id = snippet['result']['handle']['id']
     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 is_really_done and not oozie_workflow.is_running():
       if oozie_workflow.status in ('KILLED', 'FAILED'):
       if oozie_workflow.status in ('KILLED', 'FAILED'):
         raise QueryError(_('The script failed to run and was stopped'))
         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():
     elif oozie_workflow.is_running():
       status = 'running'
       status = 'running'
     else:
     else:
@@ -77,16 +87,28 @@ class PigApi(Api):
         'status': status
         '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):
   def fetch_result(self, notebook, snippet, rows, start_over):
     job_id = snippet['result']['handle']['id']
     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 {
     return {
-        'data':  [hdfs_link(output)],
+        'data':  [[line] for line in output.split('\n')], # hdfs_link()
         'meta': [{'name': 'Header', 'type': 'STRING_TYPE', 'comment': ''}],
         'meta': [{'name': 'Header', 'type': 'STRING_TYPE', 'comment': ''}],
-        'type': 'text'
+        'type': 'table',
+        'has_more': False,
     }
     }
 
 
   def cancel(self, notebook, snippet):
   def cancel(self, notebook, snippet):
@@ -101,17 +123,16 @@ class PigApi(Api):
 
 
   def get_log(self, notebook, snippet, startFrom=0, size=None):
   def get_log(self, notebook, snippet, startFrom=0, size=None):
     job_id = snippet['result']['handle']['id']
     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):
   def progress(self, snippet, logs):
     job_id = snippet['result']['handle']['id']
     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(),
     return oozie_workflow.get_progress(),
 
 
   def close_statement(self, snippet):
   def close_statement(self, snippet):
@@ -119,11 +140,3 @@ class PigApi(Api):
 
 
   def close_session(self, session):
   def close_session(self, session):
     pass
     pass
-
-
-class MockRequest():
-
-  def __init__(self, user, fs, jt):
-    self.user = user
-    self.fs = fs
-    self.js = jt