Răsfoiți Sursa

[beeswax] Raise QueryServerException on timeout in execute_and_wait, and handle

Jenny Kim 10 ani în urmă
părinte
comite
a6bb39f

+ 18 - 12
apps/beeswax/src/beeswax/api.py

@@ -114,9 +114,9 @@ def autocomplete(request, database=None, table=None, column=None, nested=None):
 
         inner_type = _get_complex_inner_type(current, extended_type, simple_type)
         response.update(inner_type)
-  except TTransportException, tx:
+  except (QueryServerException, TTransportException), e:
     response['code'] = 503
-    response['error'] = tx.message
+    response['error'] = e.message
   except Exception, e:
     LOG.warn('Autocomplete data fetching error %s.%s: %s' % (database, table, e))
     response['code'] = 500
@@ -672,14 +672,17 @@ def get_table_stats(request, database, table, column=None):
 
   response = {'status': -1, 'message': '', 'redirect': ''}
 
-  if column is not None:
-    stats = db.get_table_columns_stats(database, table, column)
-  else:
-    table = db.get_table(database, table)
-    stats = table.stats
+  try:
+    if column is not None:
+      stats = db.get_table_columns_stats(database, table, column)
+    else:
+      table = db.get_table(database, table)
+      stats = table.stats
 
-  response['stats'] = stats
-  response['status'] = 0
+    response['stats'] = stats
+    response['status'] = 0
+  except QueryServerException, e:
+    response['message'] = _('Failed to get table stats for table %s.%s: %s' % (database, table, e.message))
 
   return JsonResponse(response)
 
@@ -691,10 +694,13 @@ def get_top_terms(request, database, table, column, prefix=None):
 
   response = {'status': -1, 'message': '', 'redirect': ''}
 
-  terms = db.get_top_terms(database, table, column, prefix=prefix, limit=int(request.GET.get('limit', 30)))
+  try:
+    terms = db.get_top_terms(database, table, column, prefix=prefix, limit=int(request.GET.get('limit', 30)))
 
-  response['terms'] = terms
-  response['status'] = 0
+    response['terms'] = terms
+    response['status'] = 0
+  except QueryServerException, e:
+    response['message'] = _('Failed to get table stats for table %s.%s: %s' % (database, table, e.message))
 
   return JsonResponse(response)
 

+ 12 - 3
apps/beeswax/src/beeswax/server/dbms.py

@@ -24,7 +24,7 @@ from django.utils.encoding import force_unicode
 from django.utils.translation import ugettext as _
 
 from beeswax import hive_site
-from beeswax.conf import HIVE_SERVER_HOST, HIVE_SERVER_PORT, BROWSE_PARTITIONED_TABLE_LIMIT
+from beeswax.conf import HIVE_SERVER_HOST, HIVE_SERVER_PORT, BROWSE_PARTITIONED_TABLE_LIMIT, SERVER_CONN_TIMEOUT
 from beeswax.design import hql_query
 from beeswax.hive_site import hiveserver2_use_ssl
 from beeswax.models import QueryHistory, QUERY_TYPES
@@ -125,7 +125,9 @@ class HiveServer2Dbms(object):
   def get_tables(self, database='default', table_names='*'):
     hql = "SHOW TABLES IN `%s` '%s'" % (database, table_names) # self.client.get_tables(database, table_names) is too slow
     query = hql_query(hql)
-    handle = self.execute_and_wait(query, timeout_sec=15.0)
+    timeout = SERVER_CONN_TIMEOUT.get()
+
+    handle = self.execute_and_wait(query, timeout_sec=timeout)
 
     if handle:
       result = self.fetch(handle, rows=5000)
@@ -495,10 +497,17 @@ class HiveServer2Dbms(object):
       curr = time.time()
 
     try:
+      msg = "The query timed out after %(timeout)d seconds, canceled query [%(query)s]..." % \
+              {'timeout': timeout_sec, 'query': query.hql_query[:40]}
+      LOG.exception(msg)
       self.cancel_operation(handle)
+      raise QueryServerException(Exception(msg), message=msg)
     except:
-      LOG.exception('failed to cancel operation')
+      msg = "Failed to cancel query [%(query)s]..." % {'query': query.hql_query[:40]}
+      LOG.exception(msg)
       self.close_operation(handle)
+      raise QueryServerException(Exception(msg), message=msg)
+
     return None
 
 

+ 69 - 67
apps/impala/src/impala/dashboards.py

@@ -21,9 +21,9 @@ import json
 
 from math import log
 
+from django.utils.encoding import force_unicode
 from django.utils.translation import ugettext as _
 
-from desktop.context_processors import get_app_name
 from desktop.lib.django_util import JsonResponse, render
 from desktop.models import Document2
 
@@ -57,80 +57,82 @@ def query(request):
     'status': -1,
     'data': {}
   }
