核心流程
- 利用logstash查询Elasticsearch.
- 再利用match, mutate提取必要信息.
- 之后利用ruby执行本地shell或者命令获取输出返回值
- 利用aggregate将多个event合并为一个
- 最后发送邮件或者输出
注意, es查询到多条数据在logstash中算是多个event. 如果不做aggregate的话, 查到三条数据就会输出三次, 查到几条就输出几次. 做了aggregate后,会将多个event合并. 但是一定要配合event.cancel.这样会阻止前面的event,只保留最后一次aggregate的event. 也就是说, 你查询到了3条数据是3个event, aggregate是1个event. 整个过程中我们实际产生了4个event.前三个被cancel了, 只有最后一个aggregate没有人cancel
配置如下
1input { 2 elasticsearch { 3 hosts => ["127.0.0.1:9200"] 4 index => "test_indices" 5 query =>'{"query": { 6 "bool": { 7 "filter": [ 8 { 9 "range": { 10 "date_search": { 11 "gte": "now-5d/d", 12 "lte": "now+1d/d" 13 } 14 } 15 }, 16 { 17 "term": { 18 "tags": "errors" 19 } 20 } 21 ] 22 } 23 }}' 24 #docinfo => true 25 # schedule => "* * * * *" 26 } 27} 28filter { 29 grok { 30 match => { 31 "message" => "at (?<class>net.ray.[A-Za-z.]+)\.(?<method>[A-Za-z.]+)\([A-Za-z.]+:(?<line_num>[0-9]+)\)" 32 33 } 34 } 35 mutate { 36 gsub => [ "class", "\." , "/" ] 37 } 38 mutate { 39 update=>{ "class" => "%{class}.java" } 40 } 41 42 ruby { 43 code => " 44 cls = event.get('class') 45 line=event.get('line_num') 46 commit_log = `git -C e:/workspace loge:/workspace/src/main/java/#{cls}` 47 blame_log = `git -C e:/workspace blame -L #{line},#{line} e:/workspace/src/main/java/#{cls}` 48 event.set('commit_log',commit_log) 49 event.set('blame_log',blame_log) 50 " 51 } 52 53 aggregate { 54 task_id => "%{fields}%{log_source}" 55 code => " 56 map['result'] ||= [] 57 map['result'] << 'Log Date:' + event.get('datetime_search')+ ' \ncommit_log: \n' + event.get('commit_log') + '\n\nblame_log: \n' + event.get('blame_log') + ' \n\nJava Stack Trace:\n' + event.get('message') 58 event.cancel() 59 " 60 push_previous_map_as_event => true 61 } 62 mutate { 63 remove_field => [ 64 "tags", "host", "sequence", "@version", "@timestamp" 65 ] 66 join => {"result" => "\n\n--------------------------------\n\n"} 67 } 68} 69output { 70 stdout { 71 codec => rubydebug 72 } 73}