X-Git-Url: https://juplo.de/gitweb/?a=blobdiff_plain;f=src%2Fmain%2Fjava%2Fde%2Fjuplo%2Fkafka%2FApplicationRebalanceListener.java;fp=src%2Fmain%2Fjava%2Fde%2Fjuplo%2Fkafka%2FApplicationRebalanceListener.java;h=6776c0d6218c24f59e14ba356c52927014261730;hb=14ddf90be0adbeb8ab34b516b9b158071ae491e4;hp=32e14e8159774f1421cb65a045939ac27e7e0e28;hpb=a22627477d23328b168e06cac2a807db1f1145f2;p=demos%2Fkafka%2Ftraining diff --git a/src/main/java/de/juplo/kafka/ApplicationRebalanceListener.java b/src/main/java/de/juplo/kafka/ApplicationRebalanceListener.java index 32e14e8..6776c0d 100644 --- a/src/main/java/de/juplo/kafka/ApplicationRebalanceListener.java +++ b/src/main/java/de/juplo/kafka/ApplicationRebalanceListener.java @@ -51,14 +51,14 @@ public class ApplicationRebalanceListener implements PollIntervalAwareConsumerRe log.info("{} - removing partition: {}", id, partition); this.partitions.remove(partition); Map state = recordHandler.removePartition(partition); - for (String key : state.keySet()) + for (String user : state.keySet()) { log.info( - "{} - Seen {} messages for partition={}|key={}", + "{} - Calculations for partition={}|user={}: {}", id, - state.get(key), partition, - key); + user, + state.get(user)); } Map> results = adderResults.removePartition(partition); stateRepository.save(new StateDocument(partition, state, results));