X-Git-Url: https://juplo.de/gitweb/?a=blobdiff_plain;f=src%2Fmain%2Fjava%2Fde%2Fjuplo%2Fkafka%2FEndlessConsumer.java;h=b152310b3441316032b59ba8ff7948d3291c8c84;hb=3c1d3fa68df685146bdef7cc2e396e55fa0933dc;hp=69c82943522b50e30d228576d186d8bc86d37142;hpb=cb8505576575efa6957214c4d6cc4be777fd3b21;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 69c8294..b152310 100644 --- a/src/main/java/de/juplo/kafka/EndlessConsumer.java +++ b/src/main/java/de/juplo/kafka/EndlessConsumer.java @@ -135,6 +135,8 @@ public class EndlessConsumer implements Runnable String key = record.key() == null ? "NULL" : record.key(); seen.get(partition).increment(key); } + + seen.forEach((tp, statistics) -> repository.save(new StatisticsDocument(statistics, consumer.position(tp)))); } } catch(WakeupException e)