Forráskód Böngészése

HUE-8824 [metadata] Start refactoring of optimizer API

Very similar to the refactoring of the catalog API which can now
support Navigator, Apache Atlas or other catalogs.
Romain 6 éve
szülő
commit
73882df23c

+ 1 - 2
desktop/libs/metadata/src/metadata/catalog/dummy_client.py

@@ -26,7 +26,7 @@ from metadata.catalog.base import Api
 LOG = logging.getLogger(__name__)
 
 
-class DummyApi(Api):
+class DummyClient(Api):
 
   def __init__(self, user=None):
     self.user = user
@@ -61,4 +61,3 @@ class DummyApi(Api):
     # For updating comments of table or columns
     # Returning the entity but not used currently
     return {u'clusteredByColNames': None, u'customProperties': {}, u'owner': u'admin', u'serdeName': None, u'deleteTime': None, u'fileSystemPath': u'hdfs://self-service-analytics-1.gce.cloudera.com:8020/user/hive/warehouse/sample_07', u'sourceType': u'HIVE', u'serdeLibName': u'org.apache.hadoop.hive.serde2.lazy.LazySimpleSerDe', u'lastModifiedBy': None, u'sortByColNames': None, u'partColNames': None, u'type': u'TABLE', u'internalType': u'hv_table', u'description': u'Adding an description', u'inputFormat': u'org.apache.hadoop.mapred.TextInputFormat', u'tags': [u'usage'], u'deleted': False, u'technicalProperties': None, u'userEntity': False, u'serdeProps': None, u'originalDescription': None, u'compressed': False, u'metaClassName': u'hv_table', u'properties': {u'__cloudera_internal__hueLink': u'http://self-service-analytics-1.gce.cloudera.com:8889/metastore/table/default/sample_07'}, u'identity': u'22', u'outputFormat': u'org.apache.hadoop.hive.ql.io.HiveIgnoreKeyTextOutputFormat', u'firstClassParentId': None, u'name': None, u'extractorRunId': u'8##503', u'created': u'2018-03-30T17:14:42.000Z', u'sourceId': u'8', u'lastModified': None, u'packageName': u'nav', u'parentPath': u'/default', u'originalName': u'sample_07', u'lastAccessed': u'1970-01-01T00:00:00.000Z'}
-

+ 4 - 0
desktop/libs/metadata/src/metadata/conf.py

