Browse Source

HUE-8578 [importer] Get basic Flume ingest step integrated

Romain Rigaux 7 years ago
parent
commit
77447ff

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

@@ -334,7 +334,7 @@ def importer_submit(request):
 
 
     if destination['indexerRunJob'] or source['inputFormat'] == 'stream':
     if destination['indexerRunJob'] or source['inputFormat'] == 'stream':
       _convert_format(source["format"], inverse=True)
       _convert_format(source["format"], inverse=True)
-      job_handle = _large_indexing(request, source, index_name, start_time=start_time, lib_path=destination['indexerJobLibPath'])
+      job_handle = _large_indexing(request, source, index_name, start_time=start_time, lib_path=destination['indexerJobLibPath'], destination=destination)
     else:
     else:
       client = SolrClient(request.user)
       client = SolrClient(request.user)
       job_handle = _small_indexing(request.user, request.fs, client, source, destination, index_name)
       job_handle = _small_indexing(request.user, request.fs, client, source, destination, index_name)
@@ -454,7 +454,7 @@ def _create_table(request, source, destination, start_time=-1):
   return notebook.execute(request, batch=False)
   return notebook.execute(request, batch=False)
 
 
 
 
-def _large_indexing(request, file_format, collection_name, query=None, start_time=None, lib_path=None):
+def _large_indexing(request, file_format, collection_name, query=None, start_time=None, lib_path=None, destination=None):
   indexer = MorphlineIndexer(request.user, request.fs)
   indexer = MorphlineIndexer(request.user, request.fs)
 
 
   unique_field = indexer.get_unique_field(file_format)
   unique_field = indexer.get_unique_field(file_format)
@@ -482,7 +482,7 @@ def _large_indexing(request, file_format, collection_name, query=None, start_tim
     table_metadata = db.get_table(database=file_format['databaseName'], table_name=file_format['tableName'])
     table_metadata = db.get_table(database=file_format['databaseName'], table_name=file_format['tableName'])
     input_path = table_metadata.path_location
     input_path = table_metadata.path_location
   elif file_format['inputFormat'] == 'stream':
   elif file_format['inputFormat'] == 'stream':
-    return _envelope_job(request, file_format, collection_name, start_time=start_time, lib_path=lib_path)
+    return _envelope_job(request, file_format, destination, start_time=start_time, lib_path=lib_path)
   elif file_format['inputFormat'] == 'file':
   elif file_format['inputFormat'] == 'file':
     input_path = '${nameNode}%s' % urllib.unquote(file_format["path"])
     input_path = '${nameNode}%s' % urllib.unquote(file_format["path"])
   else:
   else:

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

@@ -1836,6 +1836,19 @@ ${ assist.assistPanel() }
         {'value': 'flume', 'name': 'Flume Agent'}
         {'value': 'flume', 'name': 'Flume Agent'}
       ]);
       ]);
       self.streamSelection = ko.observable(self.publicStreams()[0]['value']);
       self.streamSelection = ko.observable(self.publicStreams()[0]['value']);
+      self.streamSelection.subscribe(function(newValue) {
+        if (newValue == 'flume') {
+          $.post("${ url('metadata:manager_hosts') }", {
+            "service": "flume"
+          }, function (resp) {
+            if (resp.status === 0 && resp.hosts) {
+              self.channelSourceHosts(resp.hosts);
+            } else {
+              $(document).trigger("error", "${ _('Error getting hosts') }" + resp.message);
+            }
+          });
+        }
+      });
 
 
       self.kafkaTopics = ko.observableArray();
       self.kafkaTopics = ko.observableArray();
       self.kafkaSelectedTopics = ko.observable(''); // Currently designed just for one
       self.kafkaSelectedTopics = ko.observable(''); // Currently designed just for one
@@ -1868,7 +1881,7 @@ ${ assist.assistPanel() }
         {'name': '${ _("HTTP") }', 'value': 'http'}
         {'name': '${ _("HTTP") }', 'value': 'http'}
       ]);
       ]);
       self.channelSourceType = ko.observable();
       self.channelSourceType = ko.observable();
