projects
/
demos
/
kafka
/
training
/ blobdiff
commit
grep
author
committer
pickaxe
?
search:
re
summary
|
shortlog
|
log
|
commit
|
commitdiff
|
tree
raw
|
inline
| side by side
Log-Meldungen in ApplicationRebalanceListener angeglichen
[demos/kafka/training]
/
src
/
main
/
java
/
de
/
juplo
/
kafka
/
ApplicationRebalanceListener.java
diff --git
a/src/main/java/de/juplo/kafka/ApplicationRebalanceListener.java
b/src/main/java/de/juplo/kafka/ApplicationRebalanceListener.java
index
109b205
..
8e8464f
100644
(file)
--- a/
src/main/java/de/juplo/kafka/ApplicationRebalanceListener.java
+++ b/
src/main/java/de/juplo/kafka/ApplicationRebalanceListener.java
@@
-13,7
+13,7
@@
import java.util.*;
@RequiredArgsConstructor
@Slf4j
@RequiredArgsConstructor
@Slf4j
-public class ApplicationRebalanceListener implements
PollIntervalAwareConsumer
RebalanceListener
+public class ApplicationRebalanceListener implements RebalanceListener
{
private final ApplicationRecordHandler recordHandler;
private final AdderResults adderResults;
{
private final ApplicationRecordHandler recordHandler;
private final AdderResults adderResults;
@@
-22,7
+22,7
@@
public class ApplicationRebalanceListener implements PollIntervalAwareConsumerRe
private final String topic;
private final Clock clock;
private final Duration commitInterval;
private final String topic;
private final Clock clock;
private final Duration commitInterval;
- private final Consumer
<String, String>
consumer;
+ private final Consumer consumer;
private final Set<Integer> partitions = new HashSet<>();
private final Set<Integer> partitions = new HashSet<>();
@@
-97,7
+97,7
@@
public class ApplicationRebalanceListener implements PollIntervalAwareConsumerRe
}
else
{
}
else
{
- log.info("
Offset commits are disabled! Last commit: {}"
, lastCommit);
+ log.info("
{} - Offset commits are disabled! Last commit: {}", id
, lastCommit);
}
});
}
}
});
}
@@
-108,13
+108,13
@@
public class ApplicationRebalanceListener implements PollIntervalAwareConsumerRe
{
if (!commitsEnabled)
{
{
if (!commitsEnabled)
{
- log.info("
Offset commits are disabled! Last commit: {}"
, lastCommit);
+ log.info("
{} - Offset commits are disabled! Last commit: {}", id
, lastCommit);
return;
}
if (lastCommit.plus(commitInterval).isBefore(clock.instant()))
{
return;
}
if (lastCommit.plus(commitInterval).isBefore(clock.instant()))
{
- log.debug("
Storing data and offsets, last commit: {}"
, lastCommit);
+ log.debug("
{} - Storing data and offsets, last commit: {}", id
, lastCommit);
partitions.forEach(partition -> stateRepository.save(
new StateDocument(
partition,
partitions.forEach(partition -> stateRepository.save(
new StateDocument(
partition,