X-Git-Url: https://juplo.de/gitweb/?a=blobdiff_plain;f=src%2Fmain%2Fjava%2Fde%2Fjuplo%2Fkafka%2FPartitionStatistics.java;h=0e31945d0d83f4cac156a6d1262b8b5d90e06cf4;hb=915674ec49ba38b3716cc4ef53272e963f139677;hp=e47a9f98d948247e69766103e9f7d4fd92149c26;hpb=32a052b7c494009c59190857984ef3563f4f2b14;p=demos%2Fkafka%2Ftraining diff --git a/src/main/java/de/juplo/kafka/PartitionStatistics.java b/src/main/java/de/juplo/kafka/PartitionStatistics.java index e47a9f9..0e31945 100644 --- a/src/main/java/de/juplo/kafka/PartitionStatistics.java +++ b/src/main/java/de/juplo/kafka/PartitionStatistics.java @@ -2,7 +2,6 @@ package de.juplo.kafka; import lombok.EqualsAndHashCode; import lombok.Getter; -import lombok.RequiredArgsConstructor; import org.apache.kafka.common.TopicPartition; import java.util.Collection; @@ -10,15 +9,35 @@ import java.util.HashMap; import java.util.Map; -@RequiredArgsConstructor @Getter @EqualsAndHashCode(of = { "partition" }) public class PartitionStatistics { + private String id; private final TopicPartition partition; private final Map statistics = new HashMap<>(); + public PartitionStatistics(TopicPartition partition) + { + this.partition = partition; + } + + public PartitionStatistics(StatisticsDocument document) + { + this.partition = new TopicPartition(document.topic, document.partition); + document + .statistics + .entrySet() + .forEach(entry -> + { + this.statistics.put( + entry.getKey(), + new KeyCounter(entry.getKey(), entry.getValue())); + }); + } + + public KeyCounter addKey(String key) { KeyCounter counter = new KeyCounter(key);