浏览代码

HUE-8595 [flume] Collect and ingest Hue balancer logs out of the box

Romain Rigaux 7 年之前
父节点
当前提交
9c00548

+ 22 - 33
desktop/libs/indexer/src/data/morphline/hue_accesslogs_no_geo.morphline.conf

@@ -22,6 +22,8 @@ morphlines : [
 
         ## Hue HTTPD load balancer
         ## 172.18.18.3 - - [27/Aug/2018:05:47:12 -0700] "GET /static/desktop/js/jquery.rowselector.a04240f7cc48.js HTTP/1.1" 200 2321
+        ## hue access logs
+        ## [23/Sep/2018 22:11:35 -0700] INFO     172.31.114.129 -anon- - "HEAD /desktop/debug/is_alive HTTP/1.1" returned in 2ms
 
       separator:  " "
             columns:  [client_ip,C1,C2,time,dummy1,request,code,bytes]
@@ -33,43 +35,30 @@ morphlines : [
         }
     }
     {
-  split {
-    inputField : request
-    outputFields : [method, url, protocol]
-    separator : " "
-    isRegex : false
-    #separator : """\s*,\s*"""
-    #  #isRegex : true
-    addEmptyStrings : false
-    trim : true
-          }
-    }
-     {
-  split {
-    inputField : url
-    outputFields : ["", app, subapp]
-    separator : "\/"
-    isRegex : false
-    #separator : """\s*,\s*"""
-    #  #isRegex : true
-    addEmptyStrings : false
-    trim : true
-          }
+      split {
+        inputField : request
+        outputFields : [method, url, protocol]
+        separator : " "
+        isRegex : false
+        #separator : """\s*,\s*"""
+        #  #isRegex : true
+        addEmptyStrings : false
+        trim : true
+      }
     }
     {
-  userAgent {
-    inputField : user_agent
-    outputFields : {
-      user_agent_family : "@{ua_family}"
-      user_agent_major  : "@{ua_major}"
-      device_family     : "@{device_family}"
-      os_family         : "@{os_family}"
-      os_major    : "@{os_major}"
-    }
-  }
+      split {
+        inputField : url
+        outputFields : ["", app, subapp]
+        separator : "\/"
+        isRegex : false
+        #separator : """\s*,\s*"""
+        #  #isRegex : true
+        addEmptyStrings : false
+        trim : true
+      }
     }
 
