X-Git-Url: http://juplo.de/gitweb/?a=blobdiff_plain;f=src%2Fmain%2Fjava%2Fde%2Fjuplo%2Fkafka%2FSimpleConsumer.java;fp=src%2Fmain%2Fjava%2Fde%2Fjuplo%2Fkafka%2FSimpleConsumer.java;h=49eca1a1592f0c51cde678dc6f7712411382d0b5;hb=b1acb4e97b9672a682d2f6d2028a6efb003f2660;hp=72dc36efc91ddf631f8b3d5c4c6c7c445a987645;hpb=6e38da42146210b24606d6dd6262ea5ccfb15f09;p=demos%2Fkafka%2Ftraining diff --git a/src/main/java/de/juplo/kafka/SimpleConsumer.java b/src/main/java/de/juplo/kafka/SimpleConsumer.java index 72dc36e..49eca1a 100644 --- a/src/main/java/de/juplo/kafka/SimpleConsumer.java +++ b/src/main/java/de/juplo/kafka/SimpleConsumer.java @@ -31,7 +31,6 @@ public class SimpleConsumer implements Runnable { log.info("{} - Subscribing to topic test", id); consumer.subscribe(Arrays.asList(topic)); - running = true; while (true) { @@ -65,7 +64,6 @@ public class SimpleConsumer implements Runnable } finally { - running = false; log.info("{} - Closing the KafkaConsumer", id); consumer.close(); log.info("{}: Consumed {} messages in total, exiting!", id, consumed);