Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
30 changes: 17 additions & 13 deletions webinar-01/actions.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
@@ -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<Long, Person> producer = new KafkaProducer<>(KafkaConfig.getProducerConfig())) {
try (KafkaProducer<Long, Symbol> 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)
Expand All @@ -42,34 +41,44 @@ public static void main(String[] args) {
* - ProducerRecord(topic, partition, key, value, headers)
*/
long timestamp = System.currentTimeMillis();
ProducerRecord<Long, Person> producerRecord = new ProducerRecord<>(KafkaConfig.TOPIC, KafkaConfig.PARTITION,
timestamp, person.getId(), person);
ProducerRecord<Long, Symbol> producerRecordVowels = new ProducerRecord<>(KafkaConfig.TOPIC_VOWELS, KafkaConfig.PARTITION,
timestamp, vowelSymbol.getId(), vowelSymbol);
ProducerRecord<Long, Symbol> 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) {
logger.error("Ошибка при отправке сообщений в Kafka", 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);
}

}
Original file line number Diff line number Diff line change
@@ -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;

Expand All @@ -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() {
Expand Down Expand Up @@ -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 - максимальное время ожидания для успешной отправки сообщения.
* Это включает время, которое сообщение находится в очереди, а также все попытки повторной отправки.
Expand Down
Original file line number Diff line number Diff line change
@@ -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;
}
Original file line number Diff line number Diff line change
@@ -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<Symbol> {
private final ObjectMapper objectMapper = new ObjectMapper();

@Override
public void configure(Map<String, ?> 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() {
}
}