-      #{logInfo { format : "BODY : {}", args : ["@{}"] } }
     # add Unique ID, in case our message_id field from above is not present
     {
         generateUUID {

+ 22 - 6
desktop/libs/indexer/src/indexer/api3.py

@@ -32,7 +32,6 @@ from desktop.lib.exceptions_renderable import PopupException
 from desktop.lib.i18n import smart_unicode
 from desktop.models import Document2
 from kafka.kafka_api import get_topics
-from librdbms.server import dbms as rdbms
 from metadata.manager_client import ManagerApi
 from notebook.connectors.base import get_api, Notebook
 from notebook.decorators import api_error_handler
@@ -43,7 +42,7 @@ from indexer.file_format import HiveFormat
 from indexer.fields import Field
 from indexer.indexers.envelope import EnvelopeIndexer
 from indexer.indexers.morphline import MorphlineIndexer
-from indexer.indexers.rdbms import run_sqoop,  _get_api
+from indexer.indexers.rdbms import run_sqoop, _get_api
 from indexer.indexers.sql import SQLIndexer
 from indexer.solr_client import SolrClient, MAX_UPLOAD_SIZE
 from indexer.indexers.flume import FlumeIndexer
@@ -144,7 +143,7 @@ def guess_format(request):
         password=file_format['streamPassword'],
         security_token=file_format['streamToken']
     )
-    format_ = {"type": "csv", "fieldSeparator": ",", "hasHeader": True, "quoteChar": "\"", "recordSeparator": "\\n", 'objects': [sobject['name'] for sobject in sf.restful('sobjects/')['sobjects'] if sobject['queryable']]}    
+    format_ = {"type": "csv", "fieldSeparator": ",", "hasHeader": True, "quoteChar": "\"", "recordSeparator": "\\n", 'objects': [sobject['name'] for sobject in sf.restful('sobjects/')['sobjects'] if sobject['queryable']]}
 
   format_['status'] = 0
   return JsonResponse(format_)
@@ -272,11 +271,28 @@ def guess_field_types(request):
         col['keyType'] = type_mapping[col['name']]
         col['type'] = type_mapping[col['name']]
     elif file_format['streamSelection'] == 'flume':
+      if 'hue-httpd/access_log' in file_format['channelSourcePath']:
+        columns = [
+          {'name': 'id', 'type': 'string', 'unique': True},
+          {'name': 'client_ip', 'type': 'string'},
+          {'name': 'time', 'type': 'date'},
+          {'name': 'request', 'type': 'string'},
+          {'name': 'code', 'type': 'plong'},
+          {'name': 'bytes', 'type': 'plong'},
+          {'name': 'method', 'type': 'string'},
+          {'name': 'url', 'type': 'string'},
+          {'name': 'protocol', 'type': 'string'},
+          {'name': 'app', 'type': 'string'},
+          {'name': 'subapp', 'type': 'string'}
+        ]
+      else:
+        columns = [{'name': 'message', 'type': 'string'}]
+
       format_ = {
-          "sample": [['...']] * 4,
+          "sample": [['...'] * len(columns)] * 4,
           "columns": [
-              Field(col['name'], HiveFormat.FIELD_TYPE_TRANSLATE.get(col['type'], 'string')).to_dict()
-              for col in [{'name': 'message', 'type': 'string'}]
+              Field(col['name'], HiveFormat.FIELD_TYPE_TRANSLATE.get(col['type'], 'string'), unique=col.get('unique')).to_dict()
+              for col in columns
           ]
       }
   elif file_format['streamSelection'] == 'sfdc':

+ 2 - 2
desktop/libs/indexer/src/indexer/fields.py

@@ -44,13 +44,13 @@ class FieldType():
 
 class Field(object):
 
-  def __init__(self, name="new_field", field_type_name="string", operations=None, multi_valued=False):
+  def __init__(self, name="new_field", field_type_name="string", operations=None, multi_valued=False, unique=False):
     self.name = name
     self.field_type_name = field_type_name
     self.keep = True
     self.operations = operations if operations else []
     self.required = False
-    self.unique = False
+    self.unique = unique
     self.multi_valued = multi_valued
     self.show_properties = False
 

+ 13 - 6
desktop/libs/indexer/src/indexer/indexers/flume.py

@@ -47,7 +47,7 @@ class FlumeIndexer(object):
 
     responses['refresh_flume'] = api.refresh_flume(cluster_name=None, restart=True)
 
-    if file_format['ouputFormat'] == 'index':
+    if destination['ouputFormat'] == 'index':
       responses['pubSubUrl'] = 'assist.collections.refresh'
       responses['on_success_url'] = reverse('search:browse', kwargs={'name': destination_name})
 
@@ -60,9 +60,11 @@ class FlumeIndexer(object):
     if source['channelSourceType'] == 'directory':
       agent_source = '''
   tier1.sources.source1.type = exec
-  tier1.sources.source1.command = tail -F /var/log/hue-httpd/access_log
+  tier1.sources.source1.command = tail -F %(directory)s
   tier1.sources.source1.channels = channel1
-      '''
+      ''' % {
+       'directory': source['channelSourcePath']
+    }
     else:
       raise PopupException(_('Input format not recognized: %(channelSourceType)s') % source)
 
@@ -97,12 +99,16 @@ class FlumeIndexer(object):
   a1.sinks.k1.serializer.serdeSeparator = '\t'
   a1.sinks.k1.serializer.fieldnames =id,,msg'''
     elif destination['ouputFormat'] == 'kafka':
+      manager = ManagerApi()
       agent_sink = '''
       tier1.sinks.sink1.type = org.apache.flume.sink.kafka.KafkaSink
 tier1.sinks.sink1.topic = hueAccessLogs
-tier1.sinks.sink1.brokerList = spark2-envelope515-1.gce.cloudera.com:9092,spark2-envelope515-2.gce.cloudera.com:9092,spark2-envelope515-3.gce.cloudera.com:9092
+tier1.sinks.sink1.brokerList = %(brokers)s
 tier1.sinks.sink1.channel = channel1
-tier1.sinks.sink1.batchSize = 20'''
+tier1.sinks.sink1.batchSize = 20''' % {
+      'brokers': manager.get_kafka_brokers()
+    }
+
     elif destination['ouputFormat'] == 'index':
       # Morphline file
       configs.append(self.generate_morphline_config(destination))
@@ -121,6 +127,7 @@ tier1.sinks.sink1.batchSize = 20'''
   tier1.channels = channel1
   tier1.sinks = sink1
 
+  %(sources)s
 
   tier1.channels.channel1.type = memory
   tier1.channels.channel1.capacity = 10000
@@ -140,7 +147,7 @@ tier1.sinks.sink1.batchSize = 20'''
     # TODO manage generic config, cf. MorphlineIndexer
     morphline_config = open(os.path.join(config_morphline_path(), 'hue_accesslogs_no_geo.morphline.conf')).read()
     morphline_config = morphline_config.replace(
-      '${SOLR_COLLECTION}', 'log_analytics_demo'
+      '${SOLR_COLLECTION}', destination['name']
     ).replace(
       '${ZOOKEEPER_ENSEMBLE}', '%s/solr' % zkensemble()
     )

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

@@ -1892,7 +1892,7 @@ ${ assist.assistPanel() }
           viewModel.createWizard.guessFieldTypes();
         }
       })
-      self.channelSourcePath = ko.observable('/var/log/hue/access.log');
+      self.channelSourcePath = ko.observable('/var/log/hue-httpd/access_log');
 
       self.streamUsername = ko.observable('');
       self.streamPassword = ko.observable('');