-  
-  dashboard = json.loads(request.POST['dashboard'])
-  fqs = json.loads(request.POST['query'])['fqs']
 
-  database = dashboard['properties'][0]['database']
-  table = dashboard['properties'][0]['table']
-  
-  
-  if fqs:
-    filters = ' AND '.join(['%s = %s' % (fq['field'], value) for fq in fqs for value in fq['filter']])
-  else:
-    filters = ''
-
-  if 'facet' in request.POST:
-    facet = json.loads(request.POST.get('facet'))
-    slot = None
-
-    if facet['type'] == 'field':
-      template = "SELECT %(field)s, COUNT(*) AS top FROM %(database)s.%(table)s WHERE %(field)s IS NOT NULL %(filters)s GROUP BY %(field)s ORDER BY top DESC LIMIT %(limit)s"
-    elif facet['type'] == 'range':
-      slot = (facet['properties']['end'] - facet['properties']['start']) / facet['properties']['limit']
-      template = """select cast(%(field)s / %(slot)s AS int) * %(slot)s, count(*) AS top, cast(%(field)s / %(slot)s as int) as s 
-       FROM %(database)s.%(table)s WHERE %(field)s IS NOT NULL GROUP BY s ORDER BY s DESC LIMIT %(limit)s"""    
-    else:            
-      # Simple Top
-      template = "SELECT DISTINCT %(field)s FROM %(database)s.%(table)s WHERE %(field)s IS NOT NULL %(filters)s ORDER BY %(field)s DESC LIMIT %(limit)s"
-    
-    facet = json.loads(request.POST['facet'])
-    hql = template % {
-        'database': database,
-        'table': table,
-        'limit': facet['properties']['limit'],
-        'field': facet['field'],
-        'filters': (' AND ' + filters) if filters else '',
-        'slot': slot
-    }
-    result['id'] = facet['id']
-    result['field'] = facet['field']
-    fields = [fq['field'] for fq in fqs]
-    result['selected'] = facet['field'] in fields
-  else:
-    dashboard['resultsetSelectedFields'] = map(lambda f: '`%s`' % f if f in ('date',) else f, dashboard['resultsetSelectedFields'])
-    fields = ', '.join(dashboard['resultsetSelectedFields']) if dashboard['resultsetSelectedFields'] else '*'
-    hql = "SELECT %(fields)s FROM %(database)s.%(table)s" % {
-        'database': database, 
-        'table': table,
-        'fields': fields
-    }
-    if filters:
-      hql += ' WHERE ' + filters
-    hql += ' LIMIT 100'
+  try:
+    dashboard = json.loads(request.POST['dashboard'])
+    fqs = json.loads(request.POST['query'])['fqs']
 
-  query_server = get_query_server_config(name='impala')
-  db = dbms.get(request.user, query_server=query_server)
-  
-  print hql
-  query = hql_query(hql)  
-  handle = db.execute_and_wait(query, timeout_sec=35.0)
+    database = dashboard['properties'][0]['database']
+    table = dashboard['properties'][0]['table']
+
+    if fqs:
+      filters = ' AND '.join(['%s = %s' % (fq['field'], value) for fq in fqs for value in fq['filter']])
+    else:
+      filters = ''
 
-  if handle:
-    data = db.fetch(handle, rows=100)
     if 'facet' in request.POST:
       facet = json.loads(request.POST.get('facet'))
-      result['type'] = facet['type']
-      if facet['type'] == 'top':
-        result['data'] = [{"value": row[0], "count": None, "selected": False, "cat": facet['field']} for row in data.rows()]
+      slot = None
+
+      if facet['type'] == 'field':
+        template = "SELECT %(field)s, COUNT(*) AS top FROM %(database)s.%(table)s WHERE %(field)s IS NOT NULL %(filters)s GROUP BY %(field)s ORDER BY top DESC LIMIT %(limit)s"
+      elif facet['type'] == 'range':
+        slot = (facet['properties']['end'] - facet['properties']['start']) / facet['properties']['limit']
+        template = """select cast(%(field)s / %(slot)s AS int) * %(slot)s, count(*) AS top, cast(%(field)s / %(slot)s as int) as s
+         FROM %(database)s.%(table)s WHERE %(field)s IS NOT NULL GROUP BY s ORDER BY s DESC LIMIT %(limit)s"""
       else:
-        result['data'] = [{"value": row[0], "count": row[1], "selected": False, "cat": facet['field']} for row in data.rows()]
+        # Simple Top
+        template = "SELECT DISTINCT %(field)s FROM %(database)s.%(table)s WHERE %(field)s IS NOT NULL %(filters)s ORDER BY %(field)s DESC LIMIT %(limit)s"
+
+      facet = json.loads(request.POST['facet'])
+      hql = template % {
+          'database': database,
+          'table': table,
+          'limit': facet['properties']['limit'],
+          'field': facet['field'],
+          'filters': (' AND ' + filters) if filters else '',
+          'slot': slot
+      }
+      result['id'] = facet['id']
+      result['field'] = facet['field']
+      fields = [fq['field'] for fq in fqs]
+      result['selected'] = facet['field'] in fields
     else:
