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=d53f5e57477bfb9dcc85ca7c98d5d123feebb26a;hb=e016643c4e11a5b71b9de533abbbf2a8e4d66dae;hp=49eca1a1592f0c51cde678dc6f7712411382d0b5;hpb=b1acb4e97b9672a682d2f6d2028a6efb003f2660;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 49eca1a..d53f5e5 100644 --- a/src/main/java/de/juplo/kafka/SimpleConsumer.java +++ b/src/main/java/de/juplo/kafka/SimpleConsumer.java @@ -19,7 +19,7 @@ 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 long consumed = 0; @@ -34,11 +34,11 @@ public class SimpleConsumer implements Runnable 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(