From 1cbfef1ef5f134b5f1eaa0de1bfda3bbd4b59441 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 3acaebe..385df67 100644 --- a/src/main/java/de/juplo/kafka/ExampleConsumer.java +++ b/src/main/java/de/juplo/kafka/ExampleConsumer.java @@ -81,7 +81,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 -> @@ -103,8 +105,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