@@ -94,6 +94,10 @@ OPTIMIZER = ConfigSection(
   key='optimizer',
   help=_t("""Configuration options for Optimizer API"""),
   members=dict(
+    INTERFACE=Config(
+      key='interface',
+      help=_t('Type of Optimizer to connect to, e.g. Navigator Optimizer...'),
+      default='optimizer'),
     HOSTNAME=Config(
       key='hostname',
       help=_t('Hostname to Optimizer API or compatible service.'),

+ 198 - 0
desktop/libs/metadata/src/metadata/optimizer/dummy_client.py

@@ -0,0 +1,198 @@
+#!/usr/bin/env python
+# -- coding: utf-8 --
+# Licensed to Cloudera, Inc. under one
+# or more contributor license agreements.  See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership.  Cloudera, Inc. licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License.  You may obtain a copy of the License at
+#
+#     http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing, software
+# distributed under the License is distributed on an "AS IS" BASIS,
+# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+# See the License for the specific language governing permissions and
+# limitations under the License.
+
+from builtins import object
+import json
+import logging
+import os
+import time
+import uuid
+
+from tempfile import NamedTemporaryFile
+
+from django.core.cache import cache
+from django.utils.functional import wraps
+from django.utils.translation import ugettext as _
+
+from desktop.auth.backend import is_admin
+from desktop.lib.exceptions_renderable import PopupException
+from desktop.lib import export_csvxls
+from desktop.lib.i18n import smart_unicode
+from desktop.lib.rest.http_client import RestException
+from libsentry.sentry_site import get_hive_sentry_provider
+from libsentry.privilege_checker import get_checker, MissingSentryPrivilegeException
+
+from metadata.conf import OPTIMIZER, get_optimizer_url
+from metadata.optimizer.base import Api
+
+
+LOG = logging.getLogger(__name__)
+
+_JSON_CONTENT_TYPE = 'application/json'
+
+
+class DummyClient(Api):
+
+  def __init__(self, user, api_url=None, auth_key=None, auth_key_secret=None, tenant_id=None):
+    self.user = user
+    self._api_url = (api_url or get_optimizer_url()).strip('/')
+    self._auth_key = auth_key if auth_key else OPTIMIZER.AUTH_KEY_ID.get()
+    self._auth_key_secret = auth_key_secret if auth_key_secret else (OPTIMIZER.AUTH_KEY_SECRET.get() and OPTIMIZER.AUTH_KEY_SECRET.get().replace('\\n', '\n'))
+
+    self._api = None
+
+    self._tenant_id = tenant_id
+
+
+  def get_tenant(self, cluster_id='default'):
+    pass
+
+
+  def upload(self, data, data_type='queries', source_platform='generic', workload_id=None):
+    pass
+
+
+  def upload_status(self, workload_id):
+    pass
+
+  # Sentry permissions work bottom to top.
+  # @check_privileges
+  def top_tables(self, workfloadId=None, database_name='default', page_size=1000, startingToken=None):
+    data = {'results': []}
+
+    return data
+
+  @check_privileges
+  def table_details(self, database_name, table_name, page_size=100, startingToken=None):
+    return {}
+
+
+  def query_compatibility(self, source_platform, target_platform, query, page_size=100, startingToken=None):
+    return {}
+
+
+  def query_risk(self, query, source_platform, db_name, page_size=100, startingToken=None):
+    return {
+      'hints': hints,
+      'noStats': response.get('noStats', []),
+      'noDDL': response.get('noDDL', []),
+    }
+
+  def similar_queries(self, source_platform, query, page_size=100, startingToken=None):
+    raise PopupException(_('Call not supported'))
+
+
+  @check_privileges
+  def top_filters(self, db_tables=None, page_size=100, startingToken=None):
+    results = {'result': []}
+
+    return results
+
+
+  @check_privileges
+  def top_aggs(self, db_tables=None, page_size=100, startingToken=None):
+    results = {'result': []}
+
+    return results
+
+
+  @check_privileges
+  def top_columns(self, db_tables=None, page_size=100, startingToken=None):
+    results = {'results': []}
+
+    return results
+
+
+  @check_privileges
+  def top_joins(self, db_tables=None, page_size=100, startingToken=None):
+    results = {'results': []}
+
+    return results
+
+
+  def top_databases(self, page_size=100, startingToken=None):
+    args = {
+      'tenant' : self._tenant_id,
+      'pageSize': page_size,
+      'startingToken': startingToken
+    }
+
+    data = self._call('getTopDatabases', args)
+
+    if OPTIMIZER.APPLY_SENTRY_PERMISSIONS.get():
+      data['results'] = list(_secure_results(data['results'], self.user))
+
+    return data
+
+
+def OptimizerQueryDataAdapter(data):
+  headers = ['SQL_ID', 'ELAPSED_TIME', 'SQL_FULLTEXT', 'DATABASE']
+
+  if data and len(data[0]) == 4:
+    rows = data
+  else:
+    rows = ([str(uuid.uuid4()), 0.0, q, 'default'] for q in data)
+
+  yield headers, rows
+
+
+def _get_table_name(path):
+  column = None
+
+  if path.count('.') == 1:
+    database, table = path.split('.', 1)
+  elif path.count('.') == 2:
+    database, table, column = path.split('.', 2)
+  else:
+    database, table = 'default', path
+
+  name = {'database': database, 'table': table}
+  if column:
+    name['column'] = column
+  return name
+
+
+def _secure_results(results, user, action='SELECT'):
+    if OPTIMIZER.APPLY_SENTRY_PERMISSIONS.get():
+      checker = get_checker(user=user)
+
+      def getkey(result):
+        key = {'server': get_hive_sentry_provider()}
+
+        if 'dbName' in result:
+          key['db'] = result['dbName']
+        elif 'database' in result:
+          key['db'] = result['database']
+        if 'tableName' in result:
+          key['table'] = result['tableName']
+        elif 'table' in result:
+          key['table'] = result['table']
+        if 'columnName' in result:
+          key['column'] = result['columnName']
+        elif 'column' in result:
+          key['column'] = result['column']
+
+        return key
+
+      return checker.filter_objects(results, action, key=getkey)
+    else:
+      return results
+
+
+def _clean_query(query):
+  return ' '.join([line for line in query.strip().splitlines() if not line.strip().startswith('--')])

+ 3 - 23
desktop/libs/metadata/src/metadata/optimizer_client.py → desktop/libs/metadata/src/metadata/optimizer/optimizer_client.py

@@ -88,7 +88,7 @@ def check_privileges(view_func):
   return wraps(view_func)(decorate)
 
 
-class OptimizerApi(object):
+class OptimizerClient(object):
 
   def __init__(self, user, api_url=None, auth_key=None, auth_key_secret=None, tenant_id=None):
     self.user = user
@@ -190,19 +190,8 @@ class OptimizerApi(object):
   # Sentry permissions work bottom to top.
   # @check_privileges
   def top_tables(self, workfloadId=None, database_name='default', page_size=1000, startingToken=None):
-    data = self._call('getTopTables', {'tenant' : self._tenant_id, 'dbName': database_name.lower(), 'pageSize': page_size, 'startingToken': startingToken})
+    return self._call('getTopTables', {'tenant' : self._tenant_id, 'dbName': database_name.lower(), 'pageSize': page_size, 'startingToken': startingToken})
 
-    if OPTIMIZER.APPLY_SENTRY_PERMISSIONS.get():
-      checker = get_checker(user=self.user)
-      action = 'SELECT'
-
-      def getkey(table):
-        names = _get_table_name(table['name'])
-        return {'server': get_hive_sentry_provider(), 'db': names['database'], 'table': names['table']}
-
-      data['results'] = list(checker.filter_objects(data['results'], action, key=getkey))
-
-    return data
 
   @check_privileges
   def table_details(self, database_name, table_name, page_size=100, startingToken=None):
@@ -251,16 +240,7 @@ class OptimizerApi(object):
     if db_tables:
       args['dbTableList'] = [db_table.lower() for db_table in db_tables]
 
-    results = self._call('getTopFilters', args)
-
-    if OPTIMIZER.APPLY_SENTRY_PERMISSIONS.get():
-      filtered_filters = []
-      for result in results['results']:
-        cols = [_get_table_name(col['columnName']) for col in result["popularValues"][0]["group"]]
-        if len(cols) == len(list(_secure_results(cols, self.user))):
-          filtered_filters.append(result)
-      results['results'] = filtered_filters
-    return results
+    return self._call('getTopFilters', args)
 
 
   @check_privileges

+ 0 - 0
desktop/libs/metadata/src/metadata/optimizer_client_tests.py → desktop/libs/metadata/src/metadata/optimizer/optimizer_client_tests.py


+ 61 - 18
desktop/libs/metadata/src/metadata/optimizer_api.py

@@ -25,6 +25,7 @@ from django.http import Http404
 from django.utils.translation import ugettext as _
 from django.views.decorators.http import require_POST
 
+from desktop.auth.backend import is_admin
 from desktop.lib.django_util import JsonResponse
 from desktop.lib.i18n import force_unicode
 from desktop.models import Document2
@@ -32,11 +33,10 @@ from libsentry.privilege_checker import MissingSentryPrivilegeException
 from notebook.api import _get_statement
 from notebook.models import Notebook
 
-from metadata.optimizer_client import OptimizerApi, NavOptException, _get_table_name, _clean_query
+from metadata.optimizer.base import get_api
+from metadata.optimizer.optimizer_client import NavOptException, _get_table_name, _clean_query
 from metadata.conf import OPTIMIZER
 
-from desktop.auth.backend import is_admin
-
 LOG = logging.getLogger(__name__)
 
 
@@ -82,9 +82,11 @@ def error_handler(view_fn):
 def get_tenant(request):
   response = {'status': -1}
 
+  interface = request.POST.get('interface', OPTIMIZER.INTERFACE.get())
   cluster_id = request.POST.get('cluster_id')
 
-  api = OptimizerApi(request.user)
+  api = get_api(request, interface)
+
   data = api.get_tenant(cluster_id=cluster_id)
 
   if data:
@@ -101,12 +103,26 @@ def get_tenant(request):
 def top_tables(request):
   response = {'status': -1}
 
+  interface = request.POST.get('interface', OPTIMIZER.INTERFACE.get())
   database = request.POST.get('database', 'default')
   limit = request.POST.get('len', 1000)
 
-  api = OptimizerApi(user=request.user)
+  api = get_api(request, interface)
+
   data = api.top_tables(database_name=database, page_size=limit)
 
+
+  if OPTIMIZER.APPLY_SENTRY_PERMISSIONS.get():
+    checker = get_checker(user=self.user)
+    action = 'SELECT'
+
+    def getkey(table):
+      names = _get_table_name(table['name'])
+      return {'server': get_hive_sentry_provider(), 'db': names['database'], 'table': names['table']}
+
+    data['results'] = list(checker.filter_objects(data['results'], action, key=getkey))
+
+
   tables = [{
       'eid': table['eid'],
       'database': _get_table_name(table['name'])['database'],
@@ -130,10 +146,11 @@ def top_tables(request):
 def table_details(request):
   response = {'status': -1}
 
+  interface = request.POST.get('interface', OPTIMIZER.INTERFACE.get())
   database_name = request.POST.get('databaseName')
   table_name = request.POST.get('tableName')
 
-  api = OptimizerApi(request.user)
+  api = get_api(request, interface)
 
   data = api.table_details(database_name=database_name, table_name=table_name)
 
@@ -151,11 +168,12 @@ def table_details(request):
 def query_compatibility(request):
   response = {'status': -1}
 
+  interface = request.POST.get('interface', OPTIMIZER.INTERFACE.get())
   source_platform = request.POST.get('sourcePlatform')
   target_platform = request.POST.get('targetPlatform')
   query = request.POST.get('query')
 
-  api = OptimizerApi(request.user)
+  api = get_api(request, interface)
 
   data = api.query_compatibility(source_platform=source_platform, target_platform=target_platform, query=query)
 
@@ -173,11 +191,12 @@ def query_compatibility(request):
 def query_risk(request):
   response = {'status': -1}
 
+  interface = request.POST.get('interface', OPTIMIZER.INTERFACE.get())
   query = json.loads(request.POST.get('query'))
   source_platform = request.POST.get('sourcePlatform')
   db_name = request.POST.get('dbName')
 
-  api = OptimizerApi(request.user)
+  api = get_api(request, interface)
 
   data = api.query_risk(query=query, source_platform=source_platform, db_name=db_name)
 
@@ -195,10 +214,11 @@ def query_risk(request):
 def similar_queries(request):
   response = {'status': -1}
 
+  interface = request.POST.get('interface', OPTIMIZER.INTERFACE.get())
   source_platform = request.POST.get('sourcePlatform')
   query = json.loads(request.POST.get('query'))
 
-  api = OptimizerApi(request.user)
+  api = get_api(request, interface)
 
   data = api.similar_queries(source_platform=source_platform, query=query)
 
@@ -216,10 +236,12 @@ def similar_queries(request):
 def top_filters(request):
   response = {'status': -1}
 
+  interface = request.POST.get('interface', OPTIMIZER.INTERFACE.get())
   db_tables = json.loads(request.POST.get('dbTables'), '[]')
   column_name = request.POST.get('columnName') # Unused
 
-  api = OptimizerApi(request.user)
+  api = get_api(request, interface)
+
   data = api.top_filters(db_tables=db_tables)
 
   if data:
@@ -228,6 +250,14 @@ def top_filters(request):
   else:
     response['message'] = 'Optimizer: %s' % data
 
+  if OPTIMIZER.APPLY_SENTRY_PERMISSIONS.get():
+    filtered_filters = []
+    for result in results['results']:
+      cols = [_get_table_name(col['columnName']) for col in result["popularValues"][0]["group"]]
+      if len(cols) == len(list(_secure_results(cols, self.user))):
+        filtered_filters.append(result)
+    results['results'] = filtered_filters
+
   return JsonResponse(response)
 
 
@@ -236,9 +266,11 @@ def top_filters(request):
 def top_joins(request):
   response = {'status': -1}
 
+  interface = request.POST.get('interface', OPTIMIZER.INTERFACE.get())
   db_tables = json.loads(request.POST.get('dbTables'), '[]')
 
-  api = OptimizerApi(request.user)
+  api = get_api(request, interface)
+
   data = api.top_joins(db_tables=db_tables)
 
   if data:
@@ -255,9 +287,11 @@ def top_joins(request):
 def top_aggs(request):
   response = {'status': -1}
 
+  interface = request.POST.get('interface', OPTIMIZER.INTERFACE.get())
   db_tables = json.loads(request.POST.get('dbTables'), '[]')
 
-  api = OptimizerApi(request.user)
+  api = get_api(request, interface)
+
   data = api.top_aggs(db_tables=db_tables)
 
   if data:
@@ -274,7 +308,10 @@ def top_aggs(request):
 def top_databases(request):
   response = {'status': -1}
 
-  api = OptimizerApi(request.user)
+  interface = request.POST.get('interface', OPTIMIZER.INTERFACE.get())
+
+  api = get_api(request, interface)
+
   data = api.top_databases()
 
   if data:
@@ -291,9 +328,11 @@ def top_databases(request):
 def top_columns(request):
   response = {'status': -1}
 
+  interface = request.POST.get('interface', OPTIMIZER.INTERFACE.get())
   db_tables = json.loads(request.POST.get('dbTables'), '[]')
 
-  api = OptimizerApi(request.user)
+  api = get_api(request, interface)
+
   data = api.top_columns(db_tables=db_tables)
 
   if data:
@@ -328,7 +367,8 @@ def upload_history(request):
   response = {'status': -1}
 
   if is_admin(request.user):
-    api = OptimizerApi(request.user)
+    interface = request.POST.get('interface', OPTIMIZER.INTERFACE.get())
+    api = get_api(request, interface)
     histories = []
     upload_stats = {}
 
@@ -360,6 +400,7 @@ def upload_history(request):
 def upload_query(request):
   response = {'status': -1}
 
+  interface = request.POST.get('interface', OPTIMIZER.INTERFACE.get())
   source_platform = request.POST.get('sourcePlatform', 'default')
   query_id = request.POST.get('query_id')
 
@@ -371,7 +412,7 @@ def upload_query(request):
       queries = _convert_queries([query_data])
       source_platform = query_data['snippets'][0]['type']
 
-      api = OptimizerApi(request.user)
+      api = get_api(request, interface)
 
       response['query_upload'] = api.upload(data=queries, data_type='queries', source_platform=source_platform)
     except Document2.DoesNotExist:
@@ -388,6 +429,7 @@ def upload_query(request):
 def upload_table_stats(request):
   response = {'status': -1}
 
+  interface = request.POST.get('interface', OPTIMIZER.INTERFACE.get())
   db_tables = json.loads(request.POST.get('db_tables'), '[]')
   source_platform = json.loads(request.POST.get('sourcePlatform', '"hive"'))
   with_ddl = json.loads(request.POST.get('with_ddl', 'false'))
@@ -465,7 +507,7 @@ def upload_table_stats(request):
     except Exception as e:
       LOG.exception('Skipping upload of %s: %s' % (db_table, e))
 
-  api = OptimizerApi(request.user)
+  api = get_api(request, interface)
 
   response['status'] = 0
 
@@ -492,9 +534,10 @@ def upload_table_stats(request):
 def upload_status(request):
   response = {'status': -1}
 
+  interface = request.POST.get('interface', OPTIMIZER.INTERFACE.get())
   workload_id = request.POST.get('workloadId')
 
-  api = OptimizerApi(request.user)
+  api = get_api(request, interface)
 
   response['upload_status'] = api.upload_status(workload_id=workload_id)
   response['status'] = 0

+ 7 - 7
desktop/libs/notebook/src/notebook/connectors/hiveserver2.py

@@ -42,7 +42,7 @@ from desktop.lib.paths import SAFE_CHARACTERS_URI_COMPONENTS
 from desktop.lib.rest.http_client import RestException
 from desktop.lib.thrift_util import unpack_guid, unpack_guid_base64
 from desktop.models import DefaultConfiguration, Document2
-from metadata.optimizer_client import OptimizerApi
+from metadata.optimizer.optimizer_client import OptimizerClient
 
 from notebook.connectors.base import Api, QueryError, QueryExpired, OperationTimeout, OperationNotSupported, _get_snippet_name, Notebook, get_interpreter
 
@@ -595,27 +595,27 @@ DROP TABLE IF EXISTS `%(table)s`;
     response = self._get_current_statement(notebook, snippet)
     query = response['statement']
 
-    api = OptimizerApi(self.user)
+    client = OptimizerClient(self.user)
 
-    return api.query_risk(query=query, source_platform=snippet['type'], db_name=snippet.get('database') or 'default')
+    return client.query_risk(query=query, source_platform=snippet['type'], db_name=snippet.get('database') or 'default')
 
 
   def statement_compatibility(self, notebook, snippet, source_platform, target_platform):
     response = self._get_current_statement(notebook, snippet)
     query = response['statement']
 
-    api = OptimizerApi(self.user)
+    client = OptimizerClient(self.user)
 
-    return api.query_compatibility(source_platform, target_platform, query)
+    return client.query_compatibility(source_platform, target_platform, query)
 
 
   def statement_similarity(self, notebook, snippet, source_platform):
     response = self._get_current_statement(notebook, snippet)
     query = response['statement']
 
-    api = OptimizerApi(self.user)
+    client = OptimizerClient(self.user)
 
-    return api.similar_queries(source_platform, query)
+    return client.similar_queries(source_platform, query)
 
 
   def upgrade_properties(self, lang='hive', properties=None):