# Logstash kafka input to elasticsearch output stops consuming

**URL:** <https://discuss.elastic.co/t/logstash-kafka-input-to-elasticsearch-output-stops-consuming/42880>\
**Category:** Logstash\
**Created:** [February 26, 2016, 5:19pm UTC](https://discuss.elastic.co/t/logstash-kafka-input-to-elasticsearch-output-stops-consuming/42880 "2016-02-26T17:19:15Z")\
**Posts on this page:** 11\
**Page:** 1

<div class="post-metadata">

**Author:** ![zot42](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/zot42/32/8118_2.png) [@zot42](https://discuss.elastic.co/u/zot42)\
**Post date:** [February 26, 2016, 5:19pm UTC](https://discuss.elastic.co/t/logstash-kafka-input-to-elasticsearch-output-stops-consuming/42880/1 "2016-02-26T17:19:15Z")

</div>

Logstash 2.2  
elasticsearch 1.7.4

I am seeing and issue where I am consuming from kafka and sending to elasticsearch. It consumes about 500 messages at start and then stops with this messages being output to logs

`{:timestamp=>"2016-02-26T17:16:33.977000+0000", :message=>"Flushing buffer at interval", :instance=>"#<LogStash::Outputs::ElasticSearch::Buffer:0x3c356a00 @stopping=#<Concurrent::AtomicBoolean:0x56c244e1>, @last_flush=2016-02-26 17:16:32 +0000, @flush_thread=#<Thread:0x6c383c46 run>, @max_size=500, @operations_lock=#<Java::JavaUtilConcurrentLocks::ReentrantLock:0x5ffd4f2b>, @submit_proc=#<Proc:0x2bf99714@/opt/logstash/vendor/bundle/jruby/1.9/gems/logstash-output-elasticsearch-2.5.1-java/lib/logstash/outputs/elasticsearch/common.rb:57>, @flush_interval=1, @logger=#<Cabin::Channel:0x2e4d6053 @subscriber_lock=#<Mutex:0x59ef00e4>, @data={}, @metrics=#<Cabin::Metrics:0x5a23a6dc @channel=#<Cabin::Channel:0x2e4d6053 ...>, @metrics={}, @metrics_lock=#<Mutex:0x5a9fe3bf>>, @subscribers={13010=>#<Cabin::Outputs::IO:0x708f56ae @lock=#<Mutex:0x5283f33d>, @io=#<File:/mnt/log/logstash/logstash.log>>}, @level=:debug>, @buffer=[], @operations_mutex=#<Mutex:0x2b28e3f>>", :interval=>1, :level=>:debug, :file=>"logstash/outputs/elasticsearch/buffer.rb", :line=>"90", :method=>"interval_flush"}`

If i output to stdout or file , I am not seeing this issue.

---

<div class="post-metadata">

**Author:** ![Joe\_Lawson](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/joe_lawson/32/3390_2.png) [@Joe\_Lawson](https://discuss.elastic.co/u/Joe_Lawson)\
**Post date:** [February 26, 2016, 5:40pm UTC](https://discuss.elastic.co/t/logstash-kafka-input-to-elasticsearch-output-stops-consuming/42880/2 "2016-02-26T17:40:03Z")

</div>

Any errors on the Elasticsearch side?

---

<div class="post-metadata">

**Author:** ![zot42](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/zot42/32/8118_2.png) [@zot42](https://discuss.elastic.co/u/zot42)\
**Post date:** [February 26, 2016, 6:00pm UTC](https://discuss.elastic.co/t/logstash-kafka-input-to-elasticsearch-output-stops-consuming/42880/3 "2016-02-26T18:00:43Z")

</div>

no, all clean on es side

I have 10 node cluster with this as the only input, so there is no contention or load issues

---

<div class="post-metadata">

**Author:** ![Joe\_Lawson](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/joe_lawson/32/3390_2.png) [@Joe\_Lawson](https://discuss.elastic.co/u/Joe_Lawson)\
**Post date:** [February 26, 2016, 6:36pm UTC](https://discuss.elastic.co/t/logstash-kafka-input-to-elasticsearch-output-stops-consuming/42880/4 "2016-02-26T18:36:41Z")

</div>

Are you able to try 2.1.X? I wonder if there is a threading/pipeline issue. Also configs are helpful for reference.

---

<div class="post-metadata">

**Author:** ![zot42](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/zot42/32/8118_2.png) [@zot42](https://discuss.elastic.co/u/zot42)\
**Post date:** [February 26, 2016, 6:44pm UTC](https://discuss.elastic.co/t/logstash-kafka-input-to-elasticsearch-output-stops-consuming/42880/5 "2016-02-26T18:44:13Z")

</div>

not able to upgrade es yet, but doesn't seem like it's getting overwhelmed. the fact that it sends the first batch and then just stops with no error message seems like the connection between the input and output is borked.

I have the flush\_size artificially low to slow things down

```
input {
        kafka {
                topic_id => "logs"
                group_id => "log_consumers"
                zk_connect => "zoo01:2181,zoo02:2181,zoo03:2181,zoo04:2181,zoo05:2181/"
        }
}

filter {
        mutate {
                add_tag => ["kafka"]
        }
}

output {
  elasticsearch {
      hosts => ["es01:9200"]
      sniffing => true
      flush_size => 20
      idle_flush_time => 15
  }
}
```

---

<div class="post-metadata">

**Author:** ![zot42](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/zot42/32/8118_2.png) [@zot42](https://discuss.elastic.co/u/zot42)\
**Post date:** [February 26, 2016, 8:45pm UTC](https://discuss.elastic.co/t/logstash-kafka-input-to-elasticsearch-output-stops-consuming/42880/6 "2016-02-26T20:45:41Z")

</div>

So a bit more info, I downgraded to 1.5.5 logstash and it seems to be tied to the flush\_size. I set the flush\_size to 3000 and that's exactly how many messages get processed. I can see with the log in debug mode that messages are still being consumed from kafka, but they aren't making it to es. I also confirmed the same behavior with logstash 2.2 except that it is capped to 500, so if your flush\_size is below 500 it will process until 500, if it's more that 500 it stop at 500

before  
 ![](https://us1.discourse-cdn.com/elastic/original/2X/c/c0741510cd06af5f8b4c588fc774d8eedce2b1ec.png)

after  
 ![](https://us1.discourse-cdn.com/elastic/original/2X/9/919d0184e845a213c788d548a472bc91e8a8e8df.png)

---

<div class="post-metadata">

**Author:** ![Joe\_Lawson](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/joe_lawson/32/3390_2.png) [@Joe\_Lawson](https://discuss.elastic.co/u/Joe_Lawson)\
**Post date:** [February 26, 2016, 10:29pm UTC](https://discuss.elastic.co/t/logstash-kafka-input-to-elasticsearch-output-stops-consuming/42880/7 "2016-02-26T22:29:59Z")

</div>

I was searching the forums for 'interval\_flush' and it may get tripped up with networking connectivity. Could you describe the network/environment topology you are in? It could be that Logstash thinks the network connection is open but it was actually closed by a firewall which could lead to Kafka consuming messages and the output sending messages but tem not actually reaching your ES cluster.

Could you try restarting the Logstash server? Perhaps the ES cluster if needed. Also I meant trying Logstash version 2.1.X but 1.5.5 was sufficient for that. Sorry I wasn't more specific.

---

<div class="post-metadata">

**Author:** ![zot42](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/zot42/32/8118_2.png) [@zot42](https://discuss.elastic.co/u/zot42)\
**Post date:** [February 29, 2016, 4:35am UTC](https://discuss.elastic.co/t/logstash-kafka-input-to-elasticsearch-output-stops-consuming/42880/8 "2016-02-29T04:35:43Z")

</div>

I don't think that would be it, as it's very reproducible. I have fallen back to writing to elasticsearch while using the tcp/udp inputs. I haven't had time to dig into the code, but I think it's the interaction of and input and output that are both batch, because when i use what is essentially a stream input (tcp/udp/syslog) sending to elasticsearch or a kafka input and a stream output (file/stdout) things work as expected

---

<div class="post-metadata">

**Author:** ![Joe\_Lawson](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/joe_lawson/32/3390_2.png) [@Joe\_Lawson](https://discuss.elastic.co/u/Joe_Lawson)\
**Post date:** [February 29, 2016, 1:57pm UTC](https://discuss.elastic.co/t/logstash-kafka-input-to-elasticsearch-output-stops-consuming/42880/9 "2016-02-29T13:57:13Z")

</div>

do you see any errors on the kafka broker or zookeeper instances? What size are the messages? Any idea what the maximum message size is?

Sorry for all the questions bit this is puzzling. Could you share a payload of messages so that I may attempt to reproduce?

---

<div class="post-metadata">

**Author:** ![John\_Sotiropoulos](https://avatars.discourse-cdn.com/v4/letter/j/35a633/32.png) [@John\_Sotiropoulos](https://discuss.elastic.co/u/John_Sotiropoulos)\
**Post date:** [March 5, 2016, 4:53pm UTC](https://discuss.elastic.co/t/logstash-kafka-input-to-elasticsearch-output-stops-consuming/42880/10 "2016-03-05T16:53:54Z")

</div>

Hello, I have a similar problem which results to elasticsearch output stop consuming.

Specifically, I have 2x(kafka & logstash instances) =\> 3x ES instances.

Kafka and logstash are in the same node. I use kafka 0.9.0.1, logstash 2.2.1 and elasticsearch 2.2.0. (I have tried to put the zookeeper to the same node and to different node as well)

Logstash 2.2.2 kept giving some errors _[1]_ and _[2]_, so I switched to 2.2.1.

Nevertheless, there are some issues with kafka and logstash every now and then which give this warning in logstash output _[3]_ and it continues after some time. Also, while I was trying something with _geoip_ filter I got the same error with logstash 2.2.2 _[1]_.

I have looked it up on the internet but still I can't figure out what the problem is. Any input would be really appreciated!

Thank you in advance.

[1] _ERROR: org.I0Itec.zkclient.ZkEventThread: Error handling event ZkEvent[New session event sent to kafka.consumer.ZookeeperConsumerConnector$ZKSessionExpireListener@8dc07d2]_  
_kafka.common.ConsumerRebalanceFailedException: logstash\_broker1 can't rebalance after 4 retries_  
...

[2] _ERROR Error when sending message to topic connect-test with key: null, value: 128 bytes with error: Batch Expired (org.apache.kafka.clients.producer.internals.ErrorLoggingCallback)_

[3] _WARN: kafka.client.ClientUtils$: Fetching topic metadata with correlation id 21 for topics [Set(connect-test)] from broker $id:0,host:esuser02,port:9092] failed_  
_java.nio.channels.ClosedByInterruptException_  
\_ at java.nio.channels.spi.AbstractInterruptibleChannel.end(AbstractInterruptibleChannel.java:202)\_  
\_ at sun.nio.ch.SocketChannelImpl.write(SocketChannelImpl.java:511)\_  
.....

---

<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:08am UTC](https://discuss.elastic.co/t/logstash-kafka-input-to-elasticsearch-output-stops-consuming/42880/11 "2017-07-06T05:08:17Z")

</div>


