Преглед изворни кода

[notebook] Maintain underlying jobs list executed by the snippet

Maintain jobs (formerly job_urls) in notebook KO model and make API return full-set of jobs for the full execution
Jenny Kim пре 10 година
родитељ
комит
962a33465e

+ 16 - 7
desktop/libs/notebook/src/notebook/api.py

@@ -18,7 +18,6 @@
 import json
 import logging
 
-from django.core.urlresolvers import reverse
 from django.utils.translation import ugettext as _
 from django.views.decorators.http import require_GET, require_POST
 
@@ -175,12 +174,22 @@ def get_logs(request):
   size = int(size) if size else None
 
   db = get_api(request.user, snippet, request.fs, request.jt)
-  response['logs'] = db.get_log(notebook, snippet, startFrom=startFrom, size=size)
-  response['progress'] = db.progress(snippet, response['logs']) if snippet['status'] != 'available' and snippet['status'] != 'success' else 100
-  response['job_urls'] = [{
-      'name': job,
-      'url': reverse('jobbrowser.views.single_job', kwargs={'job': job})
-    } for job in db.get_jobs(response['logs'])]
+
+  logs = db.get_log(notebook, snippet, startFrom=startFrom, size=size)
+
+  jobs = json.loads(request.POST.get('jobs', '[]'))
+
+  # Get any new jobs from current logs snippet
+  new_jobs = db.get_jobs(logs)
+
+  # Append new jobs to known jobs and get the unique set
+  if new_jobs:
+    all_jobs = jobs + new_jobs
+    jobs = {job['name']: job for job in all_jobs}.values()
+
+  response['logs'] = logs
+  response['progress'] = db.progress(snippet, logs) if snippet['status'] != 'available' and snippet['status'] != 'success' else 100
+  response['jobs'] = jobs
   response['status'] = 0
 
   return JsonResponse(response)

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

@@ -150,5 +150,5 @@ class Api(object):
   def progress(self, snippet, logs=None):
     return 50
 
-  def get_jobs(self, log):
+  def get_jobs(self, logs):
     return []

+ 23 - 14
desktop/libs/notebook/src/notebook/connectors/hiveserver2.py

@@ -18,6 +18,8 @@
 import logging
 import re
 
+from django.core.urlresolvers import reverse
+
 from desktop.lib.exceptions_renderable import PopupException
 from desktop.lib.i18n import force_unicode
 
@@ -139,6 +141,20 @@ class HS2Api(Api):
     handle = self._get_handle(snippet)
     return db.get_log(handle, start_over=startFrom == 0)
 
+  @query_error_handler
+  def close_statement(self, snippet):
+    if snippet['type'] == 'impala':
+      from impala import conf as impala_conf
+
+    if (snippet['type'] == 'hive' and beeswax_conf.CLOSE_QUERIES.get()) or (snippet['type'] == 'impala' and impala_conf.CLOSE_QUERIES.get()):
+      db = self._get_db(snippet)
+
+      handle = self._get_handle(snippet)
+      db.close_operation(handle)
+      return {'status': 0}
+    else:
+      return {'status': -1}  # skipped
+
   def download(self, notebook, snippet, format):
     try:
       db = self._get_db(snippet)
@@ -168,22 +184,15 @@ class HS2Api(Api):
     else:
       return 50
 
-  @query_error_handler
-  def close_statement(self, snippet):
-    if snippet['type'] == 'impala':
-      from impala import conf as impala_conf
+  def get_jobs(self, logs):
+    job_ids = _parse_out_hadoop_jobs(logs)
 
-    if (snippet['type'] == 'hive' and beeswax_conf.CLOSE_QUERIES.get()) or (snippet['type'] == 'impala' and impala_conf.CLOSE_QUERIES.get()):
-      db = self._get_db(snippet)
-
-      handle = self._get_handle(snippet)
-      db.close_operation(handle)
-      return {'status': 0}
-    else:
-      return {'status': -1}  # skipped
+    jobs = [{
+      'name': job_id,
+      'url': reverse('jobbrowser.views.single_job', kwargs={'job': job_id})
+    } for job_id in job_ids]
 
-  def get_jobs(self, log):
-    return _parse_out_hadoop_jobs(log)
+    return jobs
 
   @query_error_handler
   def autocomplete(self, snippet, database=None, table=None, column=None, nested=None):

+ 1 - 1
desktop/libs/notebook/src/notebook/connectors/pig_batch.py

@@ -120,7 +120,7 @@ class PigApi(Api):
   def close_session(self, session):
     pass
 
-  def get_jobs(self, log):
+  def get_jobs(self, logs):
     return []
 
 

+ 1 - 1
desktop/libs/notebook/src/notebook/connectors/spark_batch.py

@@ -89,5 +89,5 @@ class SparkBatchApi(Api):
   def progress(self, snippet, logs):
     return 50
 
-  def get_jobs(self, log):
+  def get_jobs(self, logs):
     return []

+ 1 - 1
desktop/libs/notebook/src/notebook/connectors/spark_shell.py

@@ -215,5 +215,5 @@ class SparkApi(Api):
     else:
       return {'status': -1}
 
-  def get_jobs(self, log):
+  def get_jobs(self, logs):
     return []

+ 8 - 1
desktop/libs/notebook/src/notebook/static/notebook/js/notebook.ko.js

@@ -233,6 +233,7 @@ var Snippet = function (vm, notebook, snippet) {
   self.showChart = ko.observable(typeof snippet.showChart != "undefined" && snippet.showChart != null ? snippet.showChart : false);
   self.showLogs = ko.observable(typeof snippet.showLogs != "undefined" && snippet.showLogs != null ? snippet.showLogs : false);
   self.progress = ko.observable(typeof snippet.progress != "undefined" && snippet.progress != null ? snippet.progress : 0);
+  self.jobs = ko.observableArray(typeof snippet.jobs != "undefined" && snippet.jobs != null ? snippet.jobs : []);
 
   self.progress.subscribe(function (val) {
     $(document).trigger("progress", {data: val, snippet: self});
@@ -391,6 +392,7 @@ var Snippet = function (vm, notebook, snippet) {
     self.errors([]);
     self.result.logLines = 0;
     self.progress(0);
+    self.jobs([]);
 
     if (self.result.fetchedOnce()) {
       self.close();
@@ -568,7 +570,8 @@ var Snippet = function (vm, notebook, snippet) {
     $.post("/notebook/api/get_logs", {
       notebook: ko.mapping.toJSON(notebook.getContext()),
       snippet: ko.mapping.toJSON(self.getContext()),
-      from: self.result.logLines
+      from: self.result.logLines,
+      jobs: ko.mapping.toJSON(self.jobs)
     }, function (data) {
       if (data.status == 1) { // Append errors to the logs
         data.status = 0;
@@ -585,6 +588,9 @@ var Snippet = function (vm, notebook, snippet) {
             self.result.logs(oldLogs + "\n" + data.logs);
           }
         }
+        if (data.jobs.length > 0) {
+          self.jobs(data.jobs);
+        }
         self.progress(data.progress);
       } else {
         self._ajaxError(data);
@@ -609,6 +615,7 @@ var Snippet = function (vm, notebook, snippet) {
     if (self.status() == 'loading') {
       self.status('failed');
       self.progress(0);
+      self.jobs([]);
     }
   };
 };