X-Git-Url: http://juplo.de/gitweb/?a=blobdiff_plain;f=src%2Fmain%2Fjava%2Fde%2Fjuplo%2Fkafka%2FEndlessConsumer.java;h=e3a60b5b796fea7dbeb4c2434c3daeef63976f84;hb=32a052b7c494009c59190857984ef3563f4f2b14;hp=be371aec1a157202185ed19268d2c6620abfdb2c;hpb=f20231dcb36eaeee69e07192b5785bc59abc1dc1;p=demos%2Fkafka%2Ftraining diff --git a/src/main/java/de/juplo/kafka/EndlessConsumer.java b/src/main/java/de/juplo/kafka/EndlessConsumer.java index be371ae..e3a60b5 100644 --- a/src/main/java/de/juplo/kafka/EndlessConsumer.java +++ b/src/main/java/de/juplo/kafka/EndlessConsumer.java @@ -33,7 +33,7 @@ public class EndlessConsumer implements Runnable private KafkaConsumer consumer = null; private Future future = null; - private final Map> seen = new HashMap<>(); + private final Map seen = new HashMap<>(); public EndlessConsumer( @@ -77,15 +77,15 @@ public class EndlessConsumer implements Runnable partitions.forEach(tp -> { log.info("{} - removing partition: {}", id, tp); - Map removed = seen.remove(tp.partition()); - for (String key : removed.keySet()) + PartitionStatistics removed = seen.remove(tp); + for (KeyCounter counter : removed.getStatistics()) { log.info( "{} - Seen {} messages for partition={}|key={}", id, - removed.get(key), - removed, - key); + counter.getCounter(), + removed.getPartition(), + counter.getKey()); } }); } @@ -96,7 +96,7 @@ public class EndlessConsumer implements Runnable partitions.forEach(tp -> { log.info("{} - adding partition: {}", id, tp); - seen.put(tp.partition(), new HashMap<>()); + seen.put(tp, new PartitionStatistics(tp)); }); } }); @@ -121,16 +121,9 @@ public class EndlessConsumer implements Runnable record.value() ); - Integer partition = record.partition(); + TopicPartition partition = new TopicPartition(record.topic(), record.partition()); String key = record.key() == null ? "NULL" : record.key(); - Map byKey = seen.get(partition); - - if (!byKey.containsKey(key)) - byKey.put(key, 0); - - int seenByKey = byKey.get(key); - seenByKey++; - byKey.put(key, seenByKey); + seen.get(partition).increment(key); } } } @@ -151,7 +144,7 @@ public class EndlessConsumer implements Runnable } } - public Map> getSeen() + public Map getSeen() { return seen; }