瀏覽代碼

[beeswax] Add a get_settings API endpoint and add test

GET /beeswax/api/settings/
GET /impala/api/settings/
Jenny Kim 9 年之前
父節點
當前提交
60774a6

+ 93 - 76
apps/beeswax/src/beeswax/api.py

@@ -609,40 +609,7 @@ def save_results_hive_table(request, query_history_id):
   return JsonResponse(response)
 
 
-def design_to_dict(design):
-  hql_design = HQLdesign.loads(design.data)
-  return {
-    'id': design.id,
-    'query': hql_design.hql_query,
-    'name': design.name,
-    'desc': design.desc,
-    'database': hql_design.query.get('database', None),
-    'settings': hql_design.settings,
-    'file_resources': hql_design.file_resources,
-    'functions': hql_design.functions,
-    'is_parameterized': hql_design.query.get('is_parameterized', True),
-    'email_notify': hql_design.query.get('email_notify', True),
-    'is_redacted': design.is_redacted
-  }
-
-
-def query_history_to_dict(request, query_history):
-  query_history_dict = {
-    'id': query_history.id,
-    'state': query_history.last_state,
-    'query': query_history.query,
-    'has_results': query_history.has_results,
-    'statement_number': query_history.statement_number,
-    'watch_url': reverse(get_app_name(request) + ':api_watch_query_refresh_json', kwargs={'id': query_history.id}),
-    'results_url': reverse(get_app_name(request) + ':view_results', kwargs={'id': query_history.id, 'first_row': 0})
-  }
-
-  if query_history.design:
-    query_history_dict['design'] = design_to_dict(query_history.design)
-
-  return query_history_dict
-
-
+@error_handler
 def clear_history(request):
   response = {'status': -1, 'message': ''}
 
@@ -655,21 +622,11 @@ def clear_history(request):
   return JsonResponse(response)
 
 
-# Proxy API for Metastore App
-def describe_table(request, database, table):
-  try:
-    from metastore.views import describe_table
-    return describe_table(request, database, table)
-  except Exception, e:
-    LOG.exception('Describe table failed')
-    raise PopupException(_('Problem accessing table metadata'), detail=e)
-
-
 @error_handler
 def get_sample_data(request, database, table):
   query_server = dbms.get_query_server_config(get_app_name(request))
   db = dbms.get(request.user, query_server)
-  response = {'status': -1, 'error_message': ''}
+  response = {'status': -1}
 
   table_obj = db.get_table(database, table)
   sample_data = db.get_sample(database, table_obj)
@@ -678,7 +635,7 @@ def get_sample_data(request, database, table):
     response['headers'] = sample_data.cols()
     response['rows'] = escape_rows(sample_data.rows(), nulls_only=True)
   else:
-    response['error_message'] = _('Sample data took too long to be generated')
+    response['message'] = _('Failed to get sample data.')
 
   return JsonResponse(response)
 
@@ -687,7 +644,7 @@ def get_sample_data(request, database, table):
 def get_indexes(request, database, table):
   query_server = dbms.get_query_server_config(get_app_name(request))
   db = dbms.get(request.user, query_server)
-  response = {'status': -1, 'error_message': ''}
+  response = {'status': -1}
 
   indexes = db.get_indexes(database, table)
   if indexes:
@@ -695,7 +652,23 @@ def get_indexes(request, database, table):
     response['headers'] = indexes.cols()
     response['rows'] = escape_rows(indexes.rows(), nulls_only=True)
   else:
-    response['error_message'] = _('Index data took too long to be generated')
+    response['message'] = _('Failed to get indexes.')
+
+  return JsonResponse(response)
+
+
+@error_handler
+def get_settings(request):
+  query_server = dbms.get_query_server_config(get_app_name(request))
+  db = dbms.get(request.user, query_server)
+  response = {'status': -1}
+
+  settings = db.get_configuration()
+  if settings:
+    response['status'] = 0
+    response['settings'] = settings
+  else:
+    response['message'] = _('Failed to get settings.')
 
   return JsonResponse(response)
 
@@ -704,7 +677,7 @@ def get_indexes(request, database, table):
 def get_functions(request):
   query_server = dbms.get_query_server_config(get_app_name(request))
   db = dbms.get(request.user, query_server)
-  response = {'status': -1, 'error_message': ''}
+  response = {'status': -1}
 
   prefix = request.GET.get('prefix', None)
   functions = db.get_functions(prefix)
