Created
October 24, 2018 11:12
-
-
Save djptek/52175df2262072a41252d7761b8e72bc to your computer and use it in GitHub Desktop.
Query Elasticsearch with curl, parse to NDJSON and then ingest via Logstash
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| ##### Simulate scenario by indexing some example docs to local Elasticsearch | |
| POST index-2018-08-09/logs_2 | |
| { | |
| "datetime_log": "2018-08-09T00:34:36.051+02:00", | |
| "datetime_receive": "2018-08-09T00:34:36.051+02:00", | |
| "group": "DEFAUT", | |
| "ip_host": "22.33.44.55", | |
| "ip_host_pkt": "22.33.44.55", | |
| "source_msg": "22.33.44.55: -Trashback= XXXXXXXXX XXXXXXXXXX XXXXXXXXX 8002DAD8 YYYYYYYY ZZZZZZZZZ UUUUUUUU IIIIIIIIII", | |
| "unix_level": "local7", | |
| "unix_priority": "crit" | |
| } | |
| POST index-2018-08-09/logs_2 | |
| { | |
| "datetime_log": "2017-08-09T00:34:36.051+02:00", | |
| "datetime_receive": "2017-08-09T00:34:36.051+02:00", | |
| "group": "DEFAUT", | |
| "ip_host": "11.33.44.55", | |
| "ip_host_pkt": "11.33.44.55", | |
| "source_msg": "11.33.44.55: -Trashback= XXXXXXXXX XXXXXXXXXX XXXXXXXXX 8002DAD8 YYYYYYYY ZZZZZZZZZ UUUUUUUU IIIIIIIIII", | |
| "unix_level": "local7", | |
| "unix_priority": "crit" | |
| } | |
| ##### simulate extraction of the hits from ES using curl and pipe through sed+awk+grep to generate NDJSON | |
| curl -X GET http://0.0.0.0:9200/index-2018-08-09/_search | sed -e "s/^.*_source\"://" | grep -v "}]}}" | awk '/^}/ {print (NR==1?"":RS)$0;next} {printf "%s",$0}' | grep -v "^}$" | sed -e "s/$/\}/" > /Users/Shared/logs/test.json | |
| ##### expected output | |
| $ cat /Users/Shared/logs/test.json | |
| {"datetime_log": "2018-08-09T00:34:36.051+02:00","datetime_receive": "2018-08-09T00:34:36.051+02:00","group": "DEFAUT","ip_host": "22.33.44.55","ip_host_pkt": "22.33.44.55","source_msg": "22.33.44.55: -Trashback= XXXXXXXXX XXXXXXXXXX XXXXXXXXX 8002DAD8 YYYYYYYY ZZZZZZZZZ UUUUUUUU IIIIIIIIII","unix_level": "local7","unix_priority": "crit"} | |
| {"datetime_log": "2017-08-09T00:34:36.051+02:00","datetime_receive": "2017-08-09T00:34:36.051+02:00","group": "DEFAUT","ip_host": "11.33.44.55","ip_host_pkt": "11.33.44.55","source_msg": "11.33.44.55: -Trashback= XXXXXXXXX XXXXXXXXXX XXXXXXXXX 8002DAD8 YYYYYYYY ZZZZZZZZZ UUUUUUUU IIIIIIIIII","unix_level": "local7","unix_priority": "crit"} | |
| ##### configure logstash to ingest using the json codec, rectify fields using filter and index back to a new index in Elasticsearch | |
| $ cat config/ls-151840.conf | |
| input { | |
| # Read all documents from Elasticsearch matching the given query | |
| file { | |
| path => "/Users/Shared/logs/*.json" | |
| codec => json | |
| } | |
| } | |
| filter { | |
| mutate { | |
| rename => {"datetime_receive" => "date"} | |
| rename => {"datetime_log" => "datetime_log"} | |
| rename => {"group" => "group"} | |
| rename => {"ip_host_pkt" => "host"} | |
| rename => {"source_msg" => "message"} | |
| remove_field => ["unix_priority","@version","tags","@timestamp"] | |
| } | |
| } | |
| output { | |
| elasticsearch { | |
| hosts => ["http://localhost:9200"] | |
| index => "new_logs" | |
| } | |
| } | |
| ##### start logstash and wait for the documents to be ingested | |
| $ bin/logstash -f config/ls-151840.conf | |
| ##### check the new documents in Elasticsearch | |
| { | |
| "took": 4, | |
| "timed_out": false, | |
| "_shards": { | |
| "total": 5, | |
| "successful": 5, | |
| "skipped": 0, | |
| "failed": 0 | |
| }, | |
| "hits": { | |
| "total": 2, | |
| "max_score": 1, | |
| "hits": [ | |
| { | |
| "_index": "new_logs", | |
| "_type": "doc", | |
| "_id": "gZG7pWYB4d5XF6faLFF9", | |
| "_score": 1, | |
| "_source": { | |
| "unix_level": "local7", | |
| "date": "2017-08-09T00:34:36.051+02:00", | |
| "datetime_log": "2017-08-09T00:34:36.051+02:00", | |
| "ip_host": "11.33.44.55", | |
| "group": "DEFAUT", | |
| "path": "/Users/Shared/logs/test.json", | |
| "message": "11.33.44.55: -Trashback= XXXXXXXXX XXXXXXXXXX XXXXXXXXX 8002DAD8 YYYYYYYY ZZZZZZZZZ UUUUUUUU IIIIIIIIII", | |
| "host": "11.33.44.55" | |
| } | |
| }, | |
| { | |
| "_index": "new_logs", | |
| "_type": "doc", | |
| "_id": "gJG7pWYB4d5XF6faLFEa", | |
| "_score": 1, | |
| "_source": { | |
| "unix_level": "local7", | |
| "date": "2018-08-09T00:34:36.051+02:00", | |
| "datetime_log": "2018-08-09T00:34:36.051+02:00", | |
| "ip_host": "22.33.44.55", | |
| "group": "DEFAUT", | |
| "path": "/Users/Shared/logs/test.json", | |
| "message": "22.33.44.55: -Trashback= XXXXXXXXX XXXXXXXXXX XXXXXXXXX 8002DAD8 YYYYYYYY ZZZZZZZZZ UUUUUUUU IIIIIIIIII", | |
| "host": "22.33.44.55" | |
| } | |
| } | |
| ] | |
| } | |
| } |
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment