From bdd7613651de9ebb028824b54d2a192e860e75be Mon Sep 17 00:00:00 2001 From: AnamNadir Date: Tue, 2 Jul 2024 16:38:43 +0300 Subject: [PATCH] [HW_1] --- webinar-01/actions.md | 30 ++++++---- .../java/com/prosoft/KafkaProducerApp.java | 59 +++++++++++-------- .../java/com/prosoft/config/KafkaConfig.java | 7 ++- .../main/java/com/prosoft/domain/Symbol.java | 15 +++++ .../prosoft/serializer/SymbolSerializer.java | 28 +++++++++ 5 files changed, 99 insertions(+), 40 deletions(-) create mode 100644 webinar-02/producer-service/src/main/java/com/prosoft/domain/Symbol.java create mode 100644 webinar-02/producer-service/src/main/java/com/prosoft/serializer/SymbolSerializer.java diff --git a/webinar-01/actions.md b/webinar-01/actions.md index e8459a3..ed24332 100644 --- a/webinar-01/actions.md +++ b/webinar-01/actions.md @@ -13,27 +13,31 @@ docker compose ps -a docker exec -ti kafka /usr/bin/kafka-topics --list --bootstrap-server kafka:9093 ``` -`4.` Создаем новый "topic1" +`4.` Создаем новый "vowels" ```shell -docker exec -ti kafka /usr/bin/kafka-topics --create --topic topic1 --bootstrap-server localhost:9093 +docker exec -ti kafka /usr/bin/kafka-topics --create --topic vowels --bootstrap-server localhost:9093 ``` - -`5.` Отправляем сообщение "topic1": появляется консоль ввода сообщений, вводим сообщение одно за другим, разделяя Enter и в конце нажимаем в Win ctrl+D (в MacOS: control+C) +`5.` Создаем новый "consonants" ```shell -docker exec -ti kafka /usr/bin/kafka-console-producer --topic topic1 --bootstrap-server kafka:9093 +docker exec -ti kafka /usr/bin/kafka-topics --create --topic consonants --bootstrap-server localhost:9093 ``` - -`6.` Получить сообщения +`6.` Отправляем сообщение "vowels": появляется консоль ввода сообщений, вводим сообщение одно за другим, разделяя Enter и в конце нажимаем в Win ctrl+D (в MacOS: control+C) ```shell -docker exec -ti kafka /usr/bin/kafka-console-consumer --from-beginning --topic topic1 --bootstrap-server localhost:9093 +docker exec -ti kafka /usr/bin/kafka-console-producer --topic vowels --bootstrap-server kafka:9093 ``` - -`7.` Получить сообщения как consumer1 +`7.` Отправляем сообщение "consonants": появляется консоль ввода сообщений, вводим сообщение одно за другим, разделяя Enter и в конце нажимаем в Win ctrl+D (в MacOS: control+C) ```shell -docker exec -ti kafka /usr/bin/kafka-console-consumer --group consumer1 --topic topic1 --bootstrap-server localhost:9093 +docker exec -ti kafka /usr/bin/kafka-console-producer --topic consonants --bootstrap-server kafka:9093 ``` - -`8.` Останавливаем контейнеры, удаляем контейнеры, удаляем неиспользуемые тома: +`8.` Получить сообщения из vowels +```shell +docker exec -ti kafka /usr/bin/kafka-console-consumer --from-beginning --topic vowels --bootstrap-server localhost:9093 +``` +`9.` Получить сообщения из consonants +```shell +docker exec -ti kafka /usr/bin/kafka-console-consumer --from-beginning --topic consonants --bootstrap-server localhost:9093 +``` +`10.` Останавливаем контейнеры, удаляем контейнеры, удаляем неиспользуемые тома: ```shell docker compose stop docker container prune -f diff --git a/webinar-02/producer-service/src/main/java/com/prosoft/KafkaProducerApp.java b/webinar-02/producer-service/src/main/java/com/prosoft/KafkaProducerApp.java index 27eb561..e08c5f6 100644 --- a/webinar-02/producer-service/src/main/java/com/prosoft/KafkaProducerApp.java +++ b/webinar-02/producer-service/src/main/java/com/prosoft/KafkaProducerApp.java @@ -1,39 +1,38 @@ package com.prosoft; import com.prosoft.config.KafkaConfig; -import com.prosoft.domain.Person; -import org.apache.kafka.clients.producer.Callback; +import com.prosoft.domain.Symbol; import org.apache.kafka.clients.producer.KafkaProducer; import org.apache.kafka.clients.producer.ProducerRecord; -import org.apache.kafka.clients.producer.RecordMetadata; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import java.time.LocalDateTime; -import java.time.format.DateTimeFormatter; /** - * Webinar-02: Kafka producer-service (отправка объектов класса Person) + * Webinar-02: Kafka producer-service (отправка объектов класса Symbol) * Использования метода producer.send(producerRecord) с обработкой результата отправки через Callback. */ public class KafkaProducerApp { private static final Logger logger = LoggerFactory.getLogger(KafkaProducerApp.class); - private static final int MAX_MESSAGE = 10; + private static final int MAX_MESSAGE = 5; + private static final String[] VOWELS = {"a", "e", "i", "o", "u"}; + private static final String[] CONSONANTS = {"b", "c", "d", "h", "z"}; public static void main(String[] args) { - try (KafkaProducer producer = new KafkaProducer<>(KafkaConfig.getProducerConfig())) { + try (KafkaProducer producer = new KafkaProducer<>(KafkaConfig.getProducerConfig())) { for (int i = 0; i < MAX_MESSAGE; i++) { - Person person = createPerson(i); + Symbol vowelSymbol = createSymbol(i, VOWELS[i], "red", "vowel"); + Symbol consonantSymbol = createSymbol(i + 5, CONSONANTS[i], "blue", "consonant"); /** * Конструктор ProducerRecord(topic, partition, timestamp, key, value) принимает в качестве аргументов: * - topic - номер топика * - partition - номер партиции (опция) * - timestamp - время создания сообщения (опция) - * - key - ключ id экземпляра Person (опция) - * - value - объект Person + * - key - ключ id экземпляра Symbol (опция) + * - value - объект Symbol * * Варианты конструкторов: * - ProducerRecord(topic, value) @@ -42,24 +41,35 @@ public static void main(String[] args) { * - ProducerRecord(topic, partition, key, value, headers) */ long timestamp = System.currentTimeMillis(); - ProducerRecord producerRecord = new ProducerRecord<>(KafkaConfig.TOPIC, KafkaConfig.PARTITION, - timestamp, person.getId(), person); + ProducerRecord producerRecordVowels = new ProducerRecord<>(KafkaConfig.TOPIC_VOWELS, KafkaConfig.PARTITION, + timestamp, vowelSymbol.getId(), vowelSymbol); + ProducerRecord producerRecordConsonants = new ProducerRecord<>(KafkaConfig.TOPIC_CONSONANTS, KafkaConfig.PARTITION, + timestamp, consonantSymbol.getId(), consonantSymbol); /** * Анонимный внутренний класс (Callback), содержащий только один метод onCompletion(), можно записать * через лямбду */ - producer.send(producerRecord, new Callback() { - @Override - public void onCompletion(RecordMetadata recordMetadata, Exception e) { - if (e != null) { - logger.error("Error sending message: {}", e.getMessage(), e); - } else { - logger.info("Sent record: key={}, value={}, partition={}, offset={}", - person.getId(), person, recordMetadata.partition(), recordMetadata.offset()); } + producer.send(producerRecordVowels, (recordMetadata, e) -> { + if (e != null) { + logger.error("Error sending message: {}", e.getMessage(), e); + } else { + logger.info("Sent record: key={}, value={}, partition={}, offset={}", + vowelSymbol.getId(), vowelSymbol, recordMetadata.partition(), recordMetadata.offset()); } }); - logger.info("Отправлено сообщение: key-{}, value-{}", i, person); + logger.info("Отправлено сообщение: key-{}, value-{}", i, vowelSymbol); + + + producer.send(producerRecordConsonants, (recordMetadata, e) -> { + if (e != null) { + logger.error("Error sending message: {}", e.getMessage(), e); + } else { + logger.info("Sent record: key={}, value={}, partition={}, offset={}", + consonantSymbol.getId(), consonantSymbol, recordMetadata.partition(), recordMetadata.offset()); + } + }); + logger.info("Отправлено сообщение: key-{}, value-{}", i, consonantSymbol); } logger.info("Отправка завершена."); } catch (Exception e) { @@ -67,9 +77,8 @@ public void onCompletion(RecordMetadata recordMetadata, Exception e) { } } // todo показать без try -with-resources c вызовом .flush() .close() .close(Duration.ofSeconds(60)) - private static Person createPerson(int index) { - String currentTime = LocalDateTime.now().format(DateTimeFormatter.ofPattern("dd-MM-yyyy-HH-mm-ss")); - return new Person(index, "FirstName-" + currentTime, "LastName" + index, 20 + index); + private static Symbol createSymbol(int index, String value, String color, String type) { + return new Symbol(index, value, color, type); } } diff --git a/webinar-02/producer-service/src/main/java/com/prosoft/config/KafkaConfig.java b/webinar-02/producer-service/src/main/java/com/prosoft/config/KafkaConfig.java index a56b502..a83fddb 100644 --- a/webinar-02/producer-service/src/main/java/com/prosoft/config/KafkaConfig.java +++ b/webinar-02/producer-service/src/main/java/com/prosoft/config/KafkaConfig.java @@ -1,6 +1,7 @@ package com.prosoft.config; import com.prosoft.serializer.PersonSerializer; +import com.prosoft.serializer.SymbolSerializer; import org.apache.kafka.clients.producer.ProducerConfig; import org.apache.kafka.common.serialization.LongSerializer; @@ -12,9 +13,11 @@ */ public class KafkaConfig { - public static final String TOPIC = "topic2"; + public static final String TOPIC_VOWELS = "vowels"; + public static final String TOPIC_CONSONANTS = "consonants"; public static final int PARTITION = 0; + private static final String BOOTSTRAP_SERVERS = "localhost:9091, localhost:9092, localhost:9093"; private KafkaConfig() { @@ -45,7 +48,7 @@ public static Properties getProducerConfig() { properties.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, LongSerializer.class.getName()); /** Использование PersonSerializer для сериализации значения (Value) */ - properties.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, PersonSerializer.class.getName()); + properties.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, SymbolSerializer.class.getName()); /** delivery.timeout.ms - максимальное время ожидания для успешной отправки сообщения. * Это включает время, которое сообщение находится в очереди, а также все попытки повторной отправки. diff --git a/webinar-02/producer-service/src/main/java/com/prosoft/domain/Symbol.java b/webinar-02/producer-service/src/main/java/com/prosoft/domain/Symbol.java new file mode 100644 index 0000000..ef4d955 --- /dev/null +++ b/webinar-02/producer-service/src/main/java/com/prosoft/domain/Symbol.java @@ -0,0 +1,15 @@ +package com.prosoft.domain; + +import lombok.AllArgsConstructor; +import lombok.Data; +import lombok.NoArgsConstructor; + +@Data +@AllArgsConstructor +@NoArgsConstructor +public class Symbol { + private long id; + private String value; + private String color; + private String type; +} diff --git a/webinar-02/producer-service/src/main/java/com/prosoft/serializer/SymbolSerializer.java b/webinar-02/producer-service/src/main/java/com/prosoft/serializer/SymbolSerializer.java new file mode 100644 index 0000000..054960b --- /dev/null +++ b/webinar-02/producer-service/src/main/java/com/prosoft/serializer/SymbolSerializer.java @@ -0,0 +1,28 @@ +package com.prosoft.serializer; + +import com.fasterxml.jackson.databind.ObjectMapper; +import com.prosoft.domain.Symbol; +import org.apache.kafka.common.serialization.Serializer; + +import java.util.Map; + +public class SymbolSerializer implements Serializer { + private final ObjectMapper objectMapper = new ObjectMapper(); + + @Override + public void configure(Map configs, boolean isKey) { + } + + @Override + public byte[] serialize(String topic, Symbol data) { + try { + return objectMapper.writeValueAsBytes(data); + } catch (Exception e) { + throw new RuntimeException("Error serializing Person to JSON", e); + } + } + + @Override + public void close() { + } +}