浏览代码

HUE-8298 [kafka] Modularize the envelope code generation

Romain Rigaux 7 年之前
父节点
当前提交
87f536b946

+ 5 - 1
desktop/core/src/desktop/templates/hue.mako

@@ -407,6 +407,7 @@ ${ hueIcons.symbols() }
       <div id="embeddable_home" class="embeddable"></div>
       <div id="embeddable_catalog" class="embeddable"></div>
       <div id="embeddable_indexer" class="embeddable"></div>
+      <div id="embeddable_kafka" class="embeddable"></div>
       <div id="embeddable_importer" class="embeddable"></div>
       <div id="embeddable_collections" class="embeddable"></div>
       <div id="embeddable_indexes" class="embeddable"></div>
@@ -644,6 +645,7 @@ ${ smart_unicode(login_modal(request).content) | n,unicode }
         home: { url: '/home*', title: '${_('Home')}' },
         catalog: { url: '/catalog', title: '${_('Catalog')}' },
         indexer: { url: '/indexer/indexer/', title: '${_('Indexer')}' },
+        kafka: { url: '/indexer/topics/', title: '${_('Streams')}' },
         collections: { url: '/dashboard/admin/collections', title: '${_('Search')}' },
         % if hasattr(ENABLE_NEW_INDEXER, 'get') and ENABLE_NEW_INDEXER.get():
         indexes: { url: '/indexer/indexes/*', title: '${_('Indexes')}' },
@@ -683,7 +685,7 @@ ${ smart_unicode(login_modal(request).content) | n,unicode }
       };
 
       var SKIP_CACHE = [
-          'home', 'catalog', 'oozie_workflow', 'oozie_coordinator', 'oozie_bundle', 'dashboard', 'metastore',
+          'home', 'oozie_workflow', 'oozie_coordinator', 'oozie_bundle', 'dashboard', 'metastore',
           'filebrowser', 'useradmin_users', 'useradmin_groups', 'useradmin_newgroup', 'useradmin_editgroup',
           'useradmin_permissions', 'useradmin_editpermission', 'useradmin_configurations', 'useradmin_newuser',
           'useradmin_addldapusers', 'useradmin_addldapgroups', 'useradmin_edituser', 'importer',
@@ -1186,6 +1188,8 @@ ${ smart_unicode(login_modal(request).content) | n,unicode }
           }},
           { url: '/home*', app: 'home' },
           { url: '/catalog', app: 'catalog' },
+          { url: '/kafka/', app: 'kafka' },
+          { url: '/indexer/topics/*', app: 'kafka' },
           { url: '/indexer/indexes/*', app: 'indexes' },
           { url: '/indexer/', app: 'indexes' },
           { url: '/indexer/importer/', app: 'importer' },

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

