From 5dc9e85c2dcc49d73f5ff95886e358ed7134cf00 Mon Sep 17 00:00:00 2001 From: Ralf Ueberfuhr Date: Fri, 4 Aug 2023 06:02:16 +0200 Subject: [PATCH] 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