-      result['data'] = list(data.rows())
+      dashboard['resultsetSelectedFields'] = map(lambda f: '`%s`' % f if f in ('date',) else f, dashboard['resultsetSelectedFields'])
+      fields = ', '.join(dashboard['resultsetSelectedFields']) if dashboard['resultsetSelectedFields'] else '*'
+      hql = "SELECT %(fields)s FROM %(database)s.%(table)s" % {
+          'database': database,
+          'table': table,
+          'fields': fields
+      }
+      if filters:
+        hql += ' WHERE ' + filters
+      hql += ' LIMIT 100'
 
-    result['cols'] = list(data.cols())
-    result['status'] = 0
-    db.close(handle)
+    query_server = get_query_server_config(name='impala')
+    db = dbms.get(request.user, query_server=query_server)
+
+    query = hql_query(hql)
+    handle = db.execute_and_wait(query, timeout_sec=35.0)
+
+    if handle:
+      data = db.fetch(handle, rows=100)
+      if 'facet' in request.POST:
+        facet = json.loads(request.POST.get('facet'))
+        result['type'] = facet['type']
+        if facet['type'] == 'top':
+          result['data'] = [{"value": row[0], "count": None, "selected": False, "cat": facet['field']} for row in data.rows()]
+        else:
+          result['data'] = [{"value": row[0], "count": row[1], "selected": False, "cat": facet['field']} for row in data.rows()]
+      else:
+        result['data'] = list(data.rows())
+
+      result['cols'] = list(data.cols())
+      result['status'] = 0
+      db.close(handle)
+  except Exception, e:
+    LOG.exception(e)
+    result['message'] = force_unicode(e)
 
   return JsonResponse(result)
 

+ 13 - 10
apps/metastore/src/metastore/views.py

@@ -105,19 +105,22 @@ def show_tables(request, database=None):
 
   db = dbms.get(request.user)
 
-  databases = db.get_databases()
+  try:
+    databases = db.get_databases()
 
-  if database not in databases:
-    database = 'default'
+    if database not in databases:
+      database = 'default'
 
-  if request.method == 'POST':
-    db_form = DbForm(request.POST, databases=databases)
-    if db_form.is_valid():
-      database = db_form.cleaned_data['database']
-  else:
-    db_form = DbForm(initial={'database': database}, databases=databases)
+    if request.method == 'POST':
+      db_form = DbForm(request.POST, databases=databases)
+      if db_form.is_valid():
+        database = db_form.cleaned_data['database']
+    else:
+      db_form = DbForm(initial={'database': database}, databases=databases)
 
-  tables = db.get_tables(database=database)
+    tables = db.get_tables(database=database)
+  except Exception, e:
+    raise PopupException(_('Failed to retrieve tables for database % s' % database), detail=e)
 
   resp = render("tables.mako", request, {
     'breadcrumbs': [

+ 16 - 12
desktop/libs/indexer/src/indexer/controller.py

@@ -249,20 +249,24 @@ class CollectionManagerController(object):
       table = db.get_table(database, table)
       hql = "SELECT %s FROM `%s.%s` %s" % (','.join(columns), database, table.name, db._get_browse_limit_clause(table))
       query = dbms.hql_query(hql)
-      handle = db.execute_and_wait(query)
 
-      if handle:
-        result = db.fetch(handle, rows=100)
-        db.close(handle)
+      try:
+        handle = db.execute_and_wait(query)
 
-        dataset = tablib.Dataset()
-        dataset.append(columns)
-        for row in result.rows():
-          dataset.append(row)
+        if handle:
+          result = db.fetch(handle, rows=100)
+          db.close(handle)
 
-        if not api.update(collection_or_core_name, dataset.csv, content_type='csv'):
-          raise PopupException(_('Could not update index. Check error logs for more info.'))
-      else:
-        raise PopupException(_('Could not update index. Could not fetch any data from Hive.'))
+          dataset = tablib.Dataset()
+          dataset.append(columns)
+          for row in result.rows():
+            dataset.append(row)
+
+          if not api.update(collection_or_core_name, dataset.csv, content_type='csv'):
+            raise PopupException(_('Could not update index. Check error logs for more info.'))
+        else:
+          raise PopupException(_('Could not update index. Could not fetch any data from Hive.'))
+      except Exception, e:
+        raise PopupException(_('Could not update index.'), detail=e)
     else:
       raise PopupException(_('Could not update index. Indexing strategy %s not supported.') % indexing_strategy)