From 5dc9e85c2dcc49d73f5ff95886e358ed7134cf00 Mon Sep 17 00:00:00 2001 From: Ralf Ueberfuhr Date: Fri, 4 Aug 2023 06:02:16 +0200 Subject: [PATCH 1/2] BAEL-6415: Add simple Kafka client example (with BootStrap Servers setting) --- apache-kafka-3/pom.xml | 27 +++++++++++++ .../java/com/baeldung/kafka/KafkaJava.java | 39 +++++++++++++++++++ pom.xml | 2 + 3 files changed, 68 insertions(+) create mode 100644 apache-kafka-3/pom.xml create mode 100644 apache-kafka-3/src/main/java/com/baeldung/kafka/KafkaJava.java diff --git a/apache-kafka-3/pom.xml b/apache-kafka-3/pom.xml new file mode 100644 index 0000000000..6c4b5284dd --- /dev/null +++ b/apache-kafka-3/pom.xml @@ -0,0 +1,27 @@ + + + 4.0.0 + apache-kafka-3 + apache-kafka-3 + + + com.baeldung + parent-modules + 1.0.0-SNAPSHOT + + + + + org.apache.kafka + kafka-clients + ${kafka.version} + + + + + 3.5.1 + + + \ No newline at end of file diff --git a/apache-kafka-3/src/main/java/com/baeldung/kafka/KafkaJava.java b/apache-kafka-3/src/main/java/com/baeldung/kafka/KafkaJava.java new file mode 100644 index 0000000000..bc2eeadb78 --- /dev/null +++ b/apache-kafka-3/src/main/java/com/baeldung/kafka/KafkaJava.java @@ -0,0 +1,39 @@ +package com.baeldung.kafka; + +import org.apache.kafka.clients.consumer.*; +import org.apache.kafka.common.serialization.LongDeserializer; +import org.apache.kafka.common.serialization.StringDeserializer; + +import java.time.Duration; +import java.util.Arrays; +import java.util.Properties; + +public class KafkaJava { + + public static void main(String[] args) { + try(final Consumer consumer = createConsumer()) { + ConsumerRecords records = consumer.poll(Duration.ofMinutes(1)); + for(ConsumerRecord record : records) { + System.out.println(record.value()); + } + } + } + + private static Consumer createConsumer() { + final Properties props = new Properties(); + props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, + "localhost:9092,another-host.com:29092"); + props.put(ConsumerConfig.GROUP_ID_CONFIG, + "MySampleConsumer"); + props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, + LongDeserializer.class.getName()); + props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, + StringDeserializer.class.getName()); + // Create the consumer using props. + final Consumer consumer = new KafkaConsumer(props); + // Subscribe to the topic. + consumer.subscribe(Arrays.asList("samples")); + return consumer; + } + +} diff --git a/pom.xml b/pom.xml index 8fb2234a23..e11395c709 100644 --- a/pom.xml +++ b/pom.xml @@ -820,6 +820,7 @@ antlr apache-kafka apache-kafka-2 + apache-kafka-3 apache-olingo apache-poi-2 @@ -1089,6 +1090,7 @@ antlr apache-kafka apache-kafka-2 + apache-kafka-3 apache-olingo apache-poi-2 From 1df07e6696576d363890d67384e9dffb8767338d Mon Sep 17 00:00:00 2001 From: Ralf Ueberfuhr Date: Tue, 15 Aug 2023 07:47:27 +0200 Subject: [PATCH 2/2] BAEL-6415: Merged Kafka-3 project into Kafka-2 project --- .../SimpleConsumerWithBootStrapServers.java | 4 +-- apache-kafka-3/pom.xml | 27 ------------------- pom.xml | 2 -- 3 files changed, 2 insertions(+), 31 deletions(-) rename apache-kafka-3/src/main/java/com/baeldung/kafka/KafkaJava.java => apache-kafka-2/src/main/java/com/baeldung/kafka/consumer/SimpleConsumerWithBootStrapServers.java (94%) delete mode 100644 apache-kafka-3/pom.xml diff --git a/apache-kafka-3/src/main/java/com/baeldung/kafka/KafkaJava.java b/apache-kafka-2/src/main/java/com/baeldung/kafka/consumer/SimpleConsumerWithBootStrapServers.java similarity index 94% rename from apache-kafka-3/src/main/java/com/baeldung/kafka/KafkaJava.java rename to apache-kafka-2/src/main/java/com/baeldung/kafka/consumer/SimpleConsumerWithBootStrapServers.java index bc2eeadb78..7501b40056 100644 --- a/apache-kafka-3/src/main/java/com/baeldung/kafka/KafkaJava.java +++ b/apache-kafka-2/src/main/java/com/baeldung/kafka/consumer/SimpleConsumerWithBootStrapServers.java @@ -1,4 +1,4 @@ -package com.baeldung.kafka; +package com.baeldung.kafka.consumer; import org.apache.kafka.clients.consumer.*; import org.apache.kafka.common.serialization.LongDeserializer; @@ -8,7 +8,7 @@ import java.time.Duration; import java.util.Arrays; import java.util.Properties; -public class KafkaJava { +public class SimpleConsumerWithBootStrapServers { public static void main(String[] args) { try(final Consumer consumer = createConsumer()) { diff --git a/apache-kafka-3/pom.xml b/apache-kafka-3/pom.xml deleted file mode 100644 index 6c4b5284dd..0000000000 --- a/apache-kafka-3/pom.xml +++ /dev/null @@ -1,27 +0,0 @@ - - - 4.0.0 - apache-kafka-3 - apache-kafka-3 - - - com.baeldung - parent-modules - 1.0.0-SNAPSHOT - - - - - org.apache.kafka - kafka-clients - ${kafka.version} - - - - - 3.5.1 - - - \ No newline at end of file diff --git a/pom.xml b/pom.xml index e11395c709..8fb2234a23 100644 --- a/pom.xml +++ b/pom.xml @@ -820,7 +820,6 @@ antlr apache-kafka apache-kafka-2 - apache-kafka-3 apache-olingo apache-poi-2 @@ -1090,7 +1089,6 @@ antlr apache-kafka apache-kafka-2 - apache-kafka-3 apache-olingo apache-poi-2