فهرست منبع

[ksql] List topics API

Romain 5 سال پیش
والد
کامیت
5bff3cbe79

+ 1 - 1
desktop/libs/indexer/src/indexer/api3.py

@@ -168,7 +168,7 @@ def guess_format(request):
         "hasHeader": True,
         "quoteChar": "\"",
         "recordSeparator": "\\n",
-        'topics': get_topics()
+        'topics': get_topics(request.user)
       }
     elif file_format['streamSelection'] == 'flume':
       format_ = {"type": "csv", "fieldSeparator": ",", "hasHeader": True, "quoteChar": "\"", "recordSeparator": "\\n"}

+ 13 - 11
desktop/libs/kafka/src/kafka/kafka_api.py

@@ -60,7 +60,7 @@ def list_topics(request):
   return JsonResponse({
     'status': 0,
     'topics': [
-      {'name': topic} for topic in get_topics()
+      {'name': topic} for topic in get_topics(request.user)
     ]
   })
 
@@ -97,19 +97,21 @@ def create_topic(request):
   })
 
 
-def get_topics():
+def get_topics(user):
   if has_kafka_api():
     return KafkaApi().topics()
   else:
-    try:
-      manager = ManagerApi()
-      broker_host = manager.get_kafka_brokers().split(',')[0].split(':')[0]
-      return [
-        name
-        for name in list(manager.get_kafka_topics(broker_host).keys()) if not name.startswith('__')
-      ]
-    except Exception as e:
-      return ['user_behavior']
+    from metadata.models.bigquery_client import _get_notebook_api
+    data = {
+      'snippet': {},
+      'database': 'topics'
+    }
+
+    return [
+      topic['name']
+      for topic in _get_notebook_api(user, connector_id=56).autocomplete(**data)['tables_meta']
+      if not topic['name'].startswith('__')
+    ]
 
 
 def get_topic(name):

+ 7 - 3
desktop/libs/notebook/src/notebook/connectors/kafka.py

@@ -19,7 +19,6 @@ 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
@@ -51,7 +50,7 @@ class KafkaApi(Api):
       if database is None:
         response['databases'] = ['default']
       elif table is None:
-        response['tables_meta'] = get_topics()
+        response['tables_meta'] = get_topics(self.user)
       else:
         response = {
           u'status': 0,
@@ -92,7 +91,12 @@ class KafkaApi(Api):
           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'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
         }