Pārlūkot izejas kodu

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

Save as CSV format only.
Abraham Elmahrek 11 gadi atpakaļ
vecāks
revīzija
900f902

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

@@ -17,6 +17,7 @@
 
 
 import json
 import json
 import logging
 import logging
+import os
 
 
 from django.core.urlresolvers import reverse
 from django.core.urlresolvers import reverse
 from django.http import HttpResponse, Http404
 from django.http import HttpResponse, Http404
@@ -31,6 +32,7 @@ from jobsub.parameterization import substitute_variables
 import beeswax.models
 import beeswax.models
 
 
 from beeswax.forms import QueryForm
 from beeswax.forms import QueryForm
+from beeswax.data_export import upload
 from beeswax.design import HQLdesign
 from beeswax.design import HQLdesign
 from beeswax.server import dbms
 from beeswax.server import dbms
 from beeswax.server.dbms import expand_exception, get_query_server_config,\
 from beeswax.server.dbms import expand_exception, get_query_server_config,\
@@ -418,13 +420,19 @@ def save_results(request, query_history_id):
       try:
       try:
         if form.cleaned_data['save_target'] == form.SAVE_TYPE_DIR:
         if form.cleaned_data['save_target'] == form.SAVE_TYPE_DIR:
           target_dir = form.cleaned_data['target_dir']
           target_dir = form.cleaned_data['target_dir']
-          query_history = db.insert_query_into_directory(query_history, target_dir)
           response['type'] = 'hdfs'
           response['type'] = 'hdfs'
           response['id'] = query_history.id
           response['id'] = query_history.id
           response['query'] = query_history.query
           response['query'] = query_history.query
           response['path'] = target_dir
           response['path'] = target_dir
           response['success_url'] = '/filebrowser/view%s' % 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:
         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)
           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
           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__)
 LOG = logging.getLogger(__name__)
 
 
 _DATA_WAIT_SLEEP = 0.1                  # Sleep 0.1 sec before checking for data availability
 _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):
 def download(handle, format, db):
@@ -41,20 +41,42 @@ def download(handle, format, db):
     LOG.error('Unknown download format "%s"' % (format,))
     LOG.error('Unknown download format "%s"' % (format,))
     return
     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')
   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.
   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)
     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()]
   data = [results.cols()]
 
 
@@ -68,8 +90,8 @@ def HS2DataAdapter(handle, format, db, max_rows=0):
       break
       break
 
 
     if results.has_more:
     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:
     else:
       results = None
       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"),
   target_dir = PathField(label=_t("Results Location"),
                          required=False,
                          required=False,
                          help_text=_t("Empty directory in HDFS to store results."))
                          help_text=_t("Empty directory in HDFS to store results."))
+  rerun = forms.BooleanField(label=_t("Run an export query"),
+                             initial=False,
+                             required=False)
   dependencies = [
   dependencies = [
     ('save_target', SAVE_TYPE_TBL, 'target_table'),
     ('save_target', SAVE_TYPE_TBL, 'target_table'),
     ('save_target', SAVE_TYPE_DIR, 'target_dir'),
     ('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'">
             <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">
               <input data-bind="value: $root.design.results.save.path" type="text" name="target_dir" placeholder="${_('Results location')}" class="pathChooser">
             </span>
             </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>
         </div>
         </div>
       </fieldset>
       </fieldset>
@@ -1711,13 +1717,13 @@ $(document).on('error.query', function () {
     if (firstPos > -1) {
     if (firstPos > -1) {
       selectedLine = $.trim(err.substring(err.indexOf(" ", firstPos), err.indexOf(":", firstPos))) * 1;
       selectedLine = $.trim(err.substring(err.indexOf(" ", firstPos), err.indexOf(":", firstPos))) * 1;
       errorWidgets.push(
       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();
       $(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,\
 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
   HIVE_SERVER_TEST_PORT, fetch_query_result_data
 from beeswax.design import hql_query, strip_trailing_semicolon
 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.models import SavedQuery, QueryHistory, HQL
 from beeswax.server import dbms
 from beeswax.server import dbms
 from beeswax.server.dbms import QueryServerException
 from beeswax.server.dbms import QueryServerException
@@ -649,6 +649,15 @@ for x in sys.stdin:
     csv_resp = download(handle, 'csv', self.db)
     csv_resp = download(handle, 'csv', self.db)
     assert_equal(csv_resp.content.replace('.0', ''), dataset.csv.replace('.0', ''))
     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):
   def test_designs(self):
     cli = self.client
     cli = self.client
 
 
@@ -837,12 +846,13 @@ for x in sys.stdin:
 
 
   def test_save_results_to_dir(self):
   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)
       content = json.loads(select_resp.content)
       qid = content['id']
       qid = content['id']
       save_data = {
       save_data = {
         'type': 'hdfs',
         '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)
       resp = self.client.post('/beeswax/api/query/%s/results/save' % qid, save_data, follow=True)
       content = json.loads(resp.content)
       content = json.loads(resp.content)
@@ -858,6 +868,12 @@ for x in sys.stdin:
         target_ls = self.cluster.fs.listdir(target_dir)
         target_ls = self.cluster.fs.listdir(target_dir)
         assert_true(len(target_ls) >= 1)
         assert_true(len(target_ls) >= 1)
         data_buf = ""
         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:
         for target in target_ls:
           target_file = self.cluster.fs.open(target_dir + '/' + target)
           target_file = self.cluster.fs.open(target_dir + '/' + target)
           data_buf += target_file.read()
           data_buf += target_file.read()
@@ -901,6 +917,14 @@ for x in sys.stdin:
     resp = self.client.get(resp.success_url)
     resp = self.client.get(resp.success_url)
     assert_true('File Browser' in resp.content, resp.content)
     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):
   def test_save_results_to_tbl(self):
 
 

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

@@ -56,7 +56,8 @@ function BeeswaxViewModel(server) {
       'save': {
       'save': {
         'errors': null,
         'errors': null,
         'type': 'hive-table',
         'type': 'hive-table',
-        'path': null
+        'path': null,
+        'rerun': true
       }
       }
     },
     },
     'watch': {
     'watch': {
@@ -733,7 +734,8 @@ function BeeswaxViewModel(server) {
         'database': self.database(),
         'database': self.database(),
         'server': self.server(),
         'server': self.server(),
         'type': self.design.results.save.type(),
         '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 url = '/' + self.server() + '/api/query/' + self.design.history.id() + '/results/save';
       var request = {
       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.
   Return a dataset object for a csv or excel document.
   """
   """
   dataset = tablib.Dataset()
   dataset = tablib.Dataset()
-  dataset.headers = format(headers, encoding)
+
+  if headers:
+    dataset.headers = format(headers, encoding)
 
 
   for row in data:
   for row in data:
     dataset.append(format(row, encoding))
     dataset.append(format(row, encoding))