From: Kai Moritz Date: Sun, 24 Jul 2022 14:15:23 +0000 (+0200) Subject: Ausgabe der verarbeiteten Nachrichten im Revoke-Callback entfernt X-Git-Tag: endless-stream-consumer-DEPRECATED^2^2^2~1^2~10 X-Git-Url: http://juplo.de/gitweb/?a=commitdiff_plain;h=34d37c55d7cf830c6d2bdaf747f0938eb557bef3;p=demos%2Fkafka%2Ftraining Ausgabe der verarbeiteten Nachrichten im Revoke-Callback entfernt * Es musste allein für diese Ausgabe eine Map mit den zuletzt eingelesenen Offset-Positionen gepflegt werden. * Das ist zu viel Overhead, für die Randmeldung im Log. --- diff --git a/src/main/java/de/juplo/kafka/EndlessConsumer.java b/src/main/java/de/juplo/kafka/EndlessConsumer.java index 6e460b4..7e243a9 100644 --- a/src/main/java/de/juplo/kafka/EndlessConsumer.java +++ b/src/main/java/de/juplo/kafka/EndlessConsumer.java @@ -35,7 +35,6 @@ public class EndlessConsumer implements ConsumerRebalanceListener, Runnabl private long consumed = 0; private final Map> seen = new HashMap<>(); - private final Map lastOffsets = new HashMap<>(); @Override @@ -45,13 +44,10 @@ public class EndlessConsumer implements ConsumerRebalanceListener, Runnabl { Integer partition = tp.partition(); Long newOffset = consumer.position(tp); - Long oldOffset = lastOffsets.remove(partition); log.info( - "{} - removing partition: {}, consumed {} records (offset {} -> {})", + "{} - removing partition: {}, offset of next message {})", id, partition, - newOffset - oldOffset, - oldOffset, newOffset); Map removed = seen.remove(partition); for (String key : removed.keySet()) @@ -80,7 +76,6 @@ public class EndlessConsumer implements ConsumerRebalanceListener, Runnabl .findById(Integer.toString(partition)) .orElse(new StatisticsDocument(partition)); consumer.seek(tp, document.offset); - lastOffsets.put(partition, document.offset); seen.put(partition, document.statistics); }); }