-      self.channelSourceHosts = ko.observableArray(['host1.com', 'host2.com', 'host3.com', 'host4.com']);
+      self.channelSourceHosts = ko.observableArray();
       self.channelSourceSelectedHosts = ko.observableArray([]);
       self.channelSourceSelectedHosts = ko.observableArray([]);
       self.channelSourceSelectedHosts.subscribe(function(newVal) {
       self.channelSourceSelectedHosts.subscribe(function(newVal) {
         if (newVal) {
         if (newVal) {

+ 13 - 1
desktop/libs/metadata/src/metadata/manager_api.py

@@ -76,7 +76,6 @@ def admin_only_handler(view_fn):
 
 
 
 
 @error_handler
 @error_handler
-@admin_only_handler
 def hello(request):
 def hello(request):
   api = ManagerApi(request.user)
   api = ManagerApi(request.user)
 
 
@@ -85,6 +84,19 @@ def hello(request):
   return JsonResponse(response)
   return JsonResponse(response)
 
 
 
 
+@error_handler
+def get_hosts(request):
+  response = {
+    'status': 0
+  }
+  api = ManagerApi(request.user)
+
+  if request.POST.get('service', '').lower() == 'flume':
+    response['hosts'] = api.get_flume_agents()
+
+  return JsonResponse(response)
+
+
 @error_handler
 @error_handler
 @admin_only_handler
 @admin_only_handler
 def update_flume_config(request):
 def update_flume_config(request):

+ 21 - 7
desktop/libs/metadata/src/metadata/manager_client.py

@@ -94,15 +94,10 @@ class ManagerApi(object):
 
 
   def get_kafka_brokers(self, cluster_name=None):
   def get_kafka_brokers(self, cluster_name=None):
     try:
     try:
-      cluster = self._get_cluster(cluster_name)
-      services = self._root.get('clusters/%(name)s/services' % cluster)['items']
 
 
-      service = [service for service in services if service['type'] == 'KAFKA'][0]
-      broker_hosts = self._get_roles(cluster['name'], service['name'], 'KAFKA_BROKER')
-      broker_hosts_ids = [broker_host['hostRef']['hostId'] for broker_host in broker_hosts]
+      hosts = self._get_hosts('KAFKA', 'KAFKA_BROKER', cluster_name=cluster_name)
 
 
-      hosts = self._root.get('hosts')['items']
-      brokers_hosts = [host['hostname'] + ':9092' for host in hosts if host['hostId'] in broker_hosts_ids]
+      brokers_hosts = [host['hostname'] + ':9092' for host in hosts]
 
 
       return ','.join(brokers_hosts)
       return ','.join(brokers_hosts)
     except RestException, e:
     except RestException, e:
@@ -156,6 +151,25 @@ class ManagerApi(object):
     )
     )
 
 
 
 
+  def get_flume_agents(self, cluster_name=None):
+    return [host['hostname'] for host in self._get_hosts('FLUME', 'AGENT', cluster_name=cluster_name)]
+
+
+  def _get_hosts(self, service_name, role_name, cluster_name=None):
+    try:
+      cluster = self._get_cluster(cluster_name)
+      services = self._root.get('clusters/%(name)s/services' % cluster)['items']
+
+      service = [service for service in services if service['type'] == service_name][0]
+      hosts = self._get_roles(cluster['name'], service['name'], role_name)
+      hosts_ids = [host['hostRef']['hostId'] for host in hosts]
+
+      hosts = self._root.get('hosts')['items']
+      return [host for host in hosts if host['hostId'] in hosts_ids]
+    except RestException, e:
+      raise ManagerApiException(e)
+
+
   def refresh_flume(self, cluster_name, restart=False):
   def refresh_flume(self, cluster_name, restart=False):
     service = 'FLUME-1'
     service = 'FLUME-1'
     cluster = self._get_cluster(cluster_name)
     cluster = self._get_cluster(cluster_name)

+ 3 - 2
desktop/libs/metadata/src/metadata/urls.py

@@ -79,8 +79,9 @@ urlpatterns += [
 
 
 # Manager API
 # Manager API
 urlpatterns += [
 urlpatterns += [
-  url(r'^api/manager/hello/?$', metadata_manager_api.hello, name='hello'),
-  url(r'^api/manager/update_flume_config/?$', metadata_manager_api.update_flume_config, name='update_flume_config'),
+  url(r'^api/manager/hello/?$', metadata_manager_api.hello, name='manager_hello'),
+  url(r'^api/manager/get_hosts/?$', metadata_manager_api.get_hosts, name='manager_hosts'),
+  url(r'^api/manager/update_flume_config/?$', metadata_manager_api.update_flume_config, name='manager_update_flume_config'),
 ]
 ]