# Logstash error while importing datas from kafka

**URL:** <https://discuss.elastic.co/t/logstash-error-while-importing-datas-from-kafka/71855>\
**Category:** Logstash\
**Created:** [January 17, 2017, 11:58am UTC](https://discuss.elastic.co/t/logstash-error-while-importing-datas-from-kafka/71855 "2017-01-17T11:58:30Z")\
**Posts on this page:** 9\
**Page:** 1

<div class="post-metadata">

**Author:** ![Grenouille06](https://avatars.discourse-cdn.com/v4/letter/g/48db29/32.png) [@Grenouille06](https://discuss.elastic.co/u/Grenouille06)\
**Post date:** [January 17, 2017, 11:58am UTC](https://discuss.elastic.co/t/logstash-error-while-importing-datas-from-kafka/71855/1 "2017-01-17T11:58:30Z")

</div>

Hi friends, I use logstash 2.4, kafka 0.10.1 and a cluster Elasticsearch 5.1.1

I import from Active Directory three logs: Application logs, System logs and Security logs.  
I create three topics in zookeeper from that logs :  
\_ topic\_id =\> "ActiveDirectory-Application-Logs"  
\_ topic\_id =\> "ActiveDirectory-System-Logs"  
\_ topic\_id =\> "ActiveDirectory-Security-Logs"

For storing datas, I use three indices in Elasticsearch:  
\_ index =\> "logstash-application"  
\_ index =\> "logstash-system"  
\_ index =\> "logstash-security"

In /etc/logstash/conf.d/ I have three conf files:

**logstash-application.conf**  
input {  
kafka {  
zk\_connect =\> ["[192.xxx.xxx.xxx:2181](http://192.xxx.xxx.xxx:2181)"]  
group\_id =\> "logstash-application"  
topic\_id =\> "ActiveDirectory-Application-Logs"  
reset\_beginning =\> "false"  
consumer\_threads =\> 1  
codec =\> json {}  
}  
}  
output {  
elasticsearch {  
user =\> logstash\_internal  
password =\> xxxxxxxxx  
hosts =\> ["[192.xxx.xxx.xxx:9200](http://192.xxx.xxx.xxx:9200)"]  
index =\> "logstash-application"}  
}

**logstash-system.conf**  
input {  
kafka {  
zk\_connect =\> ["[192.xxx.xxx.xxx:2181](http://192.xxx.xxx.xxx:2181)"]  
group\_id =\> "logstash-system"  
topic\_id =\> "ActiveDirectory-System-Logs"  
reset\_beginning =\> "false"  
consumer\_threads =\> 1  
codec =\> json {}  
}  
}  
output {  
elasticsearch {  
user =\> logstash\_internal  
password =\> xxxxxxxxxx  
hosts =\> ["[192.xxx.xxx.xxx:9200](http://192.xxx.xxx.xxx:9200)"]  
index =\> "logstash-system"}  
}

**logstash-security.conf**  
input {  
kafka {  
zk\_connect =\> ["[192.xxx.xxx.xxx:2181](http://192.xxx.xxx.xxx:2181)"]  
group\_id =\> "logstash-security"  
topic\_id =\> "ActiveDirectory-Security-Logs"  
reset\_beginning =\> "false"  
consumer\_threads =\> 1  
codec =\> json {}  
}  
}  
output {  
elasticsearch {  
user =\> logstash\_internal  
password =\> xxxxxxxx  
hosts =\> ["[192.xxx.xxx.xxx:9200](http://192.xxx.xxx.xxx:9200)"]  
index =\> "logstash-security"}  
}

When I start the pipeline, datas are stored in elasticsearch. But in my logstash-security index, there is some mistakes, some system or application logs are present in security, or security logs in application.

I don't know how to force logs to go in the right index.

Thanks for your help, and sorry for my english.

Fayce

---

<div class="post-metadata">

**Author:** ![magnusbaeck](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/magnusbaeck/32/44943_2.png) [@magnusbaeck](https://discuss.elastic.co/u/magnusbaeck)\
**Post date:** [January 17, 2017, 12:25pm UTC](https://discuss.elastic.co/t/logstash-error-while-importing-datas-from-kafka/71855/2 "2017-01-17T12:25:10Z")

</div>

Logstash has a single event pipeline where events from all inputs go to all outputs. Separating inputs and outputs in different files does not change this. If you want an event to only reach some outputs you need to wrap the outputs in a conditionals to selects how events are routed.

This exact question comes up here quite often. I'm sure you'll find useful information in the archives.

---

<div class="post-metadata">

**Author:** ![Grenouille06](https://avatars.discourse-cdn.com/v4/letter/g/48db29/32.png) [@Grenouille06](https://discuss.elastic.co/u/Grenouille06)\
**Post date:** [January 17, 2017, 12:26pm UTC](https://discuss.elastic.co/t/logstash-error-while-importing-datas-from-kafka/71855/3 "2017-01-17T12:26:59Z")

</div>

Ok thanks, I will search again 🙂

---

<div class="post-metadata">

**Author:** ![Grenouille06](https://avatars.discourse-cdn.com/v4/letter/g/48db29/32.png) [@Grenouille06](https://discuss.elastic.co/u/Grenouille06)\
**Post date:** [January 17, 2017, 12:38pm UTC](https://discuss.elastic.co/t/logstash-error-while-importing-datas-from-kafka/71855/4 "2017-01-17T12:38:06Z")

</div>

Do you think that unique file will work ?

input {  
kafka {  
zk\_connect =\> ["[192.xxx.xxx.xxx:2181](http://192.xxx.xxx.xxx:2181)"]  
topic\_id =\> "ActiveDirectory-System-Logs"  
topic\_id =\> "ActiveDirectory-Security-Logs"  
topic\_id =\> "ActiveDirectory-Application-Logs"]  
reset\_beginning =\> "false"  
consumer\_threads =\> 1  
codec =\> json {}  
}  
}

output {  
if [log\_name] == "Application" {  
elasticsearch {  
user =\> logstash\_internal  
password =\> xxxxxxxxxx  
embedded =\> true  
index =\> "logstash-application"  
}  
} else if [log\_name] == "System"{  
elasticsearch {  
user =\> logstash\_internal  
password =\> xxxxxxxxx  
embedded =\> true  
index =\> "logstash-system"  
}  
} else if [log\_name] == "Security"{  
user =\> logstash\_internal  
password =\> xxxxxxxxx  
embedded =\> true  
index =\> "logstash-security"  
}  
}

---

<div class="post-metadata">

**Author:** ![magnusbaeck](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/magnusbaeck/32/44943_2.png) [@magnusbaeck](https://discuss.elastic.co/u/magnusbaeck)\
**Post date:** [January 17, 2017, 1:11pm UTC](https://discuss.elastic.co/t/logstash-error-while-importing-datas-from-kafka/71855/5 "2017-01-17T13:11:33Z")

</div>

Yeah, follow that pattern (but you're missing an `elasticsearch {` line in the last `else if` block).

---

<div class="post-metadata">

**Author:** ![Grenouille06](https://avatars.discourse-cdn.com/v4/letter/g/48db29/32.png) [@Grenouille06](https://discuss.elastic.co/u/Grenouille06)\
**Post date:** [January 17, 2017, 1:24pm UTC](https://discuss.elastic.co/t/logstash-error-while-importing-datas-from-kafka/71855/6 "2017-01-17T13:24:32Z")

</div>

Here is my final conf file

input {  
kafka {  
zk\_connect =\> ["[192.xxx.xxx.xxx:2181](http://192.xxx.xxx.xxx:2181)"]  
group\_id =\> "logstash-application"  
topic\_id =\> "ActiveDirectory-Application-Logs"  
reset\_beginning =\> "false"  
consumer\_threads =\> 1  
codec =\> json {}  
}  
kafka {  
zk\_connect =\> ["[192.xxx.xxx.xxx:2181](http://192.xxx.xxx.xxx:2181)"]  
group\_id =\> "logstash-system"  
topic\_id =\> "ActiveDirectory-System-Logs"  
reset\_beginning =\> "false"  
consumer\_threads =\> 1  
codec =\> json {}  
}  
kafka {  
zk\_connect =\> ["[192.xxx.xxx.xxx:2181](http://192.xxx.xxx.xxx:2181)"]  
group\_id =\> "logstash-security"  
topic\_id =\> "ActiveDirectory-Security-Logs"  
reset\_beginning =\> "false"  
consumer\_threads =\> 1  
codec =\> json {}  
}  
}

output {  
if [log\_name] == "Application" {  
elasticsearch {  
user =\> logstash\_internal  
password =\> xxxxxx  
embedded =\> true  
index =\> "logstash-application"  
}  
} else if [log\_name] == "System"{  
elasticsearch {  
user =\> logstash\_internal  
password =\> xxxxxx  
embedded =\> true  
index =\> "logstash-system"  
}  
} else if [log\_name] == "Security"{  
elasticsearch {  
user =\> logstash\_internal  
password =\> xxxxxx  
embedded =\> true  
index =\> "logstash-security"  
}  
}  
}

But when I start the instance I have that error :

{:timestamp=\>"2017-01-17T14:21:05.034000+0100", :message=\>"fetched an invalid config", :config=\>"﻿input {\r\n kafka {\r\n zk\_connect =\> ["[192.xxx.xxx.xxx:2181](http://192.xxx.xxx.xxx:2181)"]\r\n group\_id =\> "logstash-application"\r\n topic\_id =\> "ActiveDirectory-Application-Logs"\r\n reset\_beginning =\> "false"\r\n consumer\_threads =\> 1\r\n codec =\> json {}\r\n}\r\n kafka {\r\n zk\_connect =\> ["[192.xxx.xxx.xxx:2181](http://192.xxx.xxx.xxx:2181)"]\r\n group\_id =\> "logstash-system"\r\n topic\_id =\> "ActiveDirectory-System-Logs"\r\n reset\_beginning =\> "false"\r\n consumer\_threads =\> 1\r\n codec =\> json {}\r\n}\r\n kafka {\r\n zk\_connect =\> ["[192.xxx.xxx.xxx:2181](http://192.xxx.xxx.xxx:2181)"]\r\n group\_id =\> "logstash-security"\r\n topic\_id =\> "ActiveDirectory-Security-Logs"\r\n reset\_beginning =\> "false"\r\n consumer\_threads =\> 1\r\n codec =\> json {}\r\n}\r\n}\r\n\r\noutput {\r\nif [log\_name] == "Application" {\r\n elasticsearch {\r\n user =\> logstash\_internal\r\n password =\> xxxxxx\r\n embedded =\> true\r\n index =\> "logstash-application"\r\n }\r\n } else if [log\_name] == "System"{\r\n elasticsearch {\r\n user =\> logstash\_internal\r\n password =\> xxxxxx\r\n embedded =\> true\r\n index =\> "logstash-system"\r\n }\r\n } else if [log\_name] == "Security"{\r\n elasticsearch {\r\n user =\> logstash\_internal\r\n password =\> xxxxxx\r\n embedded =\> true\r\n index =\> "logstash-security"\r\n}\r\n}\r\n}\n", :reason=\>"Expected one of #, input, filter, output at line 1, column 1 (byte 1) after ", :level=\>:error}

We are close 🙂 thanks for help

---

<div class="post-metadata">

**Author:** ![magnusbaeck](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/magnusbaeck/32/44943_2.png) [@magnusbaeck](https://discuss.elastic.co/u/magnusbaeck)\
**Post date:** [January 17, 2017, 1:51pm UTC](https://discuss.elastic.co/t/logstash-error-while-importing-datas-from-kafka/71855/7 "2017-01-17T13:51:58Z")

</div>

I don't believe the `embedded` option is supported in Logstash 2.0 and later, but it appears it complains about the very beginning of the file. Check that you don't have a garbage character there. Otherwise reduce the configuration until it works to narrow it down.

---

<div class="post-metadata">

**Author:** ![Grenouille06](https://avatars.discourse-cdn.com/v4/letter/g/48db29/32.png) [@Grenouille06](https://discuss.elastic.co/u/Grenouille06)\
**Post date:** [January 17, 2017, 2:50pm UTC](https://discuss.elastic.co/t/logstash-error-while-importing-datas-from-kafka/71855/8 "2017-01-17T14:50:24Z")

</div>

It works. you're right Magnus, I deleted the embedded option and used Hexdump - for garbage character.

Thanks a lot for help 🙂 🙂 🙂

Fayce

---

<div class="post-metadata">

**Author:** ![system](https://us1.discourse-cdn.com/elastic/original/3X/1/a/1ac57faf039f6b580b3f104ef42a2a89e41014de.png) [@system](https://discuss.elastic.co/u/system)\
**Post date:** [February 14, 2017, 2:50pm UTC](https://discuss.elastic.co/t/logstash-error-while-importing-datas-from-kafka/71855/9 "2017-02-14T14:50:26Z")

</div>

This topic was automatically closed 28 days after the last reply. New replies are no longer allowed.
