# Kafka input not resuming where it stopped

**URL:** <https://discuss.elastic.co/t/kafka-input-not-resuming-where-it-stopped/31882>\
**Category:** Logstash\
**Created:** [October 8, 2015, 8:50pm UTC](https://discuss.elastic.co/t/kafka-input-not-resuming-where-it-stopped/31882 "2015-10-08T20:50:16Z")\
**Posts on this page:** 8\
**Page:** 1

<div class="post-metadata">

**Author:** ![mihir\_ray](https://avatars.discourse-cdn.com/v4/letter/m/e95f7d/32.png) [@mihir\_ray](https://discuss.elastic.co/u/mihir_ray)\
**Post date:** [October 8, 2015, 8:50pm UTC](https://discuss.elastic.co/t/kafka-input-not-resuming-where-it-stopped/31882/1 "2015-10-08T20:50:16Z")

</div>

HI,

I have a kafka cluster from where i am reading data and doing my processing. The input looks like:

input {  
kafka {  
topic\_id =\> "sb\_logs"  
codec =\> "plain"  
zk\_connect =\> "[cpanaetl01.sling.com:2181](http://cpanaetl01.sling.com:2181),[cpanaetl02.sling.com:2181](http://cpanaetl02.sling.com:2181),[cpanaetl03.sling.com:2181](http://cpanaetl03.sling.com:2181)"  
consumer\_threads =\> 20  
}  
}

If i stop logstash for few hours and start again, it reads only the recent data not resuming where it stopped.  
Is there anything wrong with my config?

---

<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:** [October 9, 2015, 11:57pm UTC](https://discuss.elastic.co/t/kafka-input-not-resuming-where-it-stopped/31882/2 "2015-10-09T23:57:27Z")

</div>

First thing, put the thread number down to 1. Kafka I'd already making  
threads per partition so there is no parallelism to be gained. We should  
remove that option from the plugin. Also set your consumer group name  
otherwise it'll assign a rand one each time which is why it isn't resuming.

Hope this helps!

Sincerely,

Joe

---

<div class="post-metadata">

**Author:** ![mihir\_ray](https://avatars.discourse-cdn.com/v4/letter/m/e95f7d/32.png) [@mihir\_ray](https://discuss.elastic.co/u/mihir_ray)\
**Post date:** [October 10, 2015, 10:17am UTC](https://discuss.elastic.co/t/kafka-input-not-resuming-where-it-stopped/31882/3 "2015-10-10T10:17:52Z")

</div>

That makes sense. Thanks.  
Please put this in your documentation as it says by default it will set the group name as logstash.

---

<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:** [October 10, 2015, 7:53pm UTC](https://discuss.elastic.co/t/kafka-input-not-resuming-where-it-stopped/31882/4 "2015-10-10T19:53:48Z")

</div>

You know, I was wrong about the group it is in fact logstash by default so it is strange that it wouldn't resume. Do you have other logstash consumers running with the same group?

---

<div class="post-metadata">

**Author:** ![mihir\_ray](https://avatars.discourse-cdn.com/v4/letter/m/e95f7d/32.png) [@mihir\_ray](https://discuss.elastic.co/u/mihir_ray)\
**Post date:** [October 11, 2015, 4:36am UTC](https://discuss.elastic.co/t/kafka-input-not-resuming-where-it-stopped/31882/5 "2015-10-11T04:36:09Z")

</div>

Yes, i have multiple consumers on the same group but for different topics.  
Sometimes i have seen, the logstash consumers do not make a offset entry in the zookeeper.

---

<div class="post-metadata">

**Author:** ![mihir\_ray](https://avatars.discourse-cdn.com/v4/letter/m/e95f7d/32.png) [@mihir\_ray](https://discuss.elastic.co/u/mihir_ray)\
**Post date:** [October 12, 2015, 7:03am UTC](https://discuss.elastic.co/t/kafka-input-not-resuming-where-it-stopped/31882/6 "2015-10-12T07:03:59Z")

</div>

HI,  
I did a validation on multi threading.

When not using consumer\_threads parameter :

/opt/kafka/kafka\_2.11-0.8.2.1/bin/kafka-run-class.sh kafka.tools.ConsumerOffsetChecker --topic test\_log --group analytics\_test\_log --zookeeper [etl01.abc.com:2181](http://etl01.abc.com:2181),[etl02.abc.com:2181](http://etl02.abc.com:2181),[etl03.abc.com:2181](http://etl03.abc.com:2181)  
Group Topic Pid Offset logSize Lag Owner  
analytics\_test\_log test\_log 0 198264791 199599448 1334657 analytics\_test\_log\_etl03.abc.com-1444631984622-cfd962f7-0  
analytics\_test\_log test\_log 1 167470850 168806166 1335316 analytics\_test\_log\_etl03.abc.com-1444631984622-cfd962f7-0  
analytics\_test\_log test\_log 2 190270643 192428588 2157945 analytics\_test\_log\_etl03.abc.com-1444631984622-cfd962f7-0  
analytics\_test\_log test\_log 3 183026805 185389546 2362741 analytics\_test\_log\_etl03.abc.com-1444631984622-cfd962f7-0  
analytics\_test\_log test\_log 4 174819301 177646047 2826746 analytics\_test\_log\_etl03.abc.com-1444631984622-cfd962f7-0  
analytics\_test\_log test\_log 5 188196243 189569061 1372818 analytics\_test\_log\_etl03.abc.com-1444631984622-cfd962f7-0  
analytics\_test\_log test\_log 6 174346801 178962394 4615593 analytics\_test\_log\_etl03.abc.com-1444631984622-cfd962f7-0  
analytics\_test\_log test\_log 7 171562338 172310342 748004 analytics\_test\_log\_etl03.abc.com-1444631984622-cfd962f7-0  
analytics\_test\_log test\_log 8 180201753 185593524 5391771 analytics\_test\_log\_etl03.abc.com-1444631984622-cfd962f7-0  
analytics\_test\_log test\_log 9 182661181 184789707 2128526 analytics\_test\_log\_etl03.abc.com-1444631984622-cfd962f7-0

When using cunsumer\_threads:

/opt/kafka/kafka\_2.11-0.8.2.1/bin/kafka-run-class.sh kafka.tools.ConsumerOffsetChecker --topic test\_log --group analytics\_test\_log --zookeeper [etl01.abc.com:2181](http://etl01.abc.com:2181),[etl02.abc.com:2181](http://etl02.abc.com:2181),[etl03.abc.com:2181](http://etl03.abc.com:2181)  
Group Topic Pid Offset logSize Lag Owner  
analytics\_test\_log test\_log 0 198347704 199599448 1251744 analytics\_test\_log\_etl03.abc.com-1444632839752-e081907b-0  
analytics\_test\_log test\_log 1 167559379 168806166 1246787 analytics\_test\_log\_etl03.abc.com-1444632839752-e081907b-1  
analytics\_test\_log test\_log 2 190354993 192428588 2073595 analytics\_test\_log\_etl03.abc.com-1444632839752-e081907b-2  
analytics\_test\_log test\_log 3 183109216 186047539 2938323 analytics\_test\_log\_etl03.abc.com-1444632839752-e081907b-3  
analytics\_test\_log test\_log 4 174903832 177646047 2742215 analytics\_test\_log\_etl03.abc.com-1444632839752-e081907b-4  
analytics\_test\_log test\_log 5 188281063 189569061 1287998 analytics\_test\_log\_etl03.abc.com-1444632839752-e081907b-5  
analytics\_test\_log test\_log 6 174432717 178962394 4529677 analytics\_test\_log\_etl03.abc.com-1444632839752-e081907b-6  
analytics\_test\_log test\_log 7 171649245 173122454 1473209 analytics\_test\_log\_etl03.abc.com-1444632839752-e081907b-7  
analytics\_test\_log test\_log 8 180285840 186082448 5796608 analytics\_test\_log\_etl03.abc.com-1444632839752-e081907b-8  
analytics\_test\_log test\_log 9 182747656 184789707 2042051 analytics\_test\_log\_etl03.abc.com-1444632839752-e081907b-9

If you see both the outputs, in the first case all the owner suffixes are same(0), in the second one it has suffixes from 1 to 10. Looks like logstash is not doing multi threading by default.

Please let me know if its not the case.

Thanks

---

<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:** [October 12, 2015, 1:13pm UTC](https://discuss.elastic.co/t/kafka-input-not-resuming-where-it-stopped/31882/7 "2015-10-12T13:13:47Z")

</div>

This is getting a little off topic but the tldr; is that the underlying  
jruby-kafka consumer thread isn't doing anything but multiplexing from the  
reader which creates a consumer thread per partition into the queue that  
logstash passes in. So yes there are more Kafka streams but each one isn't  
adding anything. Check out discussion here:

> <https://github.com/joekiller/jruby-kafka/issues/32#issuecomment-135224033>

When the latest major version of jruby-kafka is released then the  
logstash-kafka-input will need to change to pass a process into the Kafka  
stream which could add parallelism however even then it probably won't help  
much because logstash has a serialized processing chain, ie input  
queue\>filter queue\> output queue. If logstash were more like spark and  
maintained parallel threads for each chain then bumping the number of Kafka  
streams would help.

---

<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:26am UTC](https://discuss.elastic.co/t/kafka-input-not-resuming-where-it-stopped/31882/8 "2017-07-06T05:26:52Z")

</div>