@@ -713,37 +686,11 @@ def get_functions(request):
     rows = escape_rows(functions.rows(), nulls_only=True)
     response['functions'] = [row[0] for row in rows]
   else:
-    response['error_message'] = _('Fetching functions timed out.')
+    response['message'] = _('Failed to get functions.')
 
   return JsonResponse(response)
 
 
-def get_query_form(request):
-  try:
-    try:
-      # Get database choices
-      query_server = dbms.get_query_server_config(get_app_name(request))
-      db = dbms.get(request.user, query_server)
-      databases = [(database, database) for database in db.get_databases()]
-    except StructuredThriftTransportException, e:
-      # If Thrift exception was due to failed authentication, raise corresponding message
-      if 'TSocket read 0 bytes' in str(e) or 'Error validating the login' in str(e):
-        raise PopupException(_('Failed to authenticate to query server, check authentication configurations.'), detail=e)
-      else:
-        raise e
-  except Exception, e:
-    raise PopupException(_('Unable to access databases, Query Server or Metastore may be down.'), detail=e)
-
-  if not databases:
-    raise RuntimeError(_("No databases are available. Permissions could be missing."))
-
-  query_form = QueryForm()
-  query_form.bind(request.POST)
-  query_form.query.fields['database'].choices = databases # Could not do it in the form
-
-  return query_form
-
-
 @error_handler
 def analyze_table(request, database, table, columns=None):
   app_name = get_app_name(request)
@@ -855,6 +802,76 @@ def close_session(request, session_id):
   return JsonResponse(response)
 
 
+# Proxy API for Metastore App
+def describe_table(request, database, table):
+  try:
+    from metastore.views import describe_table
+    return describe_table(request, database, table)
+  except Exception, e:
+    LOG.exception('Describe table failed')
+    raise PopupException(_('Problem accessing table metadata'), detail=e)
+
+
+def design_to_dict(design):
+  hql_design = HQLdesign.loads(design.data)
+  return {
+    'id': design.id,
+    'query': hql_design.hql_query,
+    'name': design.name,
+    'desc': design.desc,
+    'database': hql_design.query.get('database', None),
+    'settings': hql_design.settings,
+    'file_resources': hql_design.file_resources,
+    'functions': hql_design.functions,
+    'is_parameterized': hql_design.query.get('is_parameterized', True),
+    'email_notify': hql_design.query.get('email_notify', True),
+    'is_redacted': design.is_redacted
+  }
+
+
+def query_history_to_dict(request, query_history):
+  query_history_dict = {
+    'id': query_history.id,
+    'state': query_history.last_state,
+    'query': query_history.query,
+    'has_results': query_history.has_results,
+    'statement_number': query_history.statement_number,
+    'watch_url': reverse(get_app_name(request) + ':api_watch_query_refresh_json', kwargs={'id': query_history.id}),
+    'results_url': reverse(get_app_name(request) + ':view_results', kwargs={'id': query_history.id, 'first_row': 0})
+  }
+
+  if query_history.design:
+    query_history_dict['design'] = design_to_dict(query_history.design)
+
+  return query_history_dict
+
+
+def get_query_form(request):
+  try:
+    try:
+      # Get database choices
+      query_server = dbms.get_query_server_config(get_app_name(request))
+      db = dbms.get(request.user, query_server)
+      databases = [(database, database) for database in db.get_databases()]
+    except StructuredThriftTransportException, e:
+      # If Thrift exception was due to failed authentication, raise corresponding message
+      if 'TSocket read 0 bytes' in str(e) or 'Error validating the login' in str(e):
+        raise PopupException(_('Failed to authenticate to query server, check authentication configurations.'), detail=e)
+      else:
+        raise e
+  except Exception, e:
+    raise PopupException(_('Unable to access databases, Query Server or Metastore may be down.'), detail=e)
+
+  if not databases:
+    raise RuntimeError(_("No databases are available. Permissions could be missing."))
+
+  query_form = QueryForm()
+  query_form.bind(request.POST)
+  query_form.query.fields['database'].choices = databases # Could not do it in the form
+
+  return query_form
+
+
 """
 Utils
 """

+ 4 - 0
apps/beeswax/src/beeswax/server/dbms.py

