From d6db20dd022821f012402f012be8d51e68f9757e Mon Sep 17 00:00:00 2001 From: panos-kakos <102670093+panos-kakos@users.noreply.github.com> Date: Mon, 26 Feb 2024 21:34:40 +0200 Subject: [PATCH] [JAVA-29500] Moved article "Implementing Retry in Kafka Consumer" to spring-kafka (#15904) --- spring-kafka-2/README.md | 1 - spring-kafka-3_/README.md | 2 - spring-kafka-3_/pom.xml | 34 ----------- .../com/baeldung/spring/kafka/SomeData.java | 37 ------------ .../spring/kafka/ListenerConfiguration.java | 42 -------------- .../spring/kafka/ProducerConfiguration.java | 40 ------------- .../spring/kafka/TrustedPackagesLiveTest.java | 57 ------------------- spring-kafka/README.md | 2 +- .../baeldung}/kafka/retryable/Farewell.java | 2 +- .../baeldung}/kafka/retryable/Greeting.java | 2 +- .../kafka/retryable/KafkaConsumerConfig.java | 3 +- .../kafka/retryable/KafkaProducerConfig.java | 2 +- .../kafka/retryable/KafkaTopicConfig.java | 2 +- .../retryable/MultiTypeKafkaListener.java | 2 +- .../RetryableApplicationKafkaApp.java | 2 +- .../resources/application-retry.properties | 0 .../KafkaRetryableIntegrationTest.java | 18 ++---- 17 files changed, 14 insertions(+), 234 deletions(-) delete mode 100644 spring-kafka-3_/README.md delete mode 100644 spring-kafka-3_/pom.xml delete mode 100644 spring-kafka-3_/src/main/java/com/baeldung/spring/kafka/SomeData.java delete mode 100644 spring-kafka-3_/src/test/java/com/baeldung/spring/kafka/ListenerConfiguration.java delete mode 100644 spring-kafka-3_/src/test/java/com/baeldung/spring/kafka/ProducerConfiguration.java delete mode 100644 spring-kafka-3_/src/test/java/com/baeldung/spring/kafka/TrustedPackagesLiveTest.java rename {spring-kafka-2/src/main/java/com/baeldung/spring => spring-kafka/src/main/java/com/baeldung}/kafka/retryable/Farewell.java (94%) rename {spring-kafka-2/src/main/java/com/baeldung/spring => spring-kafka/src/main/java/com/baeldung}/kafka/retryable/Greeting.java (92%) rename {spring-kafka-2/src/main/java/com/baeldung/spring => spring-kafka/src/main/java/com/baeldung}/kafka/retryable/KafkaConsumerConfig.java (98%) rename {spring-kafka-2/src/main/java/com/baeldung/spring => spring-kafka/src/main/java/com/baeldung}/kafka/retryable/KafkaProducerConfig.java (98%) rename {spring-kafka-2/src/main/java/com/baeldung/spring => spring-kafka/src/main/java/com/baeldung}/kafka/retryable/KafkaTopicConfig.java (97%) rename {spring-kafka-2/src/main/java/com/baeldung/spring => spring-kafka/src/main/java/com/baeldung}/kafka/retryable/MultiTypeKafkaListener.java (95%) rename {spring-kafka-2/src/main/java/com/baeldung/spring => spring-kafka/src/main/java/com/baeldung}/kafka/retryable/RetryableApplicationKafkaApp.java (91%) rename {spring-kafka-2 => spring-kafka}/src/main/resources/application-retry.properties (100%) rename {spring-kafka-2/src/test/java/com/baeldung/spring => spring-kafka/src/test/java/com/baeldung}/kafka/retryable/KafkaRetryableIntegrationTest.java (88%) diff --git a/spring-kafka-2/README.md b/spring-kafka-2/README.md index f9e07d4893..4dff7ef5db 100644 --- a/spring-kafka-2/README.md +++ b/spring-kafka-2/README.md @@ -4,7 +4,6 @@ This module contains articles about Spring with Kafka ### Relevant articles -- [Implementing Retry in Kafka Consumer](https://www.baeldung.com/spring-retry-kafka-consumer) - [Spring Kafka: Configure Multiple Listeners on Same Topic](https://www.baeldung.com/spring-kafka-multiple-listeners-same-topic) - [Understanding Kafka Topics and Partitions](https://www.baeldung.com/kafka-topics-partitions) - [How to Subscribe a Kafka Consumer to Multiple Topics](https://www.baeldung.com/kafka-subscribe-consumer-multiple-topics) diff --git a/spring-kafka-3_/README.md b/spring-kafka-3_/README.md deleted file mode 100644 index f9c0036ce3..0000000000 --- a/spring-kafka-3_/README.md +++ /dev/null @@ -1,2 +0,0 @@ -## Relevant Articles -- [Spring Kafka Trusted Packages Feature](https://www.baeldung.com/spring-kafka-trusted-packages-feature) diff --git a/spring-kafka-3_/pom.xml b/spring-kafka-3_/pom.xml deleted file mode 100644 index 058b981bc1..0000000000 --- a/spring-kafka-3_/pom.xml +++ /dev/null @@ -1,34 +0,0 @@ - - 4.0.0 - spring-kafka-3_ - jar - spring-kafka-3_ - - - com.baeldung - parent-boot-2 - 0.0.1-SNAPSHOT - ../parent-boot-2 - - - - - org.springframework.kafka - spring-kafka - - - com.fasterxml.jackson.core - jackson-databind - - - org.springframework.kafka - spring-kafka-test - test - - - - - 3.0.12 - - diff --git a/spring-kafka-3_/src/main/java/com/baeldung/spring/kafka/SomeData.java b/spring-kafka-3_/src/main/java/com/baeldung/spring/kafka/SomeData.java deleted file mode 100644 index eb2e57c33d..0000000000 --- a/spring-kafka-3_/src/main/java/com/baeldung/spring/kafka/SomeData.java +++ /dev/null @@ -1,37 +0,0 @@ -package com.baeldung.spring.kafka; - -import java.time.Instant; - -public class SomeData { - - private String id; - private String type; - private String status; - private Instant timestamp; - - public SomeData() { - } - - public SomeData(String id, String type, String status, Instant timestamp) { - this.id = id; - this.type = type; - this.status = status; - this.timestamp = timestamp; - } - - public String getId() { - return id; - } - - public String getType() { - return type; - } - - public String getStatus() { - return status; - } - - public Instant getTimestamp() { - return timestamp; - } -} diff --git a/spring-kafka-3_/src/test/java/com/baeldung/spring/kafka/ListenerConfiguration.java b/spring-kafka-3_/src/test/java/com/baeldung/spring/kafka/ListenerConfiguration.java deleted file mode 100644 index e35b1ee415..0000000000 --- a/spring-kafka-3_/src/test/java/com/baeldung/spring/kafka/ListenerConfiguration.java +++ /dev/null @@ -1,42 +0,0 @@ -package com.baeldung.spring.kafka; - -import java.util.Map; - -import org.apache.kafka.clients.consumer.ConsumerConfig; -import org.apache.kafka.common.serialization.StringDeserializer; -import org.springframework.context.annotation.Bean; -import org.springframework.context.annotation.Configuration; -import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory; -import org.springframework.kafka.core.ConsumerFactory; -import org.springframework.kafka.core.DefaultKafkaConsumerFactory; -import org.springframework.kafka.support.serializer.JsonDeserializer; - -@Configuration -public class ListenerConfiguration { - - @Bean("messageListenerContainer") - public ConcurrentKafkaListenerContainerFactory messageListenerContainer() { - ConcurrentKafkaListenerContainerFactory container = new ConcurrentKafkaListenerContainerFactory<>(); - container.setConsumerFactory(someDataConsumerFactory()); - return container; - } - - @Bean - public ConsumerFactory someDataConsumerFactory() { - JsonDeserializer payloadJsonDeserializer = new JsonDeserializer<>(); - payloadJsonDeserializer.trustedPackages("com.baeldung.spring.kafka"); - return new DefaultKafkaConsumerFactory<>( - consumerConfigs(), - new StringDeserializer(), - payloadJsonDeserializer - ); - } - - @Bean - public Map consumerConfigs() { - return Map.of( - ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "PLAINTEXT://localhost:9092", - ConsumerConfig.GROUP_ID_CONFIG, "some-group-id" - ); - } -} \ No newline at end of file diff --git a/spring-kafka-3_/src/test/java/com/baeldung/spring/kafka/ProducerConfiguration.java b/spring-kafka-3_/src/test/java/com/baeldung/spring/kafka/ProducerConfiguration.java deleted file mode 100644 index 7373424841..0000000000 --- a/spring-kafka-3_/src/test/java/com/baeldung/spring/kafka/ProducerConfiguration.java +++ /dev/null @@ -1,40 +0,0 @@ -package com.baeldung.spring.kafka; - -import org.apache.kafka.clients.producer.ProducerConfig; -import org.springframework.context.annotation.Bean; -import org.springframework.context.annotation.Configuration; -import org.springframework.kafka.core.DefaultKafkaProducerFactory; -import org.springframework.kafka.core.KafkaTemplate; -import org.springframework.kafka.core.ProducerFactory; -import org.springframework.kafka.support.serializer.JsonSerializer; -import org.springframework.kafka.support.serializer.StringOrBytesSerializer; - -import java.time.Instant; -import java.util.Map; - -@Configuration -public class ProducerConfiguration { - - @Bean - public KafkaTemplate kafkaTemplate() { - return new KafkaTemplate<>(producerFactory()); - } - - @Bean - public ProducerFactory producerFactory() { - JsonSerializer jsonSerializer = new JsonSerializer<>(); - jsonSerializer.setAddTypeInfo(true); - return new DefaultKafkaProducerFactory<>( - producerFactoryConfig(), - new StringOrBytesSerializer(), - jsonSerializer - ); - } - - @Bean - public Map producerFactoryConfig() { - return Map.of( - ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "PLAINTEXT://localhost:9092" - ); - } -} \ No newline at end of file diff --git a/spring-kafka-3_/src/test/java/com/baeldung/spring/kafka/TrustedPackagesLiveTest.java b/spring-kafka-3_/src/test/java/com/baeldung/spring/kafka/TrustedPackagesLiveTest.java deleted file mode 100644 index fa4b79cd65..0000000000 --- a/spring-kafka-3_/src/test/java/com/baeldung/spring/kafka/TrustedPackagesLiveTest.java +++ /dev/null @@ -1,57 +0,0 @@ -package com.baeldung.spring.kafka; - -import java.time.Instant; -import java.util.concurrent.CountDownLatch; -import java.util.concurrent.TimeUnit; - -import org.apache.kafka.clients.producer.ProducerRecord; -import org.junit.jupiter.api.Assertions; -import org.junit.jupiter.api.Test; -import org.mockito.Mockito; -import org.springframework.beans.factory.annotation.Autowired; -import org.springframework.boot.test.context.SpringBootTest; -import org.springframework.boot.test.mock.mockito.SpyBean; -import org.springframework.kafka.annotation.KafkaListener; -import org.springframework.kafka.core.KafkaTemplate; -import org.springframework.stereotype.Component; - -/** - * This test requires a running instance of kafka to be present - */ -@SpringBootTest -public class TrustedPackagesLiveTest { - - @Autowired - private KafkaTemplate kafkaTemplate; - - @SpyBean - TestConsumer testConsumer; - - @Test - void givenMessageInTheTopic_whenTypeInfoPackageIsTrusted_thenMessageIsSuccessfullyConsumed() throws InterruptedException { - CountDownLatch latch = new CountDownLatch(1); - - Mockito.doAnswer(invocationOnMock -> { - try { - latch.countDown(); - return invocationOnMock.callRealMethod(); - } catch (Exception e) { - return null; - } - }).when(testConsumer).onMessage(Mockito.any()); - - SomeData someData = new SomeData("1", "active", "sent", Instant.now()); - kafkaTemplate.send(new ProducerRecord<>("sourceTopic", null, someData)); - - Assertions.assertTrue(latch.await(20L, TimeUnit.SECONDS)); - } - - @Component - static class TestConsumer { - - @KafkaListener(topics = "sourceTopic", containerFactory = "messageListenerContainer") - public void onMessage(SomeData someData) { - - } - } -} diff --git a/spring-kafka/README.md b/spring-kafka/README.md index 97f459f534..6e495b1210 100644 --- a/spring-kafka/README.md +++ b/spring-kafka/README.md @@ -12,7 +12,7 @@ This module contains articles about Spring with Kafka - [Kafka Streams With Spring Boot](https://www.baeldung.com/spring-boot-kafka-streams) - [Get the Number of Messages in an Apache Kafka Topic](https://www.baeldung.com/java-kafka-count-topic-messages) - [Sending Data to a Specific Partition in Kafka](https://www.baeldung.com/kafka-send-data-partition) - +- [Implementing Retry in Kafka Consumer](https://www.baeldung.com/spring-retry-kafka-consumer) ### Intro This is a simple Spring Boot app to demonstrate sending and receiving of messages in Kafka using spring-kafka. diff --git a/spring-kafka-2/src/main/java/com/baeldung/spring/kafka/retryable/Farewell.java b/spring-kafka/src/main/java/com/baeldung/kafka/retryable/Farewell.java similarity index 94% rename from spring-kafka-2/src/main/java/com/baeldung/spring/kafka/retryable/Farewell.java rename to spring-kafka/src/main/java/com/baeldung/kafka/retryable/Farewell.java index 519d847aab..d3171540e8 100644 --- a/spring-kafka-2/src/main/java/com/baeldung/spring/kafka/retryable/Farewell.java +++ b/spring-kafka/src/main/java/com/baeldung/kafka/retryable/Farewell.java @@ -1,4 +1,4 @@ -package com.baeldung.spring.kafka.retryable; +package com.baeldung.kafka.retryable; public class Farewell { diff --git a/spring-kafka-2/src/main/java/com/baeldung/spring/kafka/retryable/Greeting.java b/spring-kafka/src/main/java/com/baeldung/kafka/retryable/Greeting.java similarity index 92% rename from spring-kafka-2/src/main/java/com/baeldung/spring/kafka/retryable/Greeting.java rename to spring-kafka/src/main/java/com/baeldung/kafka/retryable/Greeting.java index 79abeda34e..372c0db20c 100644 --- a/spring-kafka-2/src/main/java/com/baeldung/spring/kafka/retryable/Greeting.java +++ b/spring-kafka/src/main/java/com/baeldung/kafka/retryable/Greeting.java @@ -1,4 +1,4 @@ -package com.baeldung.spring.kafka.retryable; +package com.baeldung.kafka.retryable; public class Greeting { diff --git a/spring-kafka-2/src/main/java/com/baeldung/spring/kafka/retryable/KafkaConsumerConfig.java b/spring-kafka/src/main/java/com/baeldung/kafka/retryable/KafkaConsumerConfig.java similarity index 98% rename from spring-kafka-2/src/main/java/com/baeldung/spring/kafka/retryable/KafkaConsumerConfig.java rename to spring-kafka/src/main/java/com/baeldung/kafka/retryable/KafkaConsumerConfig.java index cf4b29137e..d5d8549785 100644 --- a/spring-kafka-2/src/main/java/com/baeldung/spring/kafka/retryable/KafkaConsumerConfig.java +++ b/spring-kafka/src/main/java/com/baeldung/kafka/retryable/KafkaConsumerConfig.java @@ -1,4 +1,4 @@ -package com.baeldung.spring.kafka.retryable; +package com.baeldung.kafka.retryable; import java.util.HashMap; import java.util.Map; @@ -129,7 +129,6 @@ public class KafkaConsumerConfig { public ConcurrentKafkaListenerContainerFactory greetingKafkaListenerContainerFactory() { ConcurrentKafkaListenerContainerFactory factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(multiTypeConsumerFactory()); - factory.setMessageConverter(multiTypeConverter()); factory.setCommonErrorHandler(errorHandler()); factory.getContainerProperties() .setAckMode(ContainerProperties.AckMode.RECORD); diff --git a/spring-kafka-2/src/main/java/com/baeldung/spring/kafka/retryable/KafkaProducerConfig.java b/spring-kafka/src/main/java/com/baeldung/kafka/retryable/KafkaProducerConfig.java similarity index 98% rename from spring-kafka-2/src/main/java/com/baeldung/spring/kafka/retryable/KafkaProducerConfig.java rename to spring-kafka/src/main/java/com/baeldung/kafka/retryable/KafkaProducerConfig.java index 90dcb2ccf9..7eeccd9b57 100644 --- a/spring-kafka-2/src/main/java/com/baeldung/spring/kafka/retryable/KafkaProducerConfig.java +++ b/spring-kafka/src/main/java/com/baeldung/kafka/retryable/KafkaProducerConfig.java @@ -1,4 +1,4 @@ -package com.baeldung.spring.kafka.retryable; +package com.baeldung.kafka.retryable; import java.util.HashMap; import java.util.Map; diff --git a/spring-kafka-2/src/main/java/com/baeldung/spring/kafka/retryable/KafkaTopicConfig.java b/spring-kafka/src/main/java/com/baeldung/kafka/retryable/KafkaTopicConfig.java similarity index 97% rename from spring-kafka-2/src/main/java/com/baeldung/spring/kafka/retryable/KafkaTopicConfig.java rename to spring-kafka/src/main/java/com/baeldung/kafka/retryable/KafkaTopicConfig.java index 3b4c2d9928..af0ce6b8cf 100644 --- a/spring-kafka-2/src/main/java/com/baeldung/spring/kafka/retryable/KafkaTopicConfig.java +++ b/spring-kafka/src/main/java/com/baeldung/kafka/retryable/KafkaTopicConfig.java @@ -1,4 +1,4 @@ -package com.baeldung.spring.kafka.retryable; +package com.baeldung.kafka.retryable; import java.util.HashMap; import java.util.Map; diff --git a/spring-kafka-2/src/main/java/com/baeldung/spring/kafka/retryable/MultiTypeKafkaListener.java b/spring-kafka/src/main/java/com/baeldung/kafka/retryable/MultiTypeKafkaListener.java similarity index 95% rename from spring-kafka-2/src/main/java/com/baeldung/spring/kafka/retryable/MultiTypeKafkaListener.java rename to spring-kafka/src/main/java/com/baeldung/kafka/retryable/MultiTypeKafkaListener.java index 441e564176..b73d150ce3 100644 --- a/spring-kafka-2/src/main/java/com/baeldung/spring/kafka/retryable/MultiTypeKafkaListener.java +++ b/spring-kafka/src/main/java/com/baeldung/kafka/retryable/MultiTypeKafkaListener.java @@ -1,4 +1,4 @@ -package com.baeldung.spring.kafka.retryable; +package com.baeldung.kafka.retryable; import org.springframework.kafka.annotation.KafkaHandler; import org.springframework.kafka.annotation.KafkaListener; diff --git a/spring-kafka-2/src/main/java/com/baeldung/spring/kafka/retryable/RetryableApplicationKafkaApp.java b/spring-kafka/src/main/java/com/baeldung/kafka/retryable/RetryableApplicationKafkaApp.java similarity index 91% rename from spring-kafka-2/src/main/java/com/baeldung/spring/kafka/retryable/RetryableApplicationKafkaApp.java rename to spring-kafka/src/main/java/com/baeldung/kafka/retryable/RetryableApplicationKafkaApp.java index 458ebac124..fd1fb490c6 100644 --- a/spring-kafka-2/src/main/java/com/baeldung/spring/kafka/retryable/RetryableApplicationKafkaApp.java +++ b/spring-kafka/src/main/java/com/baeldung/kafka/retryable/RetryableApplicationKafkaApp.java @@ -1,4 +1,4 @@ -package com.baeldung.spring.kafka.retryable; +package com.baeldung.kafka.retryable; import org.springframework.boot.SpringApplication; import org.springframework.boot.autoconfigure.SpringBootApplication; diff --git a/spring-kafka-2/src/main/resources/application-retry.properties b/spring-kafka/src/main/resources/application-retry.properties similarity index 100% rename from spring-kafka-2/src/main/resources/application-retry.properties rename to spring-kafka/src/main/resources/application-retry.properties diff --git a/spring-kafka-2/src/test/java/com/baeldung/spring/kafka/retryable/KafkaRetryableIntegrationTest.java b/spring-kafka/src/test/java/com/baeldung/kafka/retryable/KafkaRetryableIntegrationTest.java similarity index 88% rename from spring-kafka-2/src/test/java/com/baeldung/spring/kafka/retryable/KafkaRetryableIntegrationTest.java rename to spring-kafka/src/test/java/com/baeldung/kafka/retryable/KafkaRetryableIntegrationTest.java index 876ad0fdc7..170be53f0a 100644 --- a/spring-kafka-2/src/test/java/com/baeldung/spring/kafka/retryable/KafkaRetryableIntegrationTest.java +++ b/spring-kafka/src/test/java/com/baeldung/kafka/retryable/KafkaRetryableIntegrationTest.java @@ -1,12 +1,11 @@ -package com.baeldung.spring.kafka.retryable; +package com.baeldung.kafka.retryable; import static org.assertj.core.api.AssertionsForClassTypes.assertThat; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; -import org.junit.Before; -import org.junit.ClassRule; +import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.test.context.SpringBootTest; @@ -14,20 +13,15 @@ import org.springframework.kafka.config.KafkaListenerEndpointRegistry; import org.springframework.kafka.core.KafkaTemplate; import org.springframework.kafka.listener.AcknowledgingConsumerAwareMessageListener; import org.springframework.kafka.listener.ConcurrentMessageListenerContainer; -import org.springframework.kafka.test.EmbeddedKafkaBroker; import org.springframework.kafka.test.context.EmbeddedKafka; - -import com.baeldung.spring.kafka.retryable.Greeting; -import com.baeldung.spring.kafka.retryable.RetryableApplicationKafkaApp; -import com.fasterxml.jackson.databind.ObjectMapper; import org.springframework.test.context.ActiveProfiles; +import com.fasterxml.jackson.databind.ObjectMapper; + @SpringBootTest(classes = RetryableApplicationKafkaApp.class) @EmbeddedKafka(partitions = 1, controlledShutdown = true, brokerProperties = { "listeners=PLAINTEXT://localhost:9093", "port=9093" }) @ActiveProfiles("retry") public class KafkaRetryableIntegrationTest { - @ClassRule - public static EmbeddedKafkaBroker embeddedKafka = new EmbeddedKafkaBroker(1, true, "multitype"); @Autowired private KafkaListenerEndpointRegistry registry; @@ -41,9 +35,9 @@ public class KafkaRetryableIntegrationTest { private static final String TOPIC = "topic"; - @Before + @BeforeEach public void setup() { - System.setProperty("spring.kafka.bootstrap-servers", embeddedKafka.getBrokersAsString()); + System.setProperty("spring.kafka.bootstrap-servers", "localhost:9093"); } @Test