Explorar o código

HUE-8775 [notebook] Add HMS connector implementation

Romain Rigaux %!s(int64=6) %!d(string=hai) anos
pai
achega
f91ccffdb5

+ 16 - 1
apps/beeswax/src/beeswax/server/dbms.py

@@ -33,7 +33,7 @@ from desktop.models import Cluster
 from indexer.file_format import HiveFormat
 
 from beeswax import hive_site
-from beeswax.conf import HIVE_SERVER_HOST, HIVE_SERVER_PORT, LIST_PARTITIONS_LIMIT, SERVER_CONN_TIMEOUT, \
+from beeswax.conf import HIVE_SERVER_HOST, HIVE_SERVER_PORT, HIVE_METASTORE_HOST, HIVE_METASTORE_PORT, LIST_PARTITIONS_LIMIT, SERVER_CONN_TIMEOUT, \
   AUTH_USERNAME, AUTH_PASSWORD, APPLY_NATURAL_SORT_MAX, QUERY_PARTITIONS_LIMIT
 from beeswax.common import apply_natural_sort
 from beeswax.design import hql_query
@@ -67,6 +67,9 @@ def get(user, query_server=None, cluster=None):
         from impala.dbms import ImpalaDbms
         from impala.server import ImpalaServerClient
         DBMS_CACHE[user.id][query_server['server_name']] = ImpalaDbms(HiveServerClientCompatible(ImpalaServerClient(query_server, user)), QueryHistory.SERVER_TYPE[1][0])
+      elif query_server['server_name'] == 'hms':
+        from beeswax.server.hive_metastore_server import HiveMetastoreClient
+        DBMS_CACHE[user.id][query_server['server_name']] = HiveServer2Dbms(HiveMetastoreClient(query_server, user), QueryHistory.SERVER_TYPE[1][0])
       else:
         from beeswax.server.hive_server2_lib import HiveServerClient
         DBMS_CACHE[user.id][query_server['server_name']] = HiveServer2Dbms(HiveServerClientCompatible(HiveServerClient(query_server, user)), QueryHistory.SERVER_TYPE[1][0])
@@ -84,6 +87,18 @@ def get_query_server_config(name='beeswax', server=None, cluster=None):
   if name == 'impala':
     from impala.dbms import get_query_server_config as impala_query_server_config
     query_server = impala_query_server_config(cluster_config=cluster_config)
+  elif name == 'hms':
+    kerberos_principal = hive_site.get_hiveserver2_kerberos_principal(HIVE_SERVER_HOST.get())
+
+    query_server = {
+        'server_name': 'hms',
+        'server_host': HIVE_METASTORE_HOST.get() if not cluster_config else cluster_config.get('server_host'),
+        'server_port': HIVE_METASTORE_PORT.get(),
+        'principal': kerberos_principal,
+        'transport_mode': 'http' if hive_site.hiveserver2_transport_mode() == 'HTTP' else 'socket',
+        'auth_username': AUTH_USERNAME.get(),
+        'auth_password': AUTH_PASSWORD.get()
+    }
   else:
     kerberos_principal = hive_site.get_hiveserver2_kerberos_principal(HIVE_SERVER_HOST.get())
 

+ 90 - 0
desktop/libs/notebook/src/notebook/connectors/hive_metastore.py

@@ -0,0 +1,90 @@
+#!/usr/bin/env python
+# 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.
+
+import logging
+
+from django.urls import reverse
+from django.utils.translation import ugettext as _
+
+
+from desktop.lib.exceptions import StructuredException
+from desktop.lib.exceptions_renderable import PopupException
+from desktop.lib.i18n import force_unicode, smart_str
+from desktop.lib.rest.http_client import RestException
+
+from notebook.connectors.base import Api, QueryError, QueryExpired, OperationTimeout, OperationNotSupported
+
+
+LOG = logging.getLogger(__name__)
+
+
+try:
+  from beeswax.api import _autocomplete
+  from beeswax.models import HiveServerQueryHandle
+  from beeswax.server import dbms
+  from beeswax.server.dbms import get_query_server_config, QueryServerException
+except ImportError, e:
+  LOG.warn('Hive and HiveMetastoreServer interfaces are not enabled: %s' % e)
+  hive_settings = None
+
+
+def query_error_handler(func):
+  def decorator(*args, **kwargs):
+    try:
+      return func(*args, **kwargs)
+    except StructuredException, e:
+      message = force_unicode(str(e))
+      if 'timed out' in message:
+        raise OperationTimeout(e)
+      else:
+        raise QueryError(message)
+    except QueryServerException, e:
+      message = force_unicode(str(e))
+      if 'Invalid query handle' in message or 'Invalid OperationHandle' in message:
+        raise QueryExpired(e)
+      else:
+        raise QueryError(message)
+  return decorator
+
+
+class HiveMetastoreApi(Api):
+
+  @query_error_handler
+  def autocomplete(self, snippet, database=None, table=None, column=None, nested=None):
+    db = self._get_db(snippet, cluster=self.cluster)
+    query = None
+
+    if snippet.get('query'):
+      query = snippet.get('query')
+    elif snippet.get('source') == 'query':
+      document = Document2.objects.get(id=database)
+      document.can_read_or_exception(self.user)
+      notebook = Notebook(document=document).get_data()
+      snippet = notebook['snippets'][0]
+      query = self._get_current_statement(db, snippet)['statement']
+      database, table = '', ''
+
+    return _autocomplete(db, database, table, column, nested, query=query, cluster=self.cluster)
+
+
+  @query_error_handler
+  def get_sample_data(self, snippet, database=None, table=None, column=None, async=False, operation=None):
+    return []
+
+
+  def _get_db(self, snippet, async=False, cluster=None):
+    return dbms.get(self.user, query_server=get_query_server_config(name='hms', cluster=cluster))