@@ -787,6 +787,10 @@ class HiveServer2Dbms(object):
     return result
 
 
+  def get_configuration(self):
+    return self.client.get_configuration()
+
+
   def get_functions(self, prefix=None):
     filter = '"%s.*"' % prefix if prefix else '".*"'
     hql = 'SHOW FUNCTIONS %s' % filter

+ 15 - 9
apps/beeswax/src/beeswax/server/hive_server2_lib.py

@@ -604,8 +604,8 @@ class HiveServerClient:
                                      properties=properties)
 
     # HS2 does not return properties in TOpenSessionResp
-    if self.query_server['server_name'] == "beeswax" and not session.get_properties():
-      session.properties = json.dumps(self.get_configuration(include_hadoop=False))
+    if not session.get_properties():
+      session.properties = json.dumps(self.get_configuration())
       session.save()
 
     return session
@@ -879,17 +879,20 @@ class HiveServerClient:
     return partitions[:max_parts]
 
 
-  def get_configuration(self, include_hadoop=False):
+  def get_configuration(self):
     configuration = {}
-    query = 'SET'
-    if include_hadoop:
-      query += ' -v'
 
-    results = self.execute_query_statement(query)
-    if results:
+    if self.query_server['server_name'] == 'impala':  # Return all configuration settings
+      query = 'SET'
+      results = self.execute_query_statement(query, orientation=TFetchOrientation.FETCH_NEXT)
+      configuration = dict((row[0], row[1]) for row in results.rows())
+    else:  # For Hive, only return white-listed configurations
+      query = 'SET -v'
+      results = self.execute_query_statement(query, orientation=TFetchOrientation.FETCH_FIRST)
       config_whitelist = [config.lower() for config in CONFIG_WHITELIST.get()]
       properties = [(row[0].split('=')[0], row[0].split('=')[1]) for row in results.rows() if '=' in row[0]]
       configuration = dict((prop, value) for prop, value in properties if prop.lower() in config_whitelist)
+
     return configuration
 
 
@@ -1108,7 +1111,7 @@ class HiveServerClientCompatible(object):
 
 
   def get_default_configuration(self, *args, **kwargs):
-    return {}
+    return []
 
 
   def get_results_metadata(self, handle):
@@ -1137,3 +1140,6 @@ class HiveServerClientCompatible(object):
 
 
   def alter_partition(self, db_name, tbl_name, new_part): raise NotImplementedError()
+
+  def get_configuration(self):
+    return self._client.get_configuration()

+ 22 - 0
apps/beeswax/src/beeswax/tests.py

@@ -1967,6 +1967,28 @@ for x in sys.stdin:
     assert_equal(2, len(json_resp['rows']), json_resp['rows'])
 
 
+  def test_get_settings(self):
+    resets = [
+      beeswax.conf.CONFIG_WHITELIST.set_for_testing('hive.execution.engine,mapreduce.job.queuename'),
+    ]
+
+    try:
+      resp = self.client.get(reverse("beeswax:get_settings"))
+      json_resp = json.loads(resp.content)
+      assert_equal(0, json_resp['status'])
+      assert_equal(2, len(json_resp['settings'].items()), json_resp)
+      assert_true('hive.execution.engine' in json_resp['settings'])
+      assert_true('mapreduce.job.queuename' in json_resp['settings'])
+
+      resp = self.client.get(reverse("impala:get_settings"))
+      json_resp = json.loads(resp.content)
+      assert_equal(0, json_resp['status'])
+      assert_true('QUERY_TIMEOUT_S' in json_resp['settings'])
+    finally:
+      for reset in resets:
+        reset()
+
+
   def test_get_functions(self):
     resp = self.client.get(reverse("beeswax:get_functions"))
     json_resp = json.loads(resp.content)

+ 1 - 0
apps/beeswax/src/beeswax/urls.py

@@ -60,6 +60,7 @@ urlpatterns += patterns(
   url(r'^api/session/?$', 'get_session', name='api_get_session'),
   url(r'^api/session/(?P<session_id>\d+)/?$', 'get_session', name='api_get_session'),
   url(r'^api/session/(?P<session_id>\d+)/close/?$', 'close_session', name='api_close_session'),
+  url(r'^api/settings/?$', 'get_settings', name='get_settings'),
   url(r'^api/functions/?$', 'get_functions', name='get_functions'),
 
   # Deprecated by Notebook API