# Send data to elasticsearch index from kafka that is named the same as the topic

**URL:** <https://discuss.elastic.co/t/send-data-to-elasticsearch-index-from-kafka-that-is-named-the-same-as-the-topic/47365>\
**Category:** Logstash\
**Created:** [April 14, 2016, 9:24am UTC](https://discuss.elastic.co/t/send-data-to-elasticsearch-index-from-kafka-that-is-named-the-same-as-the-topic/47365 "2016-04-14T09:24:31Z")\
**Posts on this page:** 4\
**Page:** 1

<div class="post-metadata">

**Author:** ![amb](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/amb/32/9153_2.png) [@amb](https://discuss.elastic.co/u/amb)\
**Post date:** [April 14, 2016, 9:24am UTC](https://discuss.elastic.co/t/send-data-to-elasticsearch-index-from-kafka-that-is-named-the-same-as-the-topic/47365/1 "2016-04-14T09:24:31Z")

</div>

As it says in the title I want to pull out the topic name from the incoming event and write to an index in elasticsearchthat is named the same as that topic name.

> input {  
> kafka {  
> zk\_connect =\> "zookeeper:2181"  
> group\_id =\> "transport\_logstash\_elasticsearch\_local\_local"  
> white\_list =\> ".\*\_json"  
> decorate\_events =\> true  
> consumer\_threads =\> 1  
> queue\_size =\> 20  
> rebalance\_max\_retries =\> 4  
> rebalance\_backoff\_ms =\> 2000  
> consumer\_timeout\_ms =\> -1  
> consumer\_restart\_on\_error =\> true  
> consumer\_restart\_sleep\_ms =\> 0  
> fetch\_message\_max\_bytes =\> 1048576  
> reset\_beginning =\> false  
> type =\> "%{kafka.topic}"  
> }  
> }

> filter{

> #json {
> 
> # source =\> "kafka"
> 
> # }

> }  
> output {  
> elasticsearch{  
> action =\> "index"  
> codec =\> "plain"  
> flush\_size =\> 400  
> hosts =\> ["localhost:9200"]  
> idle\_flush\_time =\> 1  
> index =\> "%{kafka.topic}"  
> workers =\> 1  
> }  
> }

I thought the json filter plugin might be helpful (commented out in above code)but it seems to be unable to parse it:

> Error parsing json {:source=\>"kafka", :raw=\>{"msg\_size"=\>9, "topic"=\>"test\_json", "consumer\_group"=\>"transport\_logstash\_elasticsearch\_local\_local", "partition"=\>0, "offset"=\>45, "key"=\>nil}, :exception=\>java.lang.ClassCastException: org.jruby.RubyHash cannot be cast to org.jruby.RubyIO, :level=\>:warn}  
> ^CSIGINT received. Shutting down the agent. {:level=\>:warn}

Any help would be great. Thank you

---

<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:** [April 14, 2016, 9:29am UTC](https://discuss.elastic.co/t/send-data-to-elasticsearch-index-from-kafka-that-is-named-the-same-as-the-topic/47365/2 "2016-04-14T09:29:18Z")

</div>

It looks like the `kafka` field has already been parsed from JSON and that your json filter is unnecessary.

Also, please note that the notation for nested fields is `[field][subfield]` and not `field.subfield`.

---

<div class="post-metadata">

**Author:** ![amb](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/amb/32/9153_2.png) [@amb](https://discuss.elastic.co/u/amb)\
**Post date:** [April 14, 2016, 9:36am UTC](https://discuss.elastic.co/t/send-data-to-elasticsearch-index-from-kafka-that-is-named-the-same-as-the-topic/47365/3 "2016-04-14T09:36:35Z")

</div>

Thank you very much Magnus - it works perfectly now.

---

<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:** [July 6, 2017, 5:02am UTC](https://discuss.elastic.co/t/send-data-to-elasticsearch-index-from-kafka-that-is-named-the-same-as-the-topic/47365/4 "2017-07-06T05:02:18Z")

</div>


