Эх сурвалжийг харах

HUE-1686 [beeswax] Saving results should not re-execute the query

Save as CSV format only.
Abraham Elmahrek 11 жил өмнө
parent
commit
900f902

+ 10 - 2
apps/beeswax/src/beeswax/api.py

@@ -17,6 +17,7 @@
 
 import json
 import logging
+import os
 
 from django.core.urlresolvers import reverse
 from django.http import HttpResponse, Http404
@@ -31,6 +32,7 @@ from jobsub.parameterization import substitute_variables
 import beeswax.models
 
 from beeswax.forms import QueryForm
+from beeswax.data_export import upload
 from beeswax.design import HQLdesign
 from beeswax.server import dbms
 from beeswax.server.dbms import expand_exception, get_query_server_config,\
@@ -418,13 +420,19 @@ def save_results(request, query_history_id):
       try:
         if form.cleaned_data['save_target'] == form.SAVE_TYPE_DIR:
           target_dir = form.cleaned_data['target_dir']
-          query_history = db.insert_query_into_directory(query_history, target_dir)
           response['type'] = 'hdfs'
           response['id'] = query_history.id
           response['query'] = query_history.query
           response['path'] = target_dir
           response['success_url'] = '/filebrowser/view%s' % target_dir
-          response['watch_url'] = reverse(get_app_name(request) + ':api_watch_query_refresh_json', kwargs={'id': query_history.id})
+          if form.cleaned_data['rerun']:
+            query_history = db.insert_query_into_directory(query_history, target_dir)
+            response['watch_url'] = reverse(get_app_name(request) + ':api_watch_query_refresh_json', kwargs={'id': query_history.id})
+          else:
+            request.fs.do_as_user(request.user.username, request.fs.mkdir, form.cleaned_data['target_dir'])
+            path = os.path.join(form.cleaned_data['target_dir'], 'results')
+            upload(path, handle, request.user, db, request.fs)
+            response['watch_url'] = reverse(get_app_name(request) + ':api_watch_query_refresh_json', kwargs={'id': query_history.id})
         elif form.cleaned_data['save_target'] == form.SAVE_TYPE_TBL:
           query_history = db.create_table_as_a_select(request, query_history, form.target_database, form.cleaned_data['target_table'], result_meta)
           response['id'] = query_history.id

+ 31 - 9
apps/beeswax/src/beeswax/data_export.py

@@ -28,7 +28,7 @@ from beeswax import common, conf
 LOG = logging.getLogger(__name__)
 
 _DATA_WAIT_SLEEP = 0.1                  # Sleep 0.1 sec before checking for data availability
-FETCH_ROWS = 100000
+FETCH_ROWS = 10
 
 
 def download(handle, format, db):
@@ -41,20 +41,42 @@ def download(handle, format, db):
     LOG.error('Unknown download format "%s"' % (format,))
     return
 
-  data = HS2DataAdapter(handle, format, db, conf.DOWNLOAD_ROW_LIMIT.get())
+  data, has_more = HS2DataAdapter(handle, db, conf.DOWNLOAD_ROW_LIMIT.get())
   return export_csvxls.make_response(data[0], data[1:], format, 'query_result')
 
 
-def HS2DataAdapter(handle, format, db, max_rows=0):
+def upload(path, handle, user, db, fs):
   """
-  HS2DataAdapter(query_model, format, db) -> 2D array of data.
+  upload(query_model, path, user, db, fs) -> None
+
+  Retrieve the query result in the format specified and upload to hdfs.
+  """
+  has_more = True
+  start_over = True
+
+  fs.do_as_user(user.username, fs.create, path, overwrite=True)
+
+  while has_more:
+    data, has_more = HS2DataAdapter(handle, db, conf.DOWNLOAD_ROW_LIMIT.get(), start_over=start_over)
+    dataset = export_csvxls.dataset(None, data[1:])
+    fs.do_as_user(user.username, fs.append, path, dataset.csv)
+
+    if start_over:
+      start_over = False
+
+
+def HS2DataAdapter(handle, db, max_rows=0, start_over=True):
+  """
+  HS2DataAdapter(query_model, db) -> 2D array of data.
 
   First line should be the headers.
   """
-  results = db.fetch(handle, start_over=True, rows=FETCH_ROWS)
-  while not results.ready:   # For Beeswax
+  fetch_rows = max_rows if max_rows > -1 else FETCH_ROWS
+
+  results = db.fetch(handle, start_over=start_over, rows=fetch_rows)
+  while not results.ready:
     time.sleep(_DATA_WAIT_SLEEP)
-    results = db.fetch(handle, start_over=True, rows=FETCH_ROWS)
+    results = db.fetch(handle, start_over=start_over, rows=fetch_rows)
 
   data = [results.cols()]
 
@@ -68,8 +90,8 @@ def HS2DataAdapter(handle, format, db, max_rows=0):
       break
 
     if results.has_more:
-      results = db.fetch(handle, start_over=False, rows=FETCH_ROWS)
+      results = db.fetch(handle, start_over=False, rows=fetch_rows)
     else:
       results = None
 
-  return data
+  return data, results.has_more if results else False

+ 3 - 0
apps/beeswax/src/beeswax/forms.py

@@ -92,6 +92,9 @@ class SaveResultsForm(DependencyAwareForm):
   target_dir = PathField(label=_t("Results Location"),
                          required=False,
                          help_text=_t("Empty directory in HDFS to store results."))
