X-Git-Url: http://juplo.de/gitweb/?a=blobdiff_plain;f=src%2Fmain%2Fjava%2Fde%2Fjuplo%2Fkafka%2FSimpleConsumer.java;h=53bd11205692d125acfbbde85613150cb65411bf;hb=28861eab2d4da8a0594a115de989ffeb90b05cd4;hp=0d371f46caec01f20bcdd19dad6df7cdf4c86b7c;hpb=ffd5ad8116f8269ae828a7732cf2bd862f7ba095;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 0d371f4..53bd112 100644 --- a/src/main/java/de/juplo/kafka/SimpleConsumer.java +++ b/src/main/java/de/juplo/kafka/SimpleConsumer.java @@ -9,19 +9,16 @@ import org.apache.kafka.common.errors.WakeupException; import java.time.Duration; import java.util.Arrays; -import java.util.concurrent.ExecutorService; @Slf4j @RequiredArgsConstructor public class SimpleConsumer implements Runnable { - private final ExecutorService executor; private final String id; private final String topic; - private final Consumer consumer; + private final Consumer consumer; - private volatile boolean running = false; private long consumed = 0; @@ -31,16 +28,15 @@ public class SimpleConsumer implements Runnable try { log.info("{} - Subscribing to topic test", id); - consumer.subscribe(Arrays.asList("test")); - running = true; + consumer.subscribe(Arrays.asList(topic)); while (true) { - ConsumerRecords records = + ConsumerRecords records = consumer.poll(Duration.ofSeconds(1)); log.info("{} - Received {} messages", id, records.count()); - for (ConsumerRecord record : records) + for (ConsumerRecord record : records) { consumed++; log.info( @@ -66,15 +62,9 @@ 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); } } - - public void start() - { - executor.submit(this); - } }