From d8e035da08782a110ca8e75e7bdc96c90d4ae7fa Mon Sep 17 00:00:00 2001 From: s9m33r Date: Wed, 31 Jan 2024 00:40:41 +0530 Subject: [PATCH] BAEL-7351 moving to a new module --- apache-kafka-3/.gitignore | 33 +++++++ apache-kafka-3/pom.xml | 50 +++++++++++ .../apachekafka3}/groupId/KafkaConfig.java | 2 +- .../groupId/KafkaErrorHandler.java | 9 +- .../baeldung/apachekafka3}/groupId/Main.java | 28 +++--- .../groupId/MyKafkaConsumer.java | 82 ++++++++--------- .../groupId/MyKafkaProducer.java | 48 +++++----- .../src/main/resources/application.properties | 1 + .../apachekafka3}/groupId/MainLiveTest.java | 88 +++++++++---------- pom.xml | 3 +- 10 files changed, 211 insertions(+), 133 deletions(-) create mode 100644 apache-kafka-3/.gitignore create mode 100644 apache-kafka-3/pom.xml rename {spring-kafka-3/src/main/java/com/baeldung/spring/kafka => apache-kafka-3/src/main/java/com/baeldung/apachekafka3}/groupId/KafkaConfig.java (97%) rename {spring-kafka-3/src/main/java/com/baeldung/spring/kafka => apache-kafka-3/src/main/java/com/baeldung/apachekafka3}/groupId/KafkaErrorHandler.java (76%) rename {spring-kafka-3/src/main/java/com/baeldung/spring/kafka => apache-kafka-3/src/main/java/com/baeldung/apachekafka3}/groupId/Main.java (73%) rename {spring-kafka-3/src/main/java/com/baeldung/spring/kafka => apache-kafka-3/src/main/java/com/baeldung/apachekafka3}/groupId/MyKafkaConsumer.java (82%) rename {spring-kafka-3/src/main/java/com/baeldung/spring/kafka => apache-kafka-3/src/main/java/com/baeldung/apachekafka3}/groupId/MyKafkaProducer.java (91%) create mode 100644 apache-kafka-3/src/main/resources/application.properties rename {spring-kafka-3/src/test/java/com/baeldung/spring/kafka => apache-kafka-3/src/test/java/com/baeldung/apachekafka3}/groupId/MainLiveTest.java (90%) diff --git a/apache-kafka-3/.gitignore b/apache-kafka-3/.gitignore new file mode 100644 index 0000000000..549e00a2a9 --- /dev/null +++ b/apache-kafka-3/.gitignore @@ -0,0 +1,33 @@ +HELP.md +target/ +!.mvn/wrapper/maven-wrapper.jar +!**/src/main/**/target/ +!**/src/test/**/target/ + +### STS ### +.apt_generated +.classpath +.factorypath +.project +.settings +.springBeans +.sts4-cache + +### IntelliJ IDEA ### +.idea +*.iws +*.iml +*.ipr + +### NetBeans ### +/nbproject/private/ +/nbbuild/ +/dist/ +/nbdist/ +/.nb-gradle/ +build/ +!**/src/main/**/build/ +!**/src/test/**/build/ + +### VS Code ### +.vscode/ diff --git a/apache-kafka-3/pom.xml b/apache-kafka-3/pom.xml new file mode 100644 index 0000000000..9466068482 --- /dev/null +++ b/apache-kafka-3/pom.xml @@ -0,0 +1,50 @@ + + + 4.0.0 + + + com.baeldung + parent-boot-2 + 0.0.1-SNAPSHOT + ../parent-boot-2 + + + apache-kafka-3 + Apache Kafka 3 + Third module for Apache Kafka related articles + + 17 + + + + org.springframework.boot + spring-boot-starter + + + org.springframework.kafka + spring-kafka + + + + org.springframework.boot + spring-boot-starter-test + test + + + org.springframework.kafka + spring-kafka-test + test + + + + + + + org.springframework.boot + spring-boot-maven-plugin + + + + + diff --git a/spring-kafka-3/src/main/java/com/baeldung/spring/kafka/groupId/KafkaConfig.java b/apache-kafka-3/src/main/java/com/baeldung/apachekafka3/groupId/KafkaConfig.java similarity index 97% rename from spring-kafka-3/src/main/java/com/baeldung/spring/kafka/groupId/KafkaConfig.java rename to apache-kafka-3/src/main/java/com/baeldung/apachekafka3/groupId/KafkaConfig.java index 6514816424..da9a41f425 100644 --- a/spring-kafka-3/src/main/java/com/baeldung/spring/kafka/groupId/KafkaConfig.java +++ b/apache-kafka-3/src/main/java/com/baeldung/apachekafka3/groupId/KafkaConfig.java @@ -1,4 +1,4 @@ -package com.baeldung.spring.kafka.groupId; +package com.baeldung.apachekafka3.groupId; import java.util.HashMap; import java.util.Map; diff --git a/spring-kafka-3/src/main/java/com/baeldung/spring/kafka/groupId/KafkaErrorHandler.java b/apache-kafka-3/src/main/java/com/baeldung/apachekafka3/groupId/KafkaErrorHandler.java similarity index 76% rename from spring-kafka-3/src/main/java/com/baeldung/spring/kafka/groupId/KafkaErrorHandler.java rename to apache-kafka-3/src/main/java/com/baeldung/apachekafka3/groupId/KafkaErrorHandler.java index 00cc795272..3acd18a452 100644 --- a/spring-kafka-3/src/main/java/com/baeldung/spring/kafka/groupId/KafkaErrorHandler.java +++ b/apache-kafka-3/src/main/java/com/baeldung/apachekafka3/groupId/KafkaErrorHandler.java @@ -1,7 +1,6 @@ -package com.baeldung.spring.kafka.groupId; +package com.baeldung.apachekafka3.groupId; import org.apache.kafka.clients.consumer.Consumer; -import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.common.errors.RecordDeserializationException; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -13,12 +12,6 @@ class KafkaErrorHandler implements CommonErrorHandler { private static final Logger log = LoggerFactory.getLogger(KafkaErrorHandler.class); - @Override - public void handleRecord(@NonNull Exception exception, @NonNull ConsumerRecord record, @NonNull Consumer consumer, - @NonNull MessageListenerContainer container) { - handle(exception, consumer); - } - @Override public void handleOtherException(@NonNull Exception exception, @NonNull Consumer consumer, @NonNull MessageListenerContainer container, boolean batchListener) { diff --git a/spring-kafka-3/src/main/java/com/baeldung/spring/kafka/groupId/Main.java b/apache-kafka-3/src/main/java/com/baeldung/apachekafka3/groupId/Main.java similarity index 73% rename from spring-kafka-3/src/main/java/com/baeldung/spring/kafka/groupId/Main.java rename to apache-kafka-3/src/main/java/com/baeldung/apachekafka3/groupId/Main.java index 1045bc5192..e38699270f 100644 --- a/spring-kafka-3/src/main/java/com/baeldung/spring/kafka/groupId/Main.java +++ b/apache-kafka-3/src/main/java/com/baeldung/apachekafka3/groupId/Main.java @@ -1,14 +1,14 @@ -package com.baeldung.spring.kafka.groupId; - -import org.springframework.boot.SpringApplication; -import org.springframework.boot.autoconfigure.SpringBootApplication; -import org.springframework.context.annotation.ComponentScan; - -@SpringBootApplication -@ComponentScan(basePackages = "com.baeldung.spring.kafka.groupId") -public class Main { - - public static void main(String[] args) { - SpringApplication.run(Main.class, args); - } -} +package com.baeldung.apachekafka3.groupId; + +import org.springframework.boot.SpringApplication; +import org.springframework.boot.autoconfigure.SpringBootApplication; +import org.springframework.context.annotation.ComponentScan; + +@SpringBootApplication +@ComponentScan(basePackages = "com.baeldung.apachekafka3.groupId") +public class Main { + + public static void main(String[] args) { + SpringApplication.run(Main.class, args); + } +} diff --git a/spring-kafka-3/src/main/java/com/baeldung/spring/kafka/groupId/MyKafkaConsumer.java b/apache-kafka-3/src/main/java/com/baeldung/apachekafka3/groupId/MyKafkaConsumer.java similarity index 82% rename from spring-kafka-3/src/main/java/com/baeldung/spring/kafka/groupId/MyKafkaConsumer.java rename to apache-kafka-3/src/main/java/com/baeldung/apachekafka3/groupId/MyKafkaConsumer.java index e17d168338..e3c7220479 100644 --- a/spring-kafka-3/src/main/java/com/baeldung/spring/kafka/groupId/MyKafkaConsumer.java +++ b/apache-kafka-3/src/main/java/com/baeldung/apachekafka3/groupId/MyKafkaConsumer.java @@ -1,41 +1,41 @@ -package com.baeldung.spring.kafka.groupId; - -import java.util.concurrent.CountDownLatch; - -import org.apache.kafka.clients.consumer.Consumer; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; -import org.springframework.kafka.annotation.KafkaListener; -import org.springframework.messaging.handler.annotation.Payload; -import org.springframework.stereotype.Service; - -@Service -public class MyKafkaConsumer { - - private static final Logger LOGGER = LoggerFactory.getLogger(MyKafkaConsumer.class); - - private CountDownLatch latch = new CountDownLatch(1); - - private String payload; - - @KafkaListener(topics = "${kafka.topic.name:test-topic}", groupId = "${kafka.consumer.groupId:test-consumer-group}", concurrency = "4") - public void receive(@Payload String payload, Consumer consumer) { - LOGGER.info("Consumer='{}' received payload='{}'", consumer.groupMetadata() - .memberId(), payload); - this.payload = payload; - - latch.countDown(); - } - - public CountDownLatch getLatch() { - return latch; - } - - public void resetLatch() { - latch = new CountDownLatch(1); - } - - public String getPayload() { - return payload; - } -} +package com.baeldung.apachekafka3.groupId; + +import java.util.concurrent.CountDownLatch; + +import org.apache.kafka.clients.consumer.Consumer; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.kafka.annotation.KafkaListener; +import org.springframework.messaging.handler.annotation.Payload; +import org.springframework.stereotype.Service; + +@Service +public class MyKafkaConsumer { + + private static final Logger LOGGER = LoggerFactory.getLogger(MyKafkaConsumer.class); + + private CountDownLatch latch = new CountDownLatch(1); + + private String payload; + + @KafkaListener(topics = "${kafka.topic.name:test-topic}", clientIdPrefix = "neo", groupId = "${kafka.consumer.groupId:test-consumer-group}", concurrency = "4") + public void receive(@Payload String payload, Consumer consumer) { + LOGGER.info("Consumer='{}' received payload='{}'", consumer.groupMetadata() + .memberId(), payload); + this.payload = payload; + + latch.countDown(); + } + + public CountDownLatch getLatch() { + return latch; + } + + public void resetLatch() { + latch = new CountDownLatch(1); + } + + public String getPayload() { + return payload; + } +} diff --git a/spring-kafka-3/src/main/java/com/baeldung/spring/kafka/groupId/MyKafkaProducer.java b/apache-kafka-3/src/main/java/com/baeldung/apachekafka3/groupId/MyKafkaProducer.java similarity index 91% rename from spring-kafka-3/src/main/java/com/baeldung/spring/kafka/groupId/MyKafkaProducer.java rename to apache-kafka-3/src/main/java/com/baeldung/apachekafka3/groupId/MyKafkaProducer.java index c1d9a05d60..8707ce800d 100644 --- a/spring-kafka-3/src/main/java/com/baeldung/spring/kafka/groupId/MyKafkaProducer.java +++ b/apache-kafka-3/src/main/java/com/baeldung/apachekafka3/groupId/MyKafkaProducer.java @@ -1,24 +1,24 @@ -package com.baeldung.spring.kafka.groupId; - -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; -import org.springframework.beans.factory.annotation.Autowired; -import org.springframework.beans.factory.annotation.Value; -import org.springframework.kafka.core.KafkaTemplate; -import org.springframework.stereotype.Service; - -@Service -public class MyKafkaProducer { - - private static final Logger LOGGER = LoggerFactory.getLogger(MyKafkaProducer.class); - - @Value("${kafka.topic.name:test-topic}") - private String topic; - @Autowired - private KafkaTemplate kafkaTemplate; - - public void send(String payload) { - LOGGER.info("Sending payload='{}' to topic='{}'", payload, topic); - kafkaTemplate.send(topic, payload); - } -} +package com.baeldung.apachekafka3.groupId; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.annotation.Value; +import org.springframework.kafka.core.KafkaTemplate; +import org.springframework.stereotype.Service; + +@Service +public class MyKafkaProducer { + + private static final Logger LOGGER = LoggerFactory.getLogger(MyKafkaProducer.class); + + @Value("${kafka.topic.name:test-topic}") + private String topic; + @Autowired + private KafkaTemplate kafkaTemplate; + + public void send(String payload) { + LOGGER.info("Sending payload='{}' to topic='{}'", payload, topic); + kafkaTemplate.send(topic, payload); + } +} diff --git a/apache-kafka-3/src/main/resources/application.properties b/apache-kafka-3/src/main/resources/application.properties new file mode 100644 index 0000000000..0fce641151 --- /dev/null +++ b/apache-kafka-3/src/main/resources/application.properties @@ -0,0 +1 @@ +kafka.topic.name=test-topic diff --git a/spring-kafka-3/src/test/java/com/baeldung/spring/kafka/groupId/MainLiveTest.java b/apache-kafka-3/src/test/java/com/baeldung/apachekafka3/groupId/MainLiveTest.java similarity index 90% rename from spring-kafka-3/src/test/java/com/baeldung/spring/kafka/groupId/MainLiveTest.java rename to apache-kafka-3/src/test/java/com/baeldung/apachekafka3/groupId/MainLiveTest.java index c7e78f6e58..e10e41838a 100644 --- a/spring-kafka-3/src/test/java/com/baeldung/spring/kafka/groupId/MainLiveTest.java +++ b/apache-kafka-3/src/test/java/com/baeldung/apachekafka3/groupId/MainLiveTest.java @@ -1,44 +1,44 @@ -package com.baeldung.spring.kafka.groupId; - -import static org.hamcrest.MatcherAssert.assertThat; -import static org.hamcrest.Matchers.containsString; -import static org.junit.jupiter.api.Assertions.assertTrue; - -import java.util.concurrent.TimeUnit; - -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; -import org.springframework.context.annotation.ComponentScan; -import org.springframework.kafka.test.context.EmbeddedKafka; -import org.springframework.test.annotation.DirtiesContext; - -@SpringBootTest(classes = Main.class) -@ComponentScan(basePackages = "com.baeldung.spring.kafka.groupId") -@DirtiesContext -@EmbeddedKafka(partitions = 4, topics = { "${kafka.topic.name:test-topic}" }, brokerProperties = { "listeners=PLAINTEXT://localhost:9092", "port=9092" }) -public class MainLiveTest { - - @Autowired - private MyKafkaConsumer consumer; - @Autowired - private MyKafkaProducer producer; - - @BeforeEach - void setup() { - consumer.resetLatch(); - } - - @Test - public void givenEmbeddedKafkaBroker_whenSendingWithSimpleProducer_thenMessageReceived() throws Exception { - String data = "Test 123..."; - - producer.send(data); - - boolean messageConsumed = consumer.getLatch() - .await(10, TimeUnit.SECONDS); - assertTrue(messageConsumed); - assertThat(consumer.getPayload(), containsString(data)); - } -} +package com.baeldung.apachekafka3.groupId; + +import static org.hamcrest.MatcherAssert.assertThat; +import static org.hamcrest.Matchers.containsString; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import java.util.concurrent.TimeUnit; + +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; +import org.springframework.context.annotation.ComponentScan; +import org.springframework.kafka.test.context.EmbeddedKafka; +import org.springframework.test.annotation.DirtiesContext; + +@SpringBootTest(classes = Main.class) +@ComponentScan(basePackages = "com.baeldung.apachekafka3.groupId") +@DirtiesContext +@EmbeddedKafka(partitions = 4, topics = { "${kafka.topic.name:test-topic}" }, brokerProperties = { "listeners=PLAINTEXT://localhost:9092", "port=9092" }) +public class MainLiveTest { + + @Autowired + private MyKafkaConsumer consumer; + @Autowired + private MyKafkaProducer producer; + + @BeforeEach + void setup() { + consumer.resetLatch(); + } + + @Test + public void givenEmbeddedKafkaBroker_whenSendingWithSimpleProducer_thenMessageReceived() throws Exception { + String data = "Test 123..."; + + producer.send(data); + + boolean messageConsumed = consumer.getLatch() + .await(10, TimeUnit.SECONDS); + assertTrue(messageConsumed); + assertThat(consumer.getPayload(), containsString(data)); + } +} diff --git a/pom.xml b/pom.xml index 030b8e1843..3a3ba46d48 100644 --- a/pom.xml +++ b/pom.xml @@ -670,8 +670,9 @@ apache-httpclient-2 apache-httpclient4 apache-httpclient - apache-kafka-2 apache-kafka + apache-kafka-2 + apache-kafka-3 apache-libraries-2 apache-libraries apache-olingo