optimizer_api.py 14 KB


  1. #!/usr/bin/env python
  2. # Licensed to Cloudera, Inc. under one
  3. # or more contributor license agreements. See the NOTICE file
  4. # distributed with this work for additional information
  5. # regarding copyright ownership. Cloudera, Inc. licenses this file
  6. # to you under the Apache License, Version 2.0 (the
  7. # "License"); you may not use this file except in compliance
  8. # with the License. You may obtain a copy of the License at
  9. #
  10. # http://www.apache.org/licenses/LICENSE-2.0
  11. #
  12. # Unless required by applicable law or agreed to in writing, software
  13. # distributed under the License is distributed on an "AS IS" BASIS,
  14. # WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
  15. # See the License for the specific language governing permissions and
  16. # limitations under the License.
  17. import base64
  18. import json
  19. import logging
  20. import struct
  21. from django.http import Http404
  22. from django.views.decorators.http import require_POST
  23. from desktop.lib.django_util import JsonResponse
  24. from desktop.lib.i18n import force_unicode
  25. from desktop.models import Document2
  26. from libsentry.privilege_checker import MissingSentryPrivilegeException
  27. from notebook.api import _get_statement
  28. from notebook.models import Notebook
  29. from metadata.optimizer_client import OptimizerApi, NavOptException, _get_table_name, _clean_query
  30. from metadata.conf import OPTIMIZER
  31. LOG = logging.getLogger(__name__)
  32. try:
  33. from beeswax.api import get_table_stats
  34. from beeswax.design import hql_query
  35. from beeswax.server import dbms
  36. except ImportError, e:
  37. LOG.warn("Hive lib not enabled")
  38. def error_handler(view_fn):
  39. def decorator(*args, **kwargs):
  40. try:
  41. return view_fn(*args, **kwargs)
  42. except Http404, e:
  43. raise e
  44. except NavOptException, e:
  45. LOG.exception(e)
  46. response = {
  47. 'status': -1,
  48. 'message': e.message
  49. }
  50. except MissingSentryPrivilegeException, e:
  51. LOG.exception(e)
  52. response = {
  53. 'status': -1,
  54. 'message': 'Missing privileges for %s' % force_unicode(str(e))
  55. }
  56. except Exception, e:
  57. LOG.exception(e)
  58. response = {
  59. 'status': -1,
  60. 'message': force_unicode(str(e))
  61. }
  62. return JsonResponse(response, status=500)
  63. return decorator
  64. @require_POST
  65. @error_handler
  66. def get_tenant(request):
  67. response = {'status': -1}
  68. cluster_id = request.POST.get('cluster_id')
  69. api = OptimizerApi(request.user)
  70. data = api.get_tenant(cluster_id=cluster_id)
  71. if data:
  72. response['status'] = 0
  73. response['data'] = data['tenant']
  74. else:
  75. response['message'] = 'Optimizer: %s' % data['details']
  76. return JsonResponse(response)
  77. @require_POST
  78. @error_handler
  79. def top_tables(request):
  80. response = {'status': -1}
  81. database = request.POST.get('database', 'default')
  82. limit = request.POST.get('len', 1000)
  83. api = OptimizerApi(user=request.user)
  84. data = api.top_tables(database_name=database, page_size=limit)
  85. tables = [{
  86. 'eid': table['eid'],
  87. 'database': _get_table_name(table['name'])['database'],
  88. 'name': _get_table_name(table['name'])['table'],
  89. 'popularity': table['workloadPercent'],
  90. 'column_count': table['columnCount'],
  91. 'patternCount': table['patternCount'],
  92. 'total': table['total'],
  93. 'is_fact': table['type'] != 'Dimension'
  94. } for table in data['results']
  95. ]
  96. response['top_tables'] = tables
  97. response['status'] = 0
  98. return JsonResponse(response)
  99. @require_POST
  100. @error_handler
  101. def table_details(request):
  102. response = {'status': -1}
  103. database_name = request.POST.get('databaseName')
  104. table_name = request.POST.get('tableName')
  105. api = OptimizerApi(request.user)
  106. data = api.table_details(database_name=database_name, table_name=table_name)
  107. if data:
  108. response['status'] = 0
  109. response['details'] = data
  110. else:
  111. response['message'] = 'Optimizer: %s' % data['details']
  112. return JsonResponse(response)
  113. @require_POST
  114. @error_handler
  115. def query_compatibility(request):
  116. response = {'status': -1}
  117. source_platform = request.POST.get('sourcePlatform')
  118. target_platform = request.POST.get('targetPlatform')
  119. query = request.POST.get('query')
  120. api = OptimizerApi(request.user)
  121. data = api.query_compatibility(source_platform=source_platform, target_platform=target_platform, query=query)
  122. if data:
  123. response['status'] = 0
  124. response['query_compatibility'] = data
  125. else:
  126. response['message'] = 'Optimizer: %s' % data
  127. return JsonResponse(response)
  128. @require_POST
  129. @error_handler
  130. def query_risk(request):
  131. response = {'status': -1}
  132. query = json.loads(request.POST.get('query'))
  133. source_platform = request.POST.get('sourcePlatform')
  134. db_name = request.POST.get('dbName')
  135. api = OptimizerApi(request.user)
  136. data = api.query_risk(query=query, source_platform=source_platform, db_name=db_name)
  137. if data:
  138. response['status'] = 0
  139. response['query_risk'] = data
  140. else:
  141. response['message'] = 'Optimizer: %s' % data
  142. return JsonResponse(response)
  143. @require_POST
  144. @error_handler
  145. def similar_queries(request):
  146. response = {'status': -1}
  147. source_platform = request.POST.get('sourcePlatform')
  148. query = json.loads(request.POST.get('query'))
  149. api = OptimizerApi(request.user)
  150. data = api.similar_queries(source_platform=source_platform, query=query)
  151. if data:
  152. response['status'] = 0
  153. response['similar_queries'] = data
  154. else:
  155. response['message'] = 'Optimizer: %s' % data
  156. return JsonResponse(response)
  157. @require_POST
  158. @error_handler
  159. def top_filters(request):
  160. response = {'status': -1}
  161. db_tables = json.loads(request.POST.get('dbTables'), '[]')
  162. column_name = request.POST.get('columnName') # Unused
  163. api = OptimizerApi(request.user)
  164. data = api.top_filters(db_tables=db_tables)
  165. if data:
  166. response['status'] = 0
  167. response['values'] = data['results']
  168. else:
  169. response['message'] = 'Optimizer: %s' % data
  170. return JsonResponse(response)
  171. @require_POST
  172. @error_handler
  173. def top_joins(request):
  174. response = {'status': -1}
  175. db_tables = json.loads(request.POST.get('dbTables'), '[]')
  176. api = OptimizerApi(request.user)
  177. data = api.top_joins(db_tables=db_tables)
  178. if data:
  179. response['status'] = 0
  180. response['values'] = data['results']
  181. else:
  182. response['message'] = 'Optimizer: %s' % data
  183. return JsonResponse(response)
  184. @require_POST
  185. @error_handler
  186. def top_aggs(request):
  187. response = {'status': -1}
  188. db_tables = json.loads(request.POST.get('dbTables'), '[]')
  189. api = OptimizerApi(request.user)
  190. data = api.top_aggs(db_tables=db_tables)
  191. if data:
  192. response['status'] = 0
  193. response['values'] = data['results']
  194. else:
  195. response['message'] = 'Optimizer: %s' % data
  196. return JsonResponse(response)
  197. @require_POST
  198. @error_handler
  199. def top_databases(request):
  200. response = {'status': -1}
  201. api = OptimizerApi(request.user)
  202. data = api.top_databases()
  203. if data:
  204. response['status'] = 0
  205. response['values'] = data['results']
  206. else:
  207. response['message'] = 'Optimizer: %s' % data
  208. return JsonResponse(response)
  209. @require_POST
  210. @error_handler
  211. def top_columns(request):
  212. response = {'status': -1}
  213. db_tables = json.loads(request.POST.get('dbTables'), '[]')
  214. api = OptimizerApi(request.user)
  215. data = api.top_columns(db_tables=db_tables)
  216. if data:
  217. response['status'] = 0
  218. response['values'] = data
  219. else:
  220. response['message'] = 'Optimizer: %s' % data
  221. return JsonResponse(response)
  222. def _convert_queries(queries_data):
  223. queries = []
  224. for query_data in queries_data:
  225. try:
  226. snippet = query_data['snippets'][0]
  227. if 'guid' in snippet['result']['handle']: # Not failed query
  228. original_query_id = '%s:%s' % struct.unpack(b"QQ", base64.decodestring(snippet['result']['handle']['guid']))
  229. execution_time = snippet['result']['executionTime'] * 100 if snippet['status'] in ('available', 'expired') else -1
  230. statement = _clean_query(_get_statement(query_data))
  231. queries.append((original_query_id, execution_time, statement, snippet.get('database', 'default').strip()))
  232. except Exception, e:
  233. LOG.warning('Skipping upload of %s: %s' % (query_data['uuid'], e))
  234. return queries
  235. @require_POST
  236. @error_handler
  237. def upload_history(request):
  238. response = {'status': -1}
  239. if request.user.is_superuser:
  240. api = OptimizerApi(request.user)
  241. histories = []
  242. upload_stats = {}
  243. if request.POST.get('sourcePlatform'):
  244. n = min(request.POST.get('n', OPTIMIZER.QUERY_HISTORY_UPLOAD_LIMIT.get()))
  245. source_platform = request.POST.get('sourcePlatform', 'hive')
  246. histories = [(source_platform, Document2.objects.get_history(doc_type='query-%s' % source_platform, user=request.user)[:n])]
  247. elif OPTIMIZER.QUERY_HISTORY_UPLOAD_LIMIT.get() > 0:
  248. histories = [
  249. (source_platform, Document2.objects.filter(type='query-%s' % source_platform, is_history=True, is_managed=False, is_trashed=False).order_by('-last_modified')[:OPTIMIZER.QUERY_HISTORY_UPLOAD_LIMIT.get()])
  250. for source_platform in ['hive', 'impala']
  251. ]
  252. for source_platform, history in histories:
  253. queries = _convert_queries([Notebook(document=doc).get_data() for doc in history])
  254. upload_stats[source_platform] = api.upload(data=queries, data_type='queries', source_platform=source_platform)
  255. response['upload_history'] = upload_stats
  256. response['status'] = 0
  257. else:
  258. response['message'] = _('Query history upload requires Admin privileges or feature is disabled.')
  259. return JsonResponse(response)
  260. @require_POST
  261. @error_handler
  262. def upload_query(request):
  263. response = {'status': -1}
  264. if OPTIMIZER.AUTO_UPLOAD_QUERIES.get():
  265. query_id = request.POST.get('query_id')
  266. doc = Document2.objects.document(request.user, doc_id=query_id)
  267. query_data = Notebook(document=doc).get_data()
  268. queries = _convert_queries([query_data])
  269. source_platform = query_data['snippets'][0]['type']
  270. api = OptimizerApi(request.user)
  271. response['query_upload'] = api.upload(data=queries, data_type='queries', source_platform=source_platform)
  272. else:
  273. response['query_upload'] = _('Skipped')
  274. response['status'] = 0
  275. return JsonResponse(response)
  276. @require_POST
  277. @error_handler
  278. def upload_table_stats(request):
  279. response = {'status': -1}
  280. db_tables = json.loads(request.POST.get('db_tables'), '[]')
  281. source_platform = json.loads(request.POST.get('sourcePlatform', '"hive"'))
  282. with_columns = json.loads(request.POST.get('with_columns', 'false'))
  283. with_ddl = json.loads(request.POST.get('with_ddl', 'false'))
  284. table_stats = []
  285. column_stats = []
  286. table_ddls = []
  287. for db_table in db_tables:
  288. path = _get_table_name(db_table)
  289. try:
  290. if with_ddl:
  291. db = dbms.get(request.user)
  292. query = hql_query('SHOW CREATE TABLE `%(database)s`.`%(table)s`' % path)
  293. handle = db.execute_and_wait(query, timeout_sec=5.0)
  294. if handle:
  295. result = db.fetch(handle, rows=5000)
  296. db.close(handle)
  297. table_ddls.append((0, 0, ' '.join([row[0] for row in result.rows()]), path['database']))
  298. full_table_stats = json.loads(get_table_stats(request, database=path['database'], table=path['table']).content)
  299. stats = dict((stat['data_type'], stat['comment']) for stat in full_table_stats['stats'])
  300. table_stats.append({
  301. 'table_name': '%(database)s.%(table)s' % path, # DB Prefix
  302. 'num_rows': stats.get('numRows', -1),
  303. 'last_modified_time': stats.get('transient_lastDdlTime', -1),
  304. 'total_size': stats.get('totalSize', -1),
  305. 'raw_data_size': stats.get('rawDataSize', -1),
  306. 'num_files': stats.get('numFiles', -1),
  307. 'num_partitions': stats.get('numPartitions', -1),
  308. # bytes_cached
  309. # cache_replication
  310. # format
  311. })
  312. if with_columns:
  313. for col in full_table_stats['columns']:
  314. col_stats = json.loads(get_table_stats(request, database=path['database'], table=path['table'], column=col).content)['stats']
  315. col_stats = dict([(key, val) for col_stat in col_stats for key, val in col_stat.iteritems()])
  316. column_stats.append({
  317. 'table_name': '%(database)s.%(table)s' % path, # DB Prefix
  318. 'column_name': col,
  319. 'data_type': col_stats['data_type'],
  320. "num_distinct": int(col_stats.get('distinct_count')) if col_stats.get('distinct_count') != '' else -1,
  321. "num_nulls": int(col_stats['num_nulls']) if col_stats['num_nulls'] != '' else -1,
  322. "avg_col_len": int(float(col_stats['avg_col_len'])) if col_stats['avg_col_len'] != '' else -1,
  323. "max_size": int(float(col_stats['max_col_len'])) if col_stats['max_col_len'] != '' else -1,
  324. "min": col_stats['min'] if col_stats.get('min', '') != '' else -1,
  325. "max": col_stats['max'] if col_stats.get('max', '') != '' else -1,
  326. "num_trues": col_stats['num_trues'] if col_stats.get('num_trues', '') != '' else -1,
  327. "num_falses": col_stats['num_falses'] if col_stats.get('num_falses', '') != '' else -1,
  328. })
  329. except Exception, e:
  330. LOG.exception('Skipping upload of %s: %s' % (db_table, e))
  331. api = OptimizerApi(request.user)
  332. response['upload_table_stats'] = api.upload(data=table_stats, data_type='table_stats', source_platform=source_platform)
  333. response['status'] = 0 if response['upload_table_stats']['status']['state'] in ('WAITING', 'FINISHED', 'IN_PROGRESS') else -1
  334. if column_stats:
  335. response['upload_cols_stats'] = api.upload(data=column_stats, data_type='cols_stats', source_platform=source_platform)
  336. response['status'] = response['status'] if response['upload_cols_stats']['status']['state'] in ('WAITING', 'FINISHED', 'IN_PROGRESS') else -1
  337. if table_ddls:
  338. response['upload_table_ddl'] = api.upload(data=table_ddls, data_type='queries', source_platform=source_platform)
  339. return JsonResponse(response)
  340. @require_POST
  341. @error_handler
  342. def upload_status(request):
  343. response = {'status': -1}
  344. workload_id = request.POST.get('workloadId')
  345. api = OptimizerApi(request.user)
  346. response['upload_status'] = api.upload_status(workload_id=workload_id)
  347. response['status'] = 0
  348. return JsonResponse(response)