diff --git a/apache-kafka-2/src/main/java/com/baeldung/kafka/consumer/CustomKafkaListener.java b/apache-kafka-2/src/main/java/com/baeldung/kafka/consumer/CustomKafkaListener.java index 58987d473b..b432c2ba55 100644 --- a/apache-kafka-2/src/main/java/com/baeldung/kafka/consumer/CustomKafkaListener.java +++ b/apache-kafka-2/src/main/java/com/baeldung/kafka/consumer/CustomKafkaListener.java @@ -12,7 +12,7 @@ import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.consumer.KafkaConsumer; import org.apache.kafka.common.serialization.StringDeserializer; -public class CustomKafkaListener implements Runnable, Closeable { +public class CustomKafkaListener implements Runnable, AutoCloseable { private static final Logger log = Logger.getLogger(CustomKafkaListener.class.getName()); @@ -43,7 +43,7 @@ public class CustomKafkaListener implements Runnable, Closeable { return new KafkaConsumer<>(props); } - public CustomKafkaListener doForEach(Consumer newConsumer) { + public CustomKafkaListener onEach(Consumer newConsumer) { recordConsumer = recordConsumer.andThen(newConsumer); return this; } @@ -52,10 +52,9 @@ public class CustomKafkaListener implements Runnable, Closeable { public void run() { running.set(true); consumer.subscribe(Arrays.asList(topic)); - while (running.get()) { consumer.poll(Duration.ofMillis(100)) - .forEach(record -> recordConsumer.accept(record.value())); + .forEach(record -> recordConsumer.accept(record.value())); } } diff --git a/apache-kafka-2/src/test/java/com/baeldung/kafka/consumer/KafkaListenerWithoutSpringLiveTest.java b/apache-kafka-2/src/test/java/com/baeldung/kafka/consumer/CustomKafkaListenerLiveTest.java similarity index 84% rename from apache-kafka-2/src/test/java/com/baeldung/kafka/consumer/KafkaListenerWithoutSpringLiveTest.java rename to apache-kafka-2/src/test/java/com/baeldung/kafka/consumer/CustomKafkaListenerLiveTest.java index 213344d0a8..dbe7483dcd 100644 --- a/apache-kafka-2/src/test/java/com/baeldung/kafka/consumer/KafkaListenerWithoutSpringLiveTest.java +++ b/apache-kafka-2/src/test/java/com/baeldung/kafka/consumer/CustomKafkaListenerLiveTest.java @@ -7,6 +7,7 @@ import static org.assertj.core.api.Assertions.assertThat; import static org.testcontainers.shaded.org.awaitility.Awaitility.await; import java.util.ArrayList; +import java.util.Arrays; import java.util.List; import java.util.Properties; import java.util.concurrent.CompletableFuture; @@ -23,7 +24,7 @@ import org.testcontainers.shaded.org.awaitility.Awaitility; import org.testcontainers.utility.DockerImageName; @Testcontainers -class KafkaListenerWithoutSpringLiveTest { +class CustomKafkaListenerLiveTest { @Container private static final KafkaContainer KAFKA_CONTAINER = new KafkaContainer(DockerImageName.parse("confluentinc/cp-kafka:latest")); @@ -34,30 +35,28 @@ class KafkaListenerWithoutSpringLiveTest { } @Test - void test() { + void givenANewCustomKafkaListener_thenConsumesAllMessages() { // given String topic = "baeldung.articles.published"; String bootstrapServers = KAFKA_CONTAINER.getBootstrapServers(); List consumedMessages = new ArrayList<>(); // when - try (CustomKafkaListener listener = new CustomKafkaListener(topic, bootstrapServers)) { - CompletableFuture.runAsync(() -> - listener.doForEach(consumedMessages::add).run() - ); + try (CustomKafkaListener listener = new CustomKafkaListener(topic, bootstrapServers).onEach(consumedMessages::add)) { + CompletableFuture.runAsync(listener); } // and - publishArticles(topic, asList( + publishArticles(topic, "Introduction to Kafka", "Kotlin for Java Developers", "Reactive Spring Boot", "Deploying Spring Boot Applications", "Spring Security" - )); + ); // then - await().untilAsserted(() -> assertThat(consumedMessages) - .containsExactlyInAnyOrder( + await().untilAsserted(() -> + assertThat(consumedMessages).containsExactlyInAnyOrder( "Introduction to Kafka", "Kotlin for Java Developers", "Reactive Spring Boot", @@ -66,9 +65,9 @@ class KafkaListenerWithoutSpringLiveTest { )); } - private void publishArticles(String topic, List articles) { + private void publishArticles(String topic, String... articles) { try (KafkaProducer producer = testKafkaProducer()) { - articles.stream() + Arrays.stream(articles) .map(article -> new ProducerRecord(topic, article)) .forEach(producer::send); }