From 9d46bce0c4dc584a550ecb44f6399ce0feb3299d Mon Sep 17 00:00:00 2001 From: Kai Moritz Date: Mon, 28 Oct 2024 11:13:05 +0100 Subject: [PATCH] Fehler im Logging der aktiven Phase korrigiert und Meldungen verbessert --- src/main/java/de/juplo/kafka/ExampleConsumer.java | 8 +++++--- 1 file changed, 5 insertions(+), 3 deletions(-) diff --git a/src/main/java/de/juplo/kafka/ExampleConsumer.java b/src/main/java/de/juplo/kafka/ExampleConsumer.java index 2ad0a93..d845138 100644 --- a/src/main/java/de/juplo/kafka/ExampleConsumer.java +++ b/src/main/java/de/juplo/kafka/ExampleConsumer.java @@ -77,7 +77,9 @@ public class ExampleConsumer implements Runnable, ConsumerRebalanceListener ConsumerRecords records = consumer.poll(Duration.ofSeconds(1)); - log.info("{} - Received {} messages", id, records.count()); + int phase = phaser.getPhase(); + + log.info("{} - Received {} messages in phase {}", id, records.count(), phase); records .partitions() .forEach(partition -> @@ -99,8 +101,8 @@ public class ExampleConsumer implements Runnable, ConsumerRebalanceListener done[partition.partition()] = true; }); - int phase = phaser.arriveAndAwaitAdvance(); - log.info("{} - Phase {} is done!", id, phase); + int arrivedPhase = phaser.arriveAndAwaitAdvance(); + log.info("{} - Phase {} is done! Next phase: {}", id, phase, arrivedPhase); } } catch(WakeupException e) -- 2.20.1