فهرست منبع

HUE-1084 [impala] Integrate stats and terms analytics to Impala API

Romain Rigaux 10 سال پیش
والد
کامیت
7c7c1601dd

+ 55 - 0
apps/beeswax/src/beeswax/api.py

@@ -628,3 +628,58 @@ def get_query_form(request):
   query_form.query.fields['database'].choices = databases # Could not do it in the form
 
   return query_form
+
+
+def analyze_table(request, database, table, columns=None):
+  app_name = get_app_name(request)
+  query_server = get_query_server_config(app_name)
+  db = dbms.get(request.user, query_server)
+
+  response = {'status': -1, 'message': '', 'redirect': ''}
+
+  if request.method == "POST":
+    if columns is not None:
+      query_history = db.analyze_table(database, table)
+    else:
+      query_history = db.analyze_table_columns(database, table)
+
+    response['watch_url'] = reverse('beeswax:api_watch_query_refresh_json', kwargs={'id': query_history.id})
+    response['status'] = 0
+  else:
+    response['message'] = _('A POST request is required.')
+
+  return JsonResponse(response)
+
+
+def get_table_stats(request, database, table, column=None):
+  app_name = get_app_name(request)
+  query_server = get_query_server_config(app_name)
+  db = dbms.get(request.user, query_server)
+
+  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
+
+  response['stats'] = stats
+  response['status'] = 0
+
+  return JsonResponse(response)
+
+
+def get_top_terms(request, database, table, column, prefix=None):
+  app_name = get_app_name(request)
+  query_server = get_query_server_config(app_name)
+  db = dbms.get(request.user, query_server)
+
+  response = {'status': -1, 'message': '', 'redirect': ''}
+
+  terms = db.get_top_terms(database, table, column, prefix=prefix, limit=int(request.GET.get('limit', 30)))
+
+  response['terms'] = terms
+  response['status'] = 0
+
+  return JsonResponse(response)

+ 59 - 6
apps/beeswax/src/beeswax/server/dbms.py

@@ -200,19 +200,48 @@ class HiveServer2Dbms(object):
 
 
   def analyze_table(self, database, table):
-    hql = 'ANALYZE TABLE `%(database)s`.`%(table)s` COMPUTE STATISTICS' % {'database': database, 'table': table}
+    if self.server_name == 'impala':
+      hql = 'COMPUTE STATS `%(database)s`.`%(table)s`' % {'database': database, 'table': table}
+    else:
+      hql = 'ANALYZE TABLE `%(database)s`.`%(table)s` COMPUTE STATISTICS' % {'database': database, 'table': table}
 
     return self.execute_statement(hql)
 
 
   def analyze_table_columns(self, database, table):
-    hql = 'ANALYZE TABLE `%(database)s`.`%(table)s` COMPUTE STATISTICS FOR COLUMNS' % {'database': database, 'table': table}
+    if self.server_name == 'impala':
+      hql = 'COMPUTE STATS `%(database)s`.`%(table)s`' % {'database': database, 'table': table}
+    else:
+      hql = 'ANALYZE TABLE `%(database)s`.`%(table)s` COMPUTE STATISTICS FOR COLUMNS' % {'database': database, 'table': table}
 
     return self.execute_statement(hql)
 
 
+  def get_table_stats(self, database, table):
+    stats = []
+
+    if self.server_name == 'impala':
+      hql = 'SHOW TABLE STATS `%(database)s`.`%(table)s`' % {'database': database, 'table': table}
+
+      query = hql_query(hql)
+      handle = self.execute_and_wait(query, timeout_sec=5.0)
+
+      if handle:
+        result = self.fetch(handle, rows=100)
+        self.close(handle)
+        stats = list(result.rows())
+    else:
+      table = self.get_table(database, table)
+      stats = table.stats
+
+    return stats
+
+
   def get_table_columns_stats(self, database, table, column):
-    hql = 'DESCRIBE FORMATTED `%(database)s`.`%(table)s` %(column)s' % {'database': database, 'table': table, 'column': column}
+    if self.server_name == 'impala':
+      hql = 'SHOW COLUMN STATS `%(database)s`.`%(table)s`' % {'database': database, 'table': table}
+    else:
+      hql = 'DESCRIBE FORMATTED `%(database)s`.`%(table)s` %(column)s' % {'database': database, 'table': table, 'column': column}
 
     query = hql_query(hql)
     handle = self.execute_and_wait(query, timeout_sec=5.0)
