From 908b472727025a6c23abccaf4b8f1e39859baf15 Mon Sep 17 00:00:00 2001 From: Kai Moritz Date: Fri, 22 Oct 2021 14:13:42 +0200 Subject: [PATCH] WIP --- .../wordcount/recorder/SplitterStreamProcessor.java | 9 +++++++-- 1 file changed, 7 insertions(+), 2 deletions(-) diff --git a/src/main/java/de/juplo/kafka/wordcount/recorder/SplitterStreamProcessor.java b/src/main/java/de/juplo/kafka/wordcount/recorder/SplitterStreamProcessor.java index fca26d5..dd00aca 100644 --- a/src/main/java/de/juplo/kafka/wordcount/recorder/SplitterStreamProcessor.java +++ b/src/main/java/de/juplo/kafka/wordcount/recorder/SplitterStreamProcessor.java @@ -63,7 +63,7 @@ public class SplitterStreamProcessor implements ApplicationRunner try { - log.info("C - Subscribing to topic test"); + log.info("Subscribing to topic {}", inputTopic); consumer.subscribe( Arrays.asList(inputTopic), new TransactionalConsumerRebalanceListener()); @@ -131,7 +131,9 @@ public class SplitterStreamProcessor implements ApplicationRunner { log.info("Closing consumer"); consumer.close(); - log.info("C - DONE!"); + log.info("Closing producer"); + producer.close(); + log.info("Exiting!"); } } @@ -145,6 +147,9 @@ public class SplitterStreamProcessor implements ApplicationRunner private void commitTransaction() { log.info("Committing transaction"); + producer.sendOffsetsToTransaction( + consumer.po, + consumer.groupMetadata()); producer.commitTransaction(); } -- 2.20.1