@@ -431,6 +431,7 @@ def _envelope_job(request, file_format, collection_name, start_time=None, lib_pa
     manager = ManagerApi()
 
     properties = {
+      "inputFormat": file_format['inputFormat'],
       "brokers": manager.get_kafka_brokers(),
       "kudu_master": manager.get_kudu_master(),
       "output_table": "impala::%s" % collection_name,

+ 57 - 64
desktop/libs/indexer/src/indexer/indexers/envelope.py

@@ -77,41 +77,49 @@ class EnvelopeIndexer(object):
 
 
   def generate_config(self, properties):
-    return """
-application {
-    name = Traffic analysis
-    batch.milliseconds = 5000
-    executors = 1
-    executor.cores = 1
-    executor.memory = 1G
-}
-
-steps {
-    traffic {
-        input {
-            type = kafka
-            brokers = "%(brokers)s"
-            topics = %(topics)s
-            encoding = string
-            translator {
-                type = delimited
-                delimiter = ","
-                field.names = [measurement_time,number_of_vehicles]
-                field.types = [long,int]
-            }
-            window {
-                enabled = true
-                milliseconds = 60000
-            }
-        }
+    if properties['inputFormat'] == 'kafka':
+      input = """            type = kafka
+              brokers = "%(brokers)s"
+              topics = %(topics)s
+              encoding = string
+              translator {
+                  type = delimited
+                  delimiter = ","
+                  field.names = [measurement_time,number_of_vehicles]
+                  field.types = [long,int]
+              }
+              window {
+                  enabled = true
+                  milliseconds = 60000
+              }
+      """ % properties
+    else: # File
+      input = """      type = filesystem
+      path = example-input.json
+      format = json
+      """ % properties
+      
+    if properties['ouputFormat'] == 'file':
+      # parquet, 
+      output = """    dependencies = [inputdata]
+    deriver {
+      type = sql
+      query.literal = "SELECT * FROM inputdata"
     }
-
-    trafficwindow {
-        dependencies = [traffic]
+    planner = {
+      type = overwrite
+    }
+    output = {
+      type = filesystem
+      path = example-output
+      format = csv
+    }"""
+    else: # Table
+      output = """        dependencies = [inputdata]
         deriver {
             type = sql
             query.literal = \"""
-                SELECT measurement_time, number_of_vehicles FROM traffic\"""
+                SELECT measurement_time, number_of_vehicles FROM inputdata\"""
         }
         planner {
             type = upsert
@@ -120,42 +128,27 @@ steps {
             type = kudu
             connection = "%(kudu_master)s"
             table.name = "%(output_table)s"
-        }
-    }
+        }""" % properties
+      
+    return """
+application {
+    name = Traffic analysis
+    batch.milliseconds = 5000
+    executors = 1
+    executor.cores = 1
+    executor.memory = 1G
 }
 
-""" % properties
-
-
-  def generate_config_parquet(self, properties):
-    return """application {
-  name = Filesystem Example
-  executors = 1
-}
 steps {
-  fsInput {
-    input {
-      type = filesystem
-      // Be sure to load this file into HDFS first!
-      path = example-input.json
-      format = json
-    }
-  }
-  fsProcess {
-    dependencies = [fsInput]
-    deriver {
-      type = sql
-      query.literal = "SELECT foo FROM fsInput"
-    }
-    planner = {
-      type = overwrite
+    inputdata {
+        input {
+            %(input)s
+        }
     }
-    output = {
-      type = filesystem
-      // The output directory
-      path = example-output
-      format = parquet
+
+    outputdata {
+        %(output)s
     }
-  }
 }
-"""
+
+""" % {'input': input, 'output': output}

+ 10 - 6
desktop/libs/indexer/src/indexer/templates/importer.mako

@@ -332,13 +332,17 @@ ${ assist.assistPanel() }
                 ## <select data-bind="selectize: createWizard.source.kafkaTopics, value: createWizard.source.kafkaSelectedTopics" placeholder="${ _('The list of topics to consume, e.g. orders,returns') }"></select>
               </label>
               
-              <label class="control-label"><div>${ _('Field names') }</div>
-                <input type="text" class="input-xxlarge" data-bind="value: createWizard.source.kafkaFieldNames" placeholder="${ _('The list of fields to consume, e.g. orders,returns') }">
-              </label>
+              <br/>
+              
+              <span data-bind="visible: createWizard.source.kafkaSelectedTopics">
+                <label class="control-label"><div>${ _('Field names') }</div>
+                  <input type="text" class="input-xxlarge" data-bind="value: createWizard.source.kafkaFieldNames" placeholder="${ _('The list of fields to consume, e.g. orders,returns') }">
+                </label>
 
-              <label class="control-label"><div>${ _('Field types') }</div>
-                <input type="text" class="input-xxlarge" data-bind="value: createWizard.source.kafkaFieldTypes" placeholder="${ _('The list of topics to consume, e.g. orders,returns') }">
-              </label>
+                <label class="control-label"><div>${ _('Field types') }</div>
+                  <input type="text" class="input-xxlarge" data-bind="value: createWizard.source.kafkaFieldTypes" placeholder="${ _('The list of topics to consume, e.g. orders,returns') }">
+                </label>
+              </span>
             </div>
           <!-- /ko -->
 

+ 6 - 3
desktop/libs/indexer/src/indexer/templates/topics.mako

@@ -58,7 +58,7 @@ ${ assist.assistPanel() }
   <h1>
     <!-- ko with: index() -->
     <div class="inline-block pull-right">
-      <a class="btn btn-default" href="javascript:void(0)" data-bind="hueLink: '/indexer/importer/prefill/all/index/' + name(), tooltip: { placement: 'bottom', delay: 750 }" title="${_('Inport stream data into a table or file')}">
+      <a class="btn btn-default" href="javascript:void(0)" data-bind="hueLink: '/indexer/importer/prefill/kafka/' + name(), tooltip: { placement: 'bottom', delay: 750 }" title="${_('Inport stream data into a table or file')}">
         <i class="fa fa-download fa-fw"></i> ${_('Consume')}
       </a>
 
@@ -91,7 +91,7 @@ ${ assist.assistPanel() }
           <ul class="nav">
             <li class="app-header">
               <a href="/${app_name}">
-                ${ _('Index Browser') }
+                ${ _('Kafka Browser') }
               </a>
             </li>
           </ul>
@@ -153,7 +153,7 @@ ${ assist.assistPanel() }
               </%def>
 
               <%def name="creation()">
-                <a href="javascript:void(0)" class="btn" data-bind="hueLink: '/indexer/importer/prefill/all/index/'">
+                <a href="javascript:void(0)" class="btn" data-bind="hueLink: '/indexer/importer/prefill/kafka'">
                   <i class="fa fa-plus-circle"></i> ${ _('Create') }
                 </a>
               </%def>
@@ -242,6 +242,9 @@ ${ assist.assistPanel() }
     <li class="active"><a href="#index-overview" data-toggle="tab" data-bind="click: function(){ $root.tab('index-overview'); }">${_('Overview')}</a></li>
     <li><a href="#index-columns" data-toggle="tab" data-bind="click: function(){ $root.tab('index-columns'); }">${_('Partitions')} (<span data-bind="text: fields().length"></span>)</a></li>
     <li><a href="#index-sample" data-toggle="tab" data-bind="click: function(){ $root.tab('index-sample'); }">${_('Sample')} (<span data-bind="text: sample().length"></span>)</a></li>
+    <li><a href="#index-sample" data-toggle="tab" data-bind="click: function(){ $root.tab('index-sample'); }">
+      ${_('Permissions')} (<span data-bind="text: sample().length"></span>)</a>
+    </li>    
   </ul>
 
   <div class="tab-content" style="border: none; overflow: hidden">