Browse Source

HUE-8208 [importer] Kafka topic listing

Romain Rigaux 7 years ago
parent
commit
b6ee2ff702

+ 15 - 2
desktop/libs/indexer/src/indexer/api3.py

@@ -29,6 +29,9 @@ from desktop.lib.exceptions_renderable import PopupException
 from desktop.lib.i18n import smart_unicode
 from desktop.models import Document2
 from librdbms.server import dbms as rdbms
+from libsentry.conf import is_enabled
+from metadata.kafka_client import KafkaApi
+from metadata.manager_client import ManagerApi
 from notebook.connectors.base import get_api, Notebook
 from notebook.decorators import api_error_handler
 from notebook.models import make_notebook, MockedDjangoRequest, escape_rows
@@ -117,6 +120,8 @@ def guess_format(request):
     format_ = {"quoteChar": "\"", "recordSeparator": "\\n", "type": "csv", "hasHeader": False, "fieldSeparator": "\u0001"}
   elif file_format['inputFormat'] == 'rdbms':
     format_ = RdbmsIndexer(request.user, file_format['rdbmsType']).guess_format()
+  elif file_format['inputFormat'] == 'kafka':
+    format_ = {'type': 'csv', 'topics': KafkaApi().topics()}
 
   format_['status'] = 0
   return JsonResponse(format_)