@@ -220,16 +249,40 @@ class HiveServer2Dbms(object):
     if handle:
       result = self.fetch(handle, rows=100)
       self.close(handle)
-      return list(result.rows())
+      data = list(result.rows())
+
+      if self.server_name == 'impala':
+        data = [col for col in data if col[0] == column][0]
+        return [
+            {'col_name': data[0]},
+            {'data_type': data[1]},
+            {'distinct_count': data[2]},
+            {'num_nulls': data[3]},
+            {'max_col_len': data[4]},
+            {'avg_col_len': data[5]},
+        ]
+      else:
+        return [
+            {'col_name': data[2][0]},
+            {'data_type': data[2][1]},
+            {'min': data[2][2]},
+            {'max': data[2][3]},
+            {'num_nulls': data[2][4]},
+            {'distinct_count': data[2][5]},
+            {'avg_col_len': data[2][6]},
+            {'max_col_len': data[2][7]},
+            {'num_trues': data[2][8]},
+            {'num_falses': data[2][9]}
+        ]
     else:
       return []
 
 
   def get_top_terms(self, database, table, column, limit=30, prefix=None):
-    limit = max(limit, 100)
+    limit = min(limit, 100)
     prefix_match = ''
     if prefix:
-      prefix_match = "WHERE %(column)s LIKE '%(prefix)s%%'" % {'column': column, 'prefix': prefix}
+      prefix_match = "WHERE CAST(%(column)s AS STRING) LIKE '%(prefix)s%%'" % {'column': column, 'prefix': prefix}
 
     hql = 'SELECT %(column)s, COUNT(*) AS ct FROM `%(database)s`.`%(table)s` %(prefix_match)s GROUP BY %(column)s ORDER BY ct DESC LIMIT %(limit)s' % {
         'database': database, 'table': table, 'column': column, 'prefix_match': prefix_match, 'limit': limit,

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

@@ -1567,6 +1567,7 @@ for x in sys.stdin:
     assert_true('<th>foo</th>' in resp.content, resp.content)
     assert_true([0, '0x0'] in resp.context['sample'], resp.context['sample'])
 
+
   def test_redacting_queries(self):
     c = make_logged_in_client()
 
@@ -1605,6 +1606,77 @@ for x in sys.stdin:
       redaction.global_redaction_engine.policies = old_policies
 
 
+  def test_analyze_table_and_read_statistics(self):
+    # No stats
+    resp = self.client.get(reverse('beeswax:get_table_stats', kwargs={'database': 'default', 'table': 'test'}))
+    stats = json.loads(resp.content)['stats']
+    assert_equal('COLUMN_STATS_ACCURATE', stats[0]['data_type'], resp.content)
+
+    resp = self.client.get(reverse('beeswax:get_table_stats', kwargs={'database': 'default', 'table': 'test', 'column': 'foo'}))
+    stats = json.loads(resp.content)['stats']
+    assert_equal([
+          {u'col_name': u'foo'},
+          {u'data_type': u'int'},
+          {u'min': u''},
+          {u'max': u''},
+          {u'num_nulls': u''},
+          {u'distinct_count': u''},
+          {u'avg_col_len': u''},
+          {u'max_col_len': u''},
+          {u'num_trues': u''},
+          {u'num_falses': u''}
+        ],
+        stats
+    )
+
+    # Compute stats
+    response = self.client.post(reverse("beeswax:analyze_table", kwargs={'database': 'default', 'table': 'test'}), follow=True)
+    response = wait_for_query_to_finish(self.client, response, max=60.0)
+    assert_true(response, response)
+
+    response = self.client.post(reverse("beeswax:analyze_table", kwargs={'database': 'default', 'table': 'test', 'columns': True}), follow=True)
+    response = wait_for_query_to_finish(self.client, response, max=60.0)
+    assert_true(response, response)
+
+    # Retrieve stats
+    resp = self.client.get(reverse('beeswax:get_table_stats', kwargs={'database': 'default', 'table': 'test'}))
+    stats = json.loads(resp.content)['stats']
+    assert_true(any([stat for stat in stats if stat['data_type'] == 'numRows']), resp.content)
+    assert_true(any([stat for stat in stats if stat['comment'] == '256']), resp.content)
+
+    resp = self.client.get(reverse('beeswax:get_table_stats', kwargs={'database': 'default', 'table': 'test', 'column': 'foo'}))
+    stats = json.loads(resp.content)['stats']
+    assert_equal([
+          {u'col_name': u'foo'},
+          {u'data_type': u'int'},
+          {u'min': u'0'},
+          {u'max': u'255'},
+          {u'num_nulls': u'0'},
+          {u'distinct_count': u'180'},          
+          {u'avg_col_len': u''},
+          {u'max_col_len': u''},
+          {u'num_trues': u''},
+          {u'num_falses': u''}
+        ],
+        stats
+    )
+
+
+  def test_get_top_terms(self):
+    resp = self.client.get(reverse("beeswax:get_top_terms", kwargs={'database': 'default', 'table': 'test', 'column': 'foo'}))
+    terms = json.loads(resp.content)['terms']
+    assert_equal([[255, 1], [254, 1], [253, 1], [252, 1]], terms[:4])
+
+    resp = self.client.get(reverse("beeswax:get_top_terms", kwargs={'database': 'default', 'table': 'test', 'column': 'foo', 'prefix': '10'}))
+    terms = json.loads(resp.content)['terms']
+    assert_equal([[109, 1], [108, 1], [107, 1], [106, 1]], terms[:4])
+
+    resp = self.client.get(reverse("beeswax:get_top_terms", kwargs={'database': 'default', 'table': 'test', 'column': 'foo', 'prefix': '10'}) + '?limit=2')
+    terms = json.loads(resp.content)['terms']
+    assert_equal([[109, 1], [108, 1]], terms)
+
+
+
 def test_import_gzip_reader():
   """Test the gzip reader in create table"""
   # Make gzipped data

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

@@ -75,4 +75,7 @@ urlpatterns += patterns(
   url(r'^api/watch/json/(?P<id>\d+)$', 'watch_query_refresh_json', name='api_watch_query_refresh_json'),
 
   url(r'^api/table/(?P<database>\w+)/(?P<table>\w+)$', 'describe_table', name='describe_table'),
+  url(r'^api/analyze/(?P<database>\w+)/(?P<table>\w+)/(?P<columns>\w+)?$', 'analyze_table', name='analyze_table'),
+  url(r'^api/table/(?P<database>\w+)/(?P<table>\w+)/stats/(?P<column>\w+)?$', 'get_table_stats', name='get_table_stats'),
+  url(r'^api/table/(?P<database>\w+)/(?P<table>\w+)/terms/(?P<column>\w+)/(?P<prefix>\w+)?$', 'get_top_terms', name='get_top_terms'),
 )

+ 0 - 40
apps/metastore/src/metastore/tests.py

@@ -249,43 +249,3 @@ class TestMetastoreWithHadoop(BeeswaxSampleProvider):
     GroupPermission.objects.get_or_create(group=group, hue_permission=perm)
 
     check(client, [200, 302]) # Ok
-
-
-  def test_analyze_table_and_read_statistics(self):
-    # No stats
-    resp = self.client.get(reverse('metastore:get_table_stats', kwargs={'database': 'default', 'table': 'test'}))
-    stats = json.loads(resp.content)['stats']
-    assert_equal('COLUMN_STATS_ACCURATE', stats[0]['data_type'], resp.content)
-
-    resp = self.client.get(reverse('metastore:get_table_stats', kwargs={'database': 'default', 'table': 'test', 'column': 'foo'}))
-    stats = json.loads(resp.content)['stats']
-    assert_equal(["foo", "int", "", "", "", "", "", "", "", "", "from deserializer"], stats[2])
-
-    # Compute stats
-    response = self.client.post(reverse("metastore:analyze_table", kwargs={'database': 'default', 'table': 'test'}), follow=True)
-    response = wait_for_query_to_finish(self.client, response, max=60.0)
-    assert_true(response, response)
-
-    response = self.client.post(reverse("metastore:analyze_table", kwargs={'database': 'default', 'table': 'test', 'columns': True}), follow=True)
-    response = wait_for_query_to_finish(self.client, response, max=60.0)
-    assert_true(response, response)
-
-    # Retrieve stats
-    resp = self.client.get(reverse('metastore:get_table_stats', kwargs={'database': 'default', 'table': 'test'}))
-    stats = json.loads(resp.content)['stats']
-    assert_true(any([stat for stat in stats if stat['data_type'] == 'numRows']), resp.content)
-    assert_true(any([stat for stat in stats if stat['comment'] == '256']), resp.content)
-
-    resp = self.client.get(reverse('metastore:get_table_stats', kwargs={'database': 'default', 'table': 'test', 'column': 'foo'}))
-    stats = json.loads(resp.content)['stats']
-    assert_equal(["foo", "int", "0", "255", "0", "180", "", "", "", "", "from deserializer"], stats[2])
-
-
-  def test_get_top_terms(self):
-    resp = self.client.get(reverse("metastore:get_top_terms", kwargs={'database': 'default', 'table': 'test', 'column': 'foo'}))
-    terms = json.loads(resp.content)['terms']
-    assert_equal([[255, 1], [254, 1], [253, 1], [252, 1]], terms[:4])
-
-    resp = self.client.get(reverse("metastore:get_top_terms", kwargs={'database': 'default', 'table': 'test', 'column': 'foo', 'prefix': '10'}))
-    terms = json.loads(resp.content)['terms']
-    assert_equal([[109, 1], [108, 1], [107, 1], [106, 1]], terms[:4])

+ 0 - 5
apps/metastore/src/metastore/urls.py

@@ -30,9 +30,4 @@ urlpatterns = patterns('metastore.views',
   url(r'^table/(?P<database>\w+)/(?P<table>\w+)/load$', 'load_table', name='load_table'),
   url(r'^table/(?P<database>\w+)/(?P<table>\w+)/read$', 'read_table', name='read_table'),
   url(r'^table/(?P<database>\w+)/(?P<table>\w+)/partitions/(?P<partition_id>\w+)$', 'read_partition', name='read_partition'),
-
-  # API
-  url(r'^analyze/(?P<database>\w+)/(?P<table>\w+)/(?P<columns>\w+)?$', 'analyze_table', name='analyze_table'),
-  url(r'^table/(?P<database>\w+)/(?P<table>\w+)/stats/(?P<column>\w+)?$', 'get_table_stats', name='get_table_stats'),
-  url(r'^table/(?P<database>\w+)/(?P<table>\w+)/terms/(?P<column>\w+)/(?P<prefix>\w+)?$', 'get_top_terms', name='get_top_terms'),
 )

+ 0 - 55
apps/metastore/src/metastore/views.py

@@ -287,60 +287,5 @@ def describe_partitions(request, database, table):
   })
 
 
