Skip to content

Instantly share code, notes, and snippets.

@djptek
Created October 24, 2018 11:12
Show Gist options
  • Select an option

  • Save djptek/52175df2262072a41252d7761b8e72bc to your computer and use it in GitHub Desktop.

Select an option

Save djptek/52175df2262072a41252d7761b8e72bc to your computer and use it in GitHub Desktop.
Query Elasticsearch with curl, parse to NDJSON and then ingest via Logstash
##### 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