Просмотр исходного кода

HUE-2284 [impala] Provide Impala Thrift API implementation of GetExecSummary

POST /impala/api/query/<query_history_id>/exec_summary

query_history_id must reference an open HS2 session/operation
Returns a JSON response (for query='SELECT AVG(salary) as avg_salary FROM default.sample_08 ORDER by avg_salary DESC LIMIT 20'):

{
	"status": 0,
	"summary": {
		"status": null,
		"error_logs": null,
		"state": 0,
		"progress": null,
		"nodes": [{
			"exec_stats": [{
				"memory_used": 0,
				"latency_ns": 259669,
				"cardinality": 0,
				"cpu_time_ns": null
			}],
			"label_detail": "",
			"is_active": null,
			"estimated_stats": {
				"memory_used": -1,
				"latency_ns": null,
				"cardinality": 0,
				"cpu_time_ns": null
			},
			"label": "02:TOP-N",
			"is_broadcast": null,
			"num_children": 1,
			"node_id": 2,
			"fragment_id": 0
		}, {
			"exec_stats": [{
				"memory_used": 8192,
				"latency_ns": 348817669,
				"cardinality": 0,
				"cpu_time_ns": null
			}],
			"label_detail": "FINALIZE",
			"is_active": null,
			"estimated_stats": {
				"memory_used": -1,
				"latency_ns": null,
				"cardinality": 0,
				"cpu_time_ns": null
			},
			"label": "04:AGGREGATE",
			"is_broadcast": null,
			"num_children": 1,
			"node_id": 4,
			"fragment_id": 0
		}, {
			"exec_stats": [{
				"memory_used": 0,
				"latency_ns": 0,
				"cardinality": 0,
				"cpu_time_ns": null
			}],
			"label_detail": "UNPARTITIONED",
			"is_active": null,
			"estimated_stats": {
				"memory_used": -1,
				"latency_ns": null,
				"cardinality": 0,
				"cpu_time_ns": null
			},
			"label": "03:EXCHANGE",
			"is_broadcast": true,
			"num_children": 0,
			"node_id": 3,
			"fragment_id": 0
		}, {
			"exec_stats": null,
			"label_detail": "",
			"is_active": null,
			"estimated_stats": {
				"memory_used": 10485760,
				"latency_ns": null,
				"cardinality": 0,
				"cpu_time_ns": null
			},
			"label": "01:AGGREGATE",
			"is_broadcast": null,
			"num_children": 1,
			"node_id": 1,
			"fragment_id": 1
		}, {
			"exec_stats": null,
			"label_detail": "default.sample_08",
			"is_active": null,
			"estimated_stats": {
				"memory_used": 33554432,
				"latency_ns": null,
				"cardinality": 0,
				"cpu_time_ns": null
			},
			"label": "00:SCAN HDFS",
			"is_broadcast": null,
			"num_children": 0,
			"node_id": 0,
			"fragment_id": 1
		}],
		"exch_to_sender_map": {
			"2": 3
		}
	}
}
Jenny Kim 10 лет назад
Родитель
Сommit
7494aa46c4

+ 25 - 1
apps/impala/src/impala/api.py

@@ -15,7 +15,6 @@
 # See the License for the specific language governing permissions and
 # limitations under the License.
 
-
 ## Main views are inherited from Beeswax.
 
 import logging
@@ -24,9 +23,12 @@ from django.utils.translation import ugettext as _
 from django.views.decorators.http import require_POST
 
 from desktop.lib.django_util import JsonResponse
+from django.views.decorators.http import require_POST
 
 from beeswax.api import error_handler
+from beeswax.models import Session
 from beeswax.server import dbms as beeswax_dbms
+from beeswax.views import authorized_get_query_history
 
 from impala import dbms
 
@@ -64,3 +66,25 @@ def refresh_table(request, database, table):
   response['message'] = _('Successfully refreshed metadata for `%s`.`%s`') % (database, table)
 
   return JsonResponse(response)
+
+
+@require_POST
+@error_handler
+def get_exec_summary(request, query_history_id):
+  query_server = dbms.get_query_server_config()
+  db = beeswax_dbms.get(request.user, query_server=query_server)
+
+  response = {'status': -1}
+  query_history = authorized_get_query_history(request, query_history_id, must_exist=True)
+
+  if query_history is None:
+    response['message'] = _('get_exec_summary requires a valid query_history_id')
+  else:
+    session = Session.objects.get_session(request.user, query_server['server_name'])
+    operation_handle = query_history.get_handle().get_rpc_handle()
+    session_handle = session.get_handle()
+    summary = db.get_exec_summary(operation_handle, session_handle)
+    response['status'] = 0
+    response['summary'] = summary
+
+  return JsonResponse(response)

+ 4 - 0
apps/impala/src/impala/dbms.py

@@ -151,6 +151,10 @@ class ImpalaDbms(HiveServer2Dbms):
     return results
 
 
+  def get_exec_summary(self, query_handle, session_handle):
+    return self.client._client.get_exec_summary(query_handle, session_handle)
+
+
   def _get_beeswax_tables(self, database):
     beeswax_query_server = dbms.get(user=self.client.user, query_server=beeswax_query_server_config(name='beeswax'))
     return beeswax_query_server.get_tables(database=database)

+ 63 - 3
apps/impala/src/impala/server.py

@@ -17,13 +17,73 @@
 
 import logging
 
-from beeswax.server.hive_server2_lib import HiveServerClient
+from beeswax.server.dbms import QueryServerException
+from beeswax.server.hive_server2_lib import HiveServerClient, HiveServerDataTable
+from TCLIService.ttypes import TExecuteStatementReq, TFetchOrientation, TStatusCode
+
+from ImpalaService import ImpalaHiveServer2Service
 
 
 LOG = logging.getLogger(__name__)
 
 
+class ImpalaServerClientException(Exception):
+  pass
+
+
 class ImpalaServerClient(HiveServerClient):
 
-  def resetCatalog(self):
-    return self._client.ResetCatalog()
+  def get_exec_summary(self, operation_handle, session_handle):
+    """
+    Calls Impala HS2 API's GetExecSummary method on the given query handle
+    :return: TExecSummary object serialized as a dict
+    """
+    req = ImpalaHiveServer2Service.TGetExecSummaryReq(operationHandle=operation_handle, sessionHandle=session_handle)
+
+    # GetExecSummary() only works for closed queries
+    self.close_operation(operation_handle)
+
+    resp = self.call(self._client.GetExecSummary, req)
+
+    if resp.status is not None and resp.status.statusCode not in (TStatusCode.SUCCESS_STATUS,):
+      if hasattr(resp.status, 'errorMessage') and resp.status.errorMessage:
+        message = resp.status.errorMessage
+      else:
+        message = ''
+      raise QueryServerException(Exception('Bad status for request %s:\n%s' % (req, resp)), message=message)
+
+    return self._serialize_exec_summary(resp.summary)
+
+
+  def _serialize_exec_summary(self, summary):
+    try:
+      summary_dict = {
+        'state': summary.state,
+        'exch_to_sender_map': summary.exch_to_sender_map,
+        'error_logs': summary.error_logs,
+        'status': None,
+        'progress': None,
+        'nodes': [],
+      }
+
+      if summary.status is not None:
+        summary_dict['status'] = summary.status.__dict__
+
+      if summary.progress is not None:
+        summary_dict['progress'] = summary.progress.__dict__
+
+      if summary.nodes:
+        for node in summary.nodes:
+          node_dict = node.__dict__
+
+          if node.exec_stats is not None:
+            node_dict['exec_stats'] = [stat.__dict__ for stat in node.exec_stats]
+
+          if node.estimated_stats is not None:
+            node_dict['estimated_stats'] = node.estimated_stats.__dict__
+
+          summary_dict['nodes'].append(node_dict)
+
+      return summary_dict
+    except Exception, e:
+      raise ImpalaServerClientException('Failed to serialize the TExecSummary object: %s' % str(e))

+ 18 - 0
apps/impala/src/impala/tests.py

@@ -297,6 +297,24 @@ class TestImpalaIntegration:
       "\ntest_refresh_table: `%s`.`%s`\nImpala Columns: %s\nBeeswax Columns: %s" % (self.DATABASE, 'tweets', ','.join(impala_columns), ','.join(beeswax_columns)))
 
 
+  def test_get_exec_summary(self):
+    query = """
+      SELECT COUNT(1) FROM tweets;
+    """
+
+    response = _make_query(self.client, query, database=self.DATABASE, local=False, server_name='impala')
+    content = json.loads(response.content)
+    query_history = QueryHistory.get(content['id'])
+
+    wait_for_query_to_finish(self.client, response, max=180.0)
+
+    resp = self.client.post(reverse('impala:get_exec_summary', kwargs={'query_history_id': query_history.id}))
+    data = json.loads(resp.content)
+    assert_equal(0, data['status'], data)
+    assert_true('nodes' in data['summary'], data)
+    assert_true(len(data['summary']['nodes']) > 0, data['summary']['nodes'])
+
+
 # Could be refactored with SavedQuery.create_empty()
 def create_saved_query(app_name, owner):
     query_type = SavedQuery.TYPES_MAPPING[app_name]

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

@@ -23,6 +23,7 @@ from beeswax.urls import urlpatterns as beeswax_urls
 urlpatterns = patterns('impala.api',
   url(r'^api/invalidate$', 'invalidate', name='invalidate'),
   url(r'^api/refresh/(?P<database>\w+)/(?P<table>\w+)$', 'refresh_table', name='refresh_table'),
+  url(r'^api/query/(?P<query_history_id>\d+)/exec_summary$', 'get_exec_summary', name='get_exec_summary'),
 )
 
 urlpatterns += patterns('impala.dashboards',