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