Browse Source

HUE-9077 [connector] Adding initial ksql connector

Romain 6 năm trước cách đây
mục cha
commit
fde074c8c5

+ 1 - 1
desktop/libs/kafka/src/kafka/ksql_client.py

@@ -56,7 +56,7 @@ class KSqlApi(object):
     except ImportError:
       raise KSqlApiException('Module missing: pip install ksql')
 
-    self._api_url = KAFKA.API_URL.get().strip('/') if KAFKA.API_URL.get() else ''
+    self._api_url = KAFKA.KSQL_API_URL.get().strip('/') if KAFKA.KSQL_API_URL.get() else ''
 
     self.user = user
     self.client = client = KSQLAPI(self._api_url)

+ 3 - 0
desktop/libs/notebook/src/notebook/connectors/base.py

@@ -396,6 +396,9 @@ def get_api(request, snippet):
   elif interface == 'hbase':
     from notebook.connectors.hbase import HBaseApi
     return HBaseApi(request.user)
+  elif interface == 'ksql':
+    from notebook.connectors.ksql import KSqlApi
+    return KSqlApi(request.user)
   elif interface == 'kafka':
     from notebook.connectors.kafka import KafkaApi
     return KafkaApi(request.user)

+ 137 - 0
desktop/libs/notebook/src/notebook/connectors/ksql.py

@@ -0,0 +1,137 @@
+#!/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.
+
+from __future__ import absolute_import
+
+import logging
+
+from django.core.urlresolvers import reverse
+from django.utils.translation import ugettext as _
+
+from desktop.lib.i18n import force_unicode
+from kafka.ksql_client import KSqlApi
+
+from notebook.connectors.base import Api, QueryError
+
+
+LOG = logging.getLogger(__name__)
+
+
+def query_error_handler(func):
+  def decorator(*args, **kwargs):
+    try:
+      return func(*args, **kwargs)
+    except Exception as e:
+      message = force_unicode(str(e))
+      raise QueryError(message)
+  return decorator
+
+
+class KafkaApi(Api):
+
+  def __init__(self, user, interpreter=None):
+    Api.__init__(self, user, interpreter=interpreter)
+
+    self.db = KSqlApi(user=user)
+
+
+  @query_error_handler
+  def execute(self, notebook, snippet):
+    if self.db is None:
+      raise AuthenticationRequired()
+
+    data, description = query_and_fetch(self.db, snippet['statement'], 1000)
+    has_result_set = data is not None
+
+    return {
+      'sync': True,
+      'has_result_set': has_result_set,
+      'result': {
+        'has_more': False,
+        'data': data if has_result_set else [],
+        'meta': [{
+          'name': col[0],
+          'type': col[1],
+          'comment': ''
+        } for col in description] if has_result_set else [],
+        'type': 'table'
+      }
+    }
+
+  @query_error_handler
+  def check_status(self, notebook, snippet):
+    return {'status': 'available'}
+
+  @query_error_handler
+  def autocomplete(self, snippet, database=None, table=None, column=None, nested=None):
+    response = {}
+
+    try:
+      if database is None:
+        response['databases'] = ['default']
+      elif table is None:
+        response['tables_meta'] = self.db.show_tables()
+      else:
+        response = {
+          u'status': 0,
+          u'comment': u'test test test 22',
+          u'hdfs_link': u'/filebrowser/view=/user/hive/warehouse/web_logs',
+          u'extended_columns': [
+            {u'comment': u'', u'type': u'bigint', u'name': u'_version_'},
+            {u'comment': u'The app', u'type': u'string', u'name': u'app'},
+            {u'comment': u'test test   test 22', u'type': u'smallint', u'name': u'bytes'},
+            {u'comment': u'The citi', u'type': u'string', u'name': u'city'},
+            {u'comment': u'', u'type': u'string', u'name': u'client_ip'},
+            {u'comment': u'', u'type': u'tinyint', u'name': u'code'},
+            {u'comment': u'', u'type': u'string', u'name': u'country_code'},
+            {u'comment': u'', u'type': u'string', u'name': u'country_code3'},
+            {u'comment': u'', u'type': u'string', u'name': u'country_name'},
+            {u'comment': u'', u'type': u'string', u'name': u'device_family'},
+            {u'comment': u'', u'type': u'string', u'name': u'extension'},
+            {u'comment': u'', u'type': u'float', u'name': u'latitude'},
+            {u'comment': u'', u'type': u'float', u'name': u'longitude'},
+            {u'comment': u'', u'type': u'string', u'name': u'method'},
+            {u'comment': u'', u'type': u'string', u'name': u'os_family'},
+            {u'comment': u'', u'type': u'string', u'name': u'os_major'},
+            {u'comment': u'', u'type': u'string', u'name': u'protocol'},
+            {u'comment': u'', u'type': u'string', u'name': u'record'},
+            {u'comment': u'', u'type': u'string', u'name': u'referer'},
+            {u'comment': u'', u'type': u'bigint', u'name': u'region_code'},
+            {u'comment': u'', u'type': u'string', u'name': u'request'},
+            {u'comment': u'', u'type': u'string', u'name': u'subapp'},
+            {u'comment': u'', u'type': u'string', u'name': u'time'},
+            {u'comment': u'', u'type': u'string', u'name': u'url'},
+            {u'comment': u'', u'type': u'string', u'name': u'user_agent'},
+            {u'comment': u'', u'type': u'string', u'name': u'user_agent_family'},
+            {u'comment': u'', u'type': u'string', u'name': u'user_agent_major'},
+            {u'comment': u'', u'type': u'string', u'name': u'id'},
+            {u'comment': u'', u'type': u'string', u'name': u'date'}
+          ],
+          u'support_updates': False,
+          u'partition_keys': [
+            {u'type': u'string', u'name': u'date'}
+          ],
+          u'columns': [u'_version_', u'app', u'bytes', u'city', u'client_ip', u'code', u'country_code', u'country_code3', u'country_name', u'device_family', u'extension', u'latitude', u'longitude', u'method', u'os_family', u'os_major', u'protocol', u'record', u'referer', u'region_code', u'request', u'subapp', u'time', u'url', u'user_agent', u'user_agent_family', u'user_agent_major', u'id', u'date'],
+          u'is_view': False
+        }
+
+    except Exception as e:
+      LOG.warn('Autocomplete data fetching error: %s' % e)
+      response['code'] = 500
+      response['error'] = e.message
+
+    return response