9a69c8f034af2093764e0024ed9d0f5f6cc6937d
[demos/kafka/training] / src / main / java / de / juplo / kafka / WordcountRebalanceListener.java
1 package de.juplo.kafka;
2
3 import lombok.RequiredArgsConstructor;
4 import lombok.extern.slf4j.Slf4j;
5 import org.apache.kafka.clients.consumer.Consumer;
6 import org.apache.kafka.clients.consumer.ConsumerRebalanceListener;
7 import org.apache.kafka.common.TopicPartition;
8
9 import java.util.Collection;
10 import java.util.Map;
11
12
13 @RequiredArgsConstructor
14 @Slf4j
15 public class WordcountRebalanceListener implements ConsumerRebalanceListener
16 {
17   private final WordcountRecordHandler handler;
18   private final PartitionStatisticsRepository repository;
19   private final String id;
20   private final Consumer<String, String> consumer;
21
22
23   @Override
24   public void onPartitionsAssigned(Collection<TopicPartition> partitions)
25   {
26     partitions.forEach(tp ->
27     {
28       Integer partition = tp.partition();
29       Long offset = consumer.position(tp);
30       log.info("{} - adding partition: {}, offset={}", id, partition, offset);
31       StatisticsDocument document =
32           repository
33               .findById(Integer.toString(partition))
34               .orElse(new StatisticsDocument(partition));
35       if (document.offset >= 0)
36       {
37         // Only seek, if a stored offset was found
38         // Otherwise: Use initial offset, generated by Kafka
39         consumer.seek(tp, document.offset);
40       }
41       handler.addPartition(partition, document.statistics);
42     });
43   }
44
45   @Override
46   public void onPartitionsRevoked(Collection<TopicPartition> partitions)
47   {
48     partitions.forEach(tp ->
49     {
50       Integer partition = tp.partition();
51       Long newOffset = consumer.position(tp);
52       log.info(
53           "{} - removing partition: {}, offset of next message {})",
54           id,
55           partition,
56           newOffset);
57       Map<String, Map<String, Long>> removed = handler.removePartition(partition);
58       repository.save(new StatisticsDocument(partition, removed, consumer.position(tp)));
59     });
60   }
61 }