From 664bec89418bb2349fe93349edd0f46b95a82726 Mon Sep 17 00:00:00 2001 From: Kai Moritz Date: Fri, 18 Nov 2022 16:52:21 +0100 Subject: [PATCH] =?utf8?q?SimpleConsumer=20signalisiert=20Exit-Status=20im?= =?utf8?q?=20R=C3=BCckgabewert?= MIME-Version: 1.0 Content-Type: text/plain; charset=utf8 Content-Transfer-Encoding: 8bit --- src/main/java/de/juplo/kafka/SimpleConsumer.java | 7 +++++-- 1 file changed, 5 insertions(+), 2 deletions(-) diff --git a/src/main/java/de/juplo/kafka/SimpleConsumer.java b/src/main/java/de/juplo/kafka/SimpleConsumer.java index 1cf9b22..aadc11f 100644 --- a/src/main/java/de/juplo/kafka/SimpleConsumer.java +++ b/src/main/java/de/juplo/kafka/SimpleConsumer.java @@ -9,11 +9,12 @@ import org.apache.kafka.common.errors.WakeupException; import java.time.Duration; import java.util.Arrays; +import java.util.concurrent.Callable; @Slf4j @RequiredArgsConstructor -public class SimpleConsumer implements Runnable +public class SimpleConsumer implements Callable { private final String id; private final String topic; @@ -23,7 +24,7 @@ public class SimpleConsumer implements Runnable @Override - public void run() + public Integer call() { try { @@ -50,11 +51,13 @@ public class SimpleConsumer implements Runnable catch(WakeupException e) { log.info("{} - Consumer was signaled to finish its work", id); + return 0; } catch(Exception e) { log.error("{} - Unexpected error: {}, unsubscribing!", id, e.toString()); consumer.unsubscribe(); + return 1; } finally { -- 2.20.1