-def analyze_table(request, database, table, columns=None):
-  app_name = get_app_name(request)
-  query_server = get_query_server_config(app_name)
-  db = dbms.get(request.user, query_server)
-
-  response = {'status': -1, 'message': '', 'redirect': ''}
-
-  if request.method == "POST":
-    if columns is not None:
-      query_history = db.analyze_table(database, table)
-    else:
-      query_history = db.analyze_table_columns(database, table)
-
-    response['watch_url'] = reverse('beeswax:api_watch_query_refresh_json', kwargs={'id': query_history.id})
-    response['status'] = 0
-  else:
-    response['message'] = _('A POST request is required.')
-
-  return JsonResponse(response)
-
-
-def get_table_stats(request, database, table, column=None):
-  app_name = get_app_name(request)
-  query_server = get_query_server_config(app_name)
-  db = dbms.get(request.user, query_server)
-
-  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
-
-  response['stats'] = stats
-  response['status'] = 0
-
-  return JsonResponse(response)
-
-
-def get_top_terms(request, database, table, column, prefix=None):
-  app_name = get_app_name(request)
-  query_server = get_query_server_config(app_name)
-  db = dbms.get(request.user, query_server)
-
-  response = {'status': -1, 'message': '', 'redirect': ''}
-
-  terms = db.get_top_terms(database, table, column, prefix=prefix)
-
-  response['terms'] = terms
-  response['status'] = 0
-
-  return JsonResponse(response)
-
-
 def has_write_access(user):
   return user.is_superuser or user.has_hue_permission(action="write", app=DJANGO_APPS[0])