@@ -422,8 +427,16 @@ def _envelope_job(request, file_format, collection_name, start_time=None, lib_pa
     input_path = '${nameNode}%s' % file_format["path"]
   else:
     input_path = None
+    
+    manager = ManagerApi()
+
+    properties = {
+      "brokers": manager.get_kafka_brokers(),
+      "kudu_master": manager.get_kudu_master(),
+      "output_table": "impala::%s" % collection_name,
+      "topics": file_format['kafkaSelectedTopics']
+    }
 
-  morphline = indexer.generate_config()
+  morphline = indexer.generate_config(properties)
 
   return indexer.run(request, collection_name, morphline, input_path, start_time=start_time, lib_path=lib_path)
-

+ 1 - 8
desktop/libs/indexer/src/indexer/indexers/envelope.py

@@ -76,14 +76,7 @@ class EnvelopeIndexer(object):
     return task.execute(request, batch=True)
 
 
-  def generate_config(self):
-    properties = {
-      "brokers": "self-service-analytics-1.gce.cloudera.com:9092,self-service-analytics-2.gce.cloudera.com:9092,self-service-analytics-3.gce.cloudera.com:9092",
-      "kudu_master": "self-service-analytics-1.gce.cloudera.com:7051",
-      "output_table": "impala::default.traffic_conditions",
-      "topics": "traffic"
-    }
-
+  def generate_config(self, properties):
     return """
 application {
     name = Traffic analysis

+ 24 - 14
desktop/libs/indexer/src/indexer/templates/importer.mako

@@ -323,17 +323,20 @@ ${ assist.assistPanel() }
           <!-- /ko -->
 
           <!-- ko if: createWizard.source.inputFormat() == 'kafka' -->
-            ## Service
+            ##<div class="control-group">
+            ##  <label for="rdbmsHostname" class="control-label"><div>${ _('Brokers') }</div>
+            ##    <input type="text" class="input-xxlarge" data-bind="value: createWizard.source.kafkaBrokers" placeholder="${ _('Enter a csv list of brokers, e.g.brokers1:9092,brokers2:9092') }">
+            ##  </label>
+            ##</div>
 
             <div class="control-group">
-              <label for="rdbmsHostname" class="control-label"><div>${ _('Brokers') }</div>
-                <input type="text" class="input-xxlarge" data-bind="value: createWizard.source.kafkaBrokers" placeholder="${ _('Enter a csv list of brokers, e.g.brokers1:9092,brokers2:9092') }">
-              </label>
-            </div>
-
-            <div class="control-group">
-              <label for="rdbmsHostname" class="control-label"><div>${ _('Topic') }</div>
-                <input type="text" class="input-xxlarge" data-bind="value: createWizard.source.kafkaTopics" placeholder="${ _('The list of topics to consume, e.g. orders,returns') }">
+              <label class="control-label"><div>${ _('Topics') }</div>
+                ##<input type="text" class="input-xxlarge" data-bind="value: createWizard.source.kafkaTopics">
+                <select data-bind="options: createWizard.source.kafkaTopics,
+                       value: createWizard.source.kafkaSelectedTopics,
+                       optionsCaption: 'Choose...'"
+                       placeholder="${ _('The list of topics to consume, e.g. orders,returns') }"></select>
+                ##<select data-bind="selectize: createWizard.source.kafkaTopics, value: createWizard.source.kafkaSelectedTopics" placeholder="${ _('The list of topics to consume, e.g. orders,returns') }"></select>
               </label>
             </div>
           <!-- /ko -->
@@ -1286,6 +1289,9 @@ ${ assist.assistPanel() }
         self.path('');
         resizeElements();
         self.rdbmsMode('customRdbms');
+        if (val == 'kafka') {
+          wizard.guessFormat();
+        }
       });
       self.inputFormatsAll = ko.observableArray([
           {'value': 'file', 'name': 'File'},
@@ -1294,7 +1300,7 @@ ${ assist.assistPanel() }
           {'value': 'rdbms', 'name': 'External Database'},
           % endif
           % if ENABLE_KAFKA.get():
-          {'value': 'kafka', 'name': 'Kafka Stream'},
+          {'value': 'kafka', 'name': 'Internal Stream'},
           % endif
           % if ENABLE_SQL_INDEXER.get():
           {'value': 'query', 'name': 'SQL Query'},
@@ -1490,9 +1496,10 @@ ${ assist.assistPanel() }
       self.draggedQuery = ko.observable();
 
       // Kafka
-      self.kafkaBrokers = ko.observable('brokers1:9092,brokers2:9092');
-      self.kafkaTopics = ko.observable('');
-      self.kafkaTopics.subscribe(function(newValue) {
+      self.kafkaBrokers = ko.observable('brokers1:9092,brokers2:9092'); // Unused
+      self.kafkaTopics = ko.observable([]);
+      self.kafkaSelectedTopics = ko.observable('');
+      self.kafkaSelectedTopics.subscribe(function(newValue) {
         if (newValue) {
           viewModel.createWizard.guessFieldTypes();
         }
@@ -1528,7 +1535,7 @@ ${ assist.assistPanel() }
         } else if (self.inputFormat() == 'manual') {
           return true;
         } else if (self.inputFormat() == 'kafka') {
-          return self.kafkaBrokers().length > 0 && self.kafkaTopics().length > 0;
+          return self.kafkaBrokers().length > 0 && self.kafkaSelectedTopics().length > 0;
         } else if (self.inputFormat() == 'rdbms') {
           return self.rdbmsDatabaseName().length > 0 && (self.rdbmsTableName().length > 0 || self.rdbmsAllTablesSelected());
         }
@@ -1949,6 +1956,9 @@ ${ assist.assistPanel() }
           } else {
             var newFormat = ko.mapping.fromJS(new FileType(resp['type'], resp));
             self.source.format(newFormat);
+            if (self.source.inputFormat() == 'kafka') {
+              self.source.kafkaTopics(resp['topics']);
+            }
             self.guessFieldTypes();
           }
 

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

@@ -392,6 +392,17 @@ MANAGER = ConfigSection(
 )
 
 
+KAFKA = ConfigSection(
+  key='kafka',
+  help=_t("""Configuration options for Kafka API"""),
+  members=dict(
+    API_URL=Config(
+      key='api_url',
+      help=_t('Base URL to API.'),
+      default=None),
+  )
+)
+
 
 def test_metadata_configurations(user):
   from libsentry.conf import is_enabled

+ 5 - 14
desktop/libs/metadata/src/metadata/kafka_client.py

@@ -17,6 +17,7 @@
 # limitations under the License.
 
 import logging
+import json
 
 from django.core.cache import cache
 from django.utils.translation import ugettext as _
@@ -24,7 +25,7 @@ from django.utils.translation import ugettext as _
 from desktop.lib.rest.http_client import RestException, HttpClient
 from desktop.lib.rest.resource import Resource
 from desktop.lib.i18n import smart_unicode
-from metadata.conf import MANAGER
+from metadata.conf import KAFKA
 
 
 LOG = logging.getLogger(__name__)
@@ -47,26 +48,16 @@ class KafkaApi(object):
   """
 
   def __init__(self, user=None, security_enabled=False, ssl_cert_ca_verify=False):
-    self._api_url = '%s/%s' % (MANAGER.API_URL.get().strip('/'))
-    # localhost:8082/topics
-    
-    self._username = 'hue' #get_navigator_auth_username()
-    self._password = 'hue' #get_navigator_auth_password()
+    self._api_url = KAFKA.API_URL.get().strip('/')
 
     self.user = user
     self._client = HttpClient(self._api_url, logger=LOG)
-
-    if security_enabled:
-      self._client.set_kerberos_auth()
-    else:
-      self._client.set_basic_auth(self._username, self._password)
-
-    self._client.set_verify(ssl_cert_ca_verify)
     self._root = Resource(self._client)
 
 
   def topics(self):
     try:
-      return self._root.get('topics')
+      response = self._root.get('topics')
+      return json.loads(response)
     except RestException, e:
       raise KafkaApiException(e)