+  rerun = forms.BooleanField(label=_t("Run an export query"),
+                             initial=False,
+                             required=False)
   dependencies = [
     ('save_target', SAVE_TYPE_TBL, 'target_table'),
     ('save_target', SAVE_TYPE_DIR, 'target_dir'),

+ 13 - 7
apps/beeswax/src/beeswax/templates/execute.mako

@@ -614,6 +614,12 @@ ${layout.menubar(section='query')}
             <span data-bind="visible: $root.design.results.save.type() == 'hdfs'">
               <input data-bind="value: $root.design.results.save.path" type="text" name="target_dir" placeholder="${_('Results location')}" class="pathChooser">
             </span>
+            % if app_name != 'impala':
+            <label class="radio" data-bind="visible: $root.design.results.save.type() == 'hdfs'">
+              <input data-bind="checked: $root.design.results.save.rerun" type="checkbox" name="rerun">
+              ${ _('Run an export query') }
+            </label>
+            % endif
           </div>
         </div>
       </fieldset>
@@ -1711,13 +1717,13 @@ $(document).on('error.query', function () {
     if (firstPos > -1) {
       selectedLine = $.trim(err.substring(err.indexOf(" ", firstPos), err.indexOf(":", firstPos))) * 1;
       errorWidgets.push(
-         codeMirror.addLineWidget(
-             selectedLine - 1,
-             $("<div>").addClass("editorError").html("<i class='fa fa-exclamation-circle'></i> " + err)[0], {
-                 coverGutter: true,
-                 noHScroll: true
-             }
-         )
+        codeMirror.addLineWidget(
+          selectedLine - 1,
+          $("<div>").addClass("editorError").html("<i class='fa fa-exclamation-circle'></i> " + err)[0], {
+            coverGutter: true,
+            noHScroll: true
+          }
+        )
       );
       $(el).hide();
     }

+ 27 - 3
apps/beeswax/src/beeswax/tests.py

@@ -55,7 +55,7 @@ from beeswax.views import collapse_whitespace
 from beeswax.test_base import make_query, wait_for_query_to_finish, verify_history, get_query_server_config,\
   HIVE_SERVER_TEST_PORT, fetch_query_result_data
 from beeswax.design import hql_query, strip_trailing_semicolon
-from beeswax.data_export import download
+from beeswax.data_export import upload
 from beeswax.models import SavedQuery, QueryHistory, HQL
 from beeswax.server import dbms
 from beeswax.server.dbms import QueryServerException
@@ -649,6 +649,15 @@ for x in sys.stdin:
     csv_resp = download(handle, 'csv', self.db)
     assert_equal(csv_resp.content.replace('.0', ''), dataset.csv.replace('.0', ''))
 
+  def test_data_upload(self):
+    hql = 'SELECT * FROM test'
+    query = hql_query(hql)
+
+    handle = self.db.execute_and_wait(query)
+    upload('/tmp/test_data_upload.csv', handle, self.db, self.cluster.fs)
+
+    assert_true(self.cluster.fs.exists('/tmp/test_data_upload.csv'))
+
   def test_designs(self):
     cli = self.client
 
@@ -837,12 +846,13 @@ for x in sys.stdin:
 
   def test_save_results_to_dir(self):
 
-    def save_and_verify(select_resp, target_dir, verify=True):
+    def save_and_verify(select_resp, target_dir, rerun=True, verify=True):
       content = json.loads(select_resp.content)
       qid = content['id']
       save_data = {
         'type': 'hdfs',
-        'path': target_dir
+        'path': target_dir,
+        'rerun': rerun
       }
       resp = self.client.post('/beeswax/api/query/%s/results/save' % qid, save_data, follow=True)
       content = json.loads(resp.content)
@@ -858,6 +868,12 @@ for x in sys.stdin:
         target_ls = self.cluster.fs.listdir(target_dir)
         assert_true(len(target_ls) >= 1)
         data_buf = ""
+
+        if not rerun:
+          assert_equal(len(target_ls), 1)
+          # filename is 'results'
+          assert_equal(target_ls[0], 'results')
+
         for target in target_ls:
           target_file = self.cluster.fs.open(target_dir + '/' + target)
           data_buf += target_file.read()
@@ -901,6 +917,14 @@ for x in sys.stdin:
     resp = self.client.get(resp.success_url)
     assert_true('File Browser' in resp.content, resp.content)
 
+    # SELECT columns. (Result dir is in /tmp.)
+    # Do not rerun
+    hql = "SELECT foo, bar FROM test"
+    resp = _make_query(self.client, hql, wait=True, local=False, max=180.0)
+    resp = save_and_verify(resp, TARGET_DIR_ROOT + '/4', rerun=False, verify=True)
+    resp = self.client.get(resp.success_url)
+    assert_true('File Browser' in resp.content, resp.content)
+
 
   def test_save_results_to_tbl(self):
 

+ 4 - 2
apps/beeswax/static/js/beeswax.vm.js

@@ -56,7 +56,8 @@ function BeeswaxViewModel(server) {
       'save': {
         'errors': null,
         'type': 'hive-table',
-        'path': null
+        'path': null,
+        'rerun': true
       }
     },
     'watch': {
@@ -733,7 +734,8 @@ function BeeswaxViewModel(server) {
         'database': self.database(),
         'server': self.server(),
         'type': self.design.results.save.type(),
-        'path': self.design.results.save.path()
+        'path': self.design.results.save.path(),
+        'rerun': self.design.results.save.rerun()
       };
       var url = '/' + self.server() + '/api/query/' + self.design.history.id() + '/results/save';
       var request = {

+ 3 - 1
desktop/core/src/desktop/lib/export_csvxls.py

@@ -43,7 +43,9 @@ def dataset(headers, data, encoding=None):
   Return a dataset object for a csv or excel document.
   """
   dataset = tablib.Dataset()
-  dataset.headers = format(headers, encoding)
+
+  if headers:
+    dataset.headers = format(headers, encoding)
 
   for row in data:
     dataset.append(format(row, encoding))