fix: In `onPartitionsAssigned()` wurde der Kafka-Offset ausgegeben
[demos/kafka/training] / src / main / java / de / juplo / kafka / AdderRebalanceListener.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.common.TopicPartition;
7
8 import java.time.Clock;
9 import java.time.Duration;
10 import java.time.Instant;
11 import java.util.Collection;
12
13
14 @RequiredArgsConstructor
15 @Slf4j
16 public class AdderRebalanceListener implements PollIntervalAwareConsumerRebalanceListener
17 {
18   private final AdderRecordHandler handler;
19   private final PartitionStatisticsRepository repository;
20   private final String id;
21   private final String topic;
22   private final Clock clock;
23   private final Duration commitInterval;
24   private final Consumer<String, String> consumer;
25
26   private Instant lastCommit = Instant.EPOCH;
27   private boolean commitsEnabled = true;
28
29   @Override
30   public void onPartitionsAssigned(Collection<TopicPartition> partitions)
31   {
32     partitions.forEach(tp ->
33     {
34       Integer partition = tp.partition();
35       StateDocument document =
36           repository
37               .findById(Integer.toString(partition))
38               .orElse(new StateDocument(partition));
39       log.info("{} - adding partition: {}, offset={}", id, partition, document.offset);
40       if (document.offset >= 0)
41       {
42         // Only seek, if a stored offset was found
43         // Otherwise: Use initial offset, generated by Kafka
44         consumer.seek(tp, document.offset);
45       }
46       handler.addPartition(partition, document.state);
47     });
48   }
49
50   @Override
51   public void onPartitionsRevoked(Collection<TopicPartition> partitions)
52   {
53     partitions.forEach(tp ->
54     {
55       Integer partition = tp.partition();
56       Long offset = consumer.position(tp);
57       log.info(
58           "{} - removing partition: {}, offset of next message {})",
59           id,
60           partition,
61           offset);
62       if (commitsEnabled)
63       {
64         repository.save(new StateDocument(partition, handler.removePartition(partition), offset));
65       }
66       else
67       {
68         log.info("Offset commits are disabled! Last commit: {}", lastCommit);
69       }
70     });
71   }
72
73
74   @Override
75   public void beforeNextPoll()
76   {
77     if (!commitsEnabled)
78     {
79       log.info("Offset commits are disabled! Last commit: {}", lastCommit);
80       return;
81     }
82
83     if (lastCommit.plus(commitInterval).isBefore(clock.instant()))
84     {
85       log.debug("Storing data and offsets, last commit: {}", lastCommit);
86       handler.getState().forEach((partiton, sumBusinessLogic) -> repository.save(
87           new StateDocument(
88               partiton,
89               sumBusinessLogic.getState(),
90               consumer.position(new TopicPartition(topic, partiton)))));
91       lastCommit = clock.instant();
92     }
93   }
94
95   @Override
96   public void enableCommits()
97   {
98     commitsEnabled = true;
99   }
100
101   @Override
102   public void disableCommits()
103   {
104     commitsEnabled = false;
105   }
106 }