# Pushing data to kafka topic using logstash

**URL:** <https://discuss.elastic.co/t/pushing-data-to-kafka-topic-using-logstash/46766>\
**Category:** Logstash\
**Created:** [April 8, 2016, 6:30am UTC](https://discuss.elastic.co/t/pushing-data-to-kafka-topic-using-logstash/46766 "2016-04-08T06:30:31Z")\
**Posts on this page:** 8\
**Page:** 1

<div class="post-metadata">

**Author:** ![rohit\_prusty](https://avatars.discourse-cdn.com/v4/letter/r/b38774/32.png) [@rohit\_prusty](https://discuss.elastic.co/u/rohit_prusty)\
**Post date:** [April 8, 2016, 6:30am UTC](https://discuss.elastic.co/t/pushing-data-to-kafka-topic-using-logstash/46766/1 "2016-04-08T06:30:31Z")

</div>

We are using kafka-0.8.2.1 and integrating it with logstash-2.3.0 and elasticsearch-2.3.0.

We are able to read from the kafka topic using "logstash-input-kafka (2.0.6)" input plugin and store data in elasticsearch.

But, we are not able to push data to kafka topic using "logstash-output-kafka (2.0.3)" output plugin. We have tried with the below command. It is not showing any error but data is not pushed to kafka topic from console.

logstash -e "input { stdin {} } output { kafka { topic\_id =\> 'logstash\_logs' } }"

Please help on this.

---

<div class="post-metadata">

**Author:** ![warkolm](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/warkolm/32/39224_2.png) [@warkolm](https://discuss.elastic.co/u/warkolm)\
**Post date:** [April 9, 2016, 7:28am UTC](https://discuss.elastic.co/t/pushing-data-to-kafka-topic-using-logstash/46766/2 "2016-04-09T07:28:28Z")

</div>

Did you try with `--debug` to see what is happening?

---

<div class="post-metadata">

**Author:** ![rohit\_prusty](https://avatars.discourse-cdn.com/v4/letter/r/b38774/32.png) [@rohit\_prusty](https://discuss.elastic.co/u/rohit_prusty)\
**Post date:** [April 11, 2016, 4:15am UTC](https://discuss.elastic.co/t/pushing-data-to-kafka-topic-using-logstash/46766/3 "2016-04-11T04:15:09Z")

</div>

I have tried with "--debug" option and below is the result. I don't see any error and pipeline is getting created, but don't know why it is not writing to kafka topic. I am able to write to kafka topic by using [kafka-console-producer.sh](http://kafka-console-producer.sh).

bin/logstash -e "input { stdin {} } output { kafka { topic\_id =\> 'logstash\_logs' } }" --debug

Plugin not defined in namespace, checking for plugin file {:type=\>"input", :name=\>"stdin", :path=\>"logstash/inputs/stdin", :level=\>:debug, :file=\>"logstash/plugin.rb", :line=\>"76", :method=\>"lookup"}  
Plugin not defined in namespace, checking for plugin file {:type=\>"codec", :name=\>"line", :path=\>"logstash/codecs/line", :level=\>:debug, :file=\>"logstash/plugin.rb", :line=\>"76", :method=\>"lookup"}  
config LogStash::Codecs::Line/@charset = "UTF-8" {:level=\>:debug, :file=\>"logstash/config/mixin.rb", :line=\>"141", :method=\>"config\_init"}  
config LogStash::Codecs::Line/@delimiter = "\n" {:level=\>:debug, :file=\>"logstash/config/mixin.rb", :line=\>"141", :method=\>"config\_init"}  
config LogStash::Inputs::Stdin/@codec = \<LogStash::Codecs::Line charset=\>"UTF-8", delimiter=\>"\n"\> {:level=\>:debug, :file=\>"logstash/config/mixin.rb", :line=\>"141", :method=\>"config\_init"}  
config LogStash::Inputs::Stdin/@add\_field = {} {:level=\>:debug, :file=\>"logstash/config/mixin.rb", :line=\>"141", :method=\>"config\_init"}  
Plugin not defined in namespace, checking for plugin file {:type=\>"output", :name=\>"kafka", :path=\>"logstash/outputs/kafka", :level=\>:debug, :file=\>"logstash/plugin.rb", :line=\>"76", :method=\>"lookup"}  
starting agent {:level=\>:info, :file=\>"logstash/agent.rb", :line=\>"190", :method=\>"execute"}  
starting pipeline {:id=\>"main", :level=\>:info, :file=\>"logstash/agent.rb", :line=\>"444", :method=\>"start\_pipeline"}  
Settings: Default pipeline workers: 16  
config LogStash::Codecs::JSON/@charset = "UTF-8" {:level=\>:debug, :file=\>"logstash/config/mixin.rb", :line=\>"141", :method=\>"config\_init"}  
config LogStash::Outputs::Kafka/@topic\_id = "logstash\_logs" {:level=\>:debug, :file=\>"logstash/config/mixin.rb", :line=\>"141", :method=\>"config\_init"}  
log4j java properties setup {:log4j\_level=\>"DEBUG", :level=\>:debug, :file=\>"logstash/logging.rb", :line=\>"89", :method=\>"setup\_log4j"}  
Registering kafka producer {:topic\_id=\>"logstash\_logs", :bootstrap\_servers=\>"localhost:9092", :level=\>:info, :file=\>"logstash/outputs/kafka.rb", :line=\>"128", :method=\>"register"}  
Will start workers for output {:worker\_count=\>1, :class=\>LogStash::Outputs::Kafka, :level=\>:debug, :file=\>"logstash/output\_delegator.rb", :line=\>"77", :method=\>"register"}  
Starting pipeline {:id=\>"main", :pipeline\_workers=\>16, :batch\_size=\>125, :batch\_delay=\>5, :max\_inflight=\>2000, :level=\>:info, :file=\>"logstash/pipeline.rb", :line=\>"192", :method=\>"start\_workers"}  
Pipeline main started {:file=\>"logstash/agent.rb", :line=\>"448", :method=\>"start\_pipeline"}  
Pushing flush onto pipeline {:level=\>:debug, :file=\>"logstash/pipeline.rb", :line=\>"461", :method=\>"flush"}  
this is alPushing flush onto pipeline {:level=\>:debug, :file=\>"logstash/pipeline.rb", :line=\>"461", :method=\>"flush"}  
so not working  
filter received {:event=\>{"@timestamp"=\>2016-04-11T03:59:25.444Z, "@version"=\>"1", "host"=\>"[VPPCAZUSW00123.cs-DNA02-BIBDTZ1-RQ29519712.d9.internal.cloudapp.net](http://VPPCAZUSW00123.cs-DNA02-BIBDTZ1-RQ29519712.d9.internal.cloudapp.net)", "message"=\>"this is also not working"}, :level=\>:debug, :file=\>"(eval)", :line=\>"17", :method=\>"filter\_func"}  
output received {:event=\>{"@timestamp"=\>2016-04-11T03:59:25.444Z, "@version"=\>"1", "host"=\>"[VPPCAZUSW00123.cs-DNA02-BIBDTZ1-RQ29519712.d9.internal.cloudapp.net](http://VPPCAZUSW00123.cs-DNA02-BIBDTZ1-RQ29519712.d9.internal.cloudapp.net)", "message"=\>"this is also not working"}, :level=\>:debug, :file=\>"(eval)", :line=\>"22", :method=\>"output\_func"}  
Pushing flush onto pipeline {:level=\>:debug, :file=\>"logstash/pipeline.rb", :line=\>"461", :method=\>"flush"}

---

<div class="post-metadata">

**Author:** ![niharvarma](https://avatars.discourse-cdn.com/v4/letter/n/3bc359/32.png) [@niharvarma](https://discuss.elastic.co/u/niharvarma)\
**Post date:** [July 8, 2016, 7:22pm UTC](https://discuss.elastic.co/t/pushing-data-to-kafka-topic-using-logstash/46766/4 "2016-07-08T19:22:46Z")

</div>

Hi rohit, i'm also stuck with the same issue. Any luck so far ..?

---

<div class="post-metadata">

**Author:** ![allenmchan](https://avatars.discourse-cdn.com/v4/letter/a/ec9cab/32.png) [@allenmchan](https://discuss.elastic.co/u/allenmchan)\
**Post date:** [July 9, 2016, 2:15am UTC](https://discuss.elastic.co/t/pushing-data-to-kafka-topic-using-logstash/46766/5 "2016-07-09T02:15:48Z")

</div>

Neither of you mentioned if Kafka broker is on the same server as logstash output.

---

<div class="post-metadata">

**Author:** ![niharvarma](https://avatars.discourse-cdn.com/v4/letter/n/3bc359/32.png) [@niharvarma](https://discuss.elastic.co/u/niharvarma)\
**Post date:** [July 13, 2016, 5:46am UTC](https://discuss.elastic.co/t/pushing-data-to-kafka-topic-using-logstash/46766/6 "2016-07-13T05:46:41Z")

</div>

@allenmchan, No the kafka broker is different server

---

<div class="post-metadata">

**Author:** ![Mick\_Mahoney](https://avatars.discourse-cdn.com/v4/letter/m/cab0a1/32.png) [@Mick\_Mahoney](https://discuss.elastic.co/u/Mick_Mahoney)\
**Post date:** [July 13, 2016, 7:53am UTC](https://discuss.elastic.co/t/pushing-data-to-kafka-topic-using-logstash/46766/7 "2016-07-13T07:53:40Z")

</div>

You need to check that the logstash/kafka plugin is at the correct version for your kafka version

[https://www.elastic.co/guide/en/logstash/current/plugins-inputs-kafka.html](https://www.elastic.co/guide/en/logstash/current/plugins-inputs-kafka.html)

Additionally you need to specify the name of the remote server and the port number to connect to.  
Looks like it is defaulting to :bootstrap\_servers=\>"localhost:9092"  
So you need to specify your remote server and port in the logstash config

Cheers,  
Mick

---

<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, 4:48am UTC](https://discuss.elastic.co/t/pushing-data-to-kafka-topic-using-logstash/46766/8 "2017-07-06T04:48:17Z")

</div>


