Apache Kafka – це одна з найпопулярніших платформ обміну повідомленнями та потокової обробки даних. Вона дозволяє обмінюватися великими обсягами даних між різними мікросервісами чи системами. Інтеграція Kafka зі Spring Boot є популярним вибором розробників, оскільки Spring Boot спрощує розробку та забезпечує потужну підтримку Kafka через Spring for Apache Kafka.
Apache Kafka — це розподілена система обміну повідомленнями з відкритим вихідним кодом, розроблена компанією LinkedIn у 2011 році і пізніше передана до Apache Software Foundation. Вона дозволяє передавати дані між різними додатками та системами в реальному часі, забезпечуючи зв’язок за допомогою потоків подій.
Kafka працює як шина даних або процесингова платформа для подій (event streaming platform), об’єднуючи продуцентів (відправників) даних і консюмерів (споживачів) через потоки повідомлень (topics). Це робить її ідеальним інструментом для побудови архітектури мікросервісів, аналітики даних, систем моніторингу та багатьох інших задач.
У цій статті ми розглянемо, як підключити Kafka до Spring Boot і налаштувати простий приклад із продюсером (producer) та консюмером (consumer).
Основні компоненти Apache Kafka
Apache Kafka складається з кількох ключових компонентів:
- Producer (відправник) — відповідає за надсилання даних до тем у Kafka. Дані надсилаються у формі повідомлень (events).
- Consumer (споживач) — отримує дані з теми Kafka і реагує на них. Споживачі можуть бути налаштовані для отримання історичних, так і нових даних.
- Broker (брокер) — основний компонент системи, що відповідає за зберігання та обробку повідомлень. Брокери забезпечують масштабованість і балансування навантаження.
- Topic (тема) — логічний канал, через який передаються повідомлення. Відправники надсилають дані до тем, а споживачі отримують їх.
- Zookeeper — координатор кластеру Kafka, який керує розподіленням даних між брокерами, конфігураціями та іншими аспектами.
- Partition (розділ) — кожна тема розподілена на кілька частин (partition), що дозволяє масштабувати обробку даних і забезпечувати високу продуктивність системи.
Apache Kafka часто використовується у таких сценаріях:
- Потокова аналітика: Збір даних у реальному часі з сенсорів, додатків або систем моніторингу для швидкої обробки.
- Обробка великих даних: Kafka часто використовується як канал зв’язку між різними компонентами системи великих даних.
- Архітектура мікросервісів: Забезпечення ефективної комунікації між мікросервісами через події.
- Системи моніторингу та логування: Обробка логів і оперативне визначення проблем у системі.
Крок 1: Підключення Kafka до проекту
Перед інтеграцією Spring Boot із Kafka необхідно переконатися, що Kafka встановлена та працює у вашому середовищі.
- Завантаження та встановлення Kafka:
- Завантажте Apache Kafka з офіційного сайту kafka.apache.org.
- Розпакуйте архів і перейдіть до папки Kafka.
- Запуск Zookeeper та Kafka Broker:
- Zookeeper:
Виконайте команду для запуску Zookeeper (Kafka використовує Zookeeper для координації):bin/zookeeper-server-start.sh config/zookeeper.properties - Kafka Broker:
Запустіть Kafka Broker:bin/kafka-server-start.sh config/server.properties
- Zookeeper:
- Створення топіка (topic): Для обміну повідомленнями використовується концепція “топіків”. Створіть топік (наприклад,
test-topic) за допомогою команди:bin/kafka-topics.sh --create --topic test-topic --bootstrap-server localhost:9092 --partitions 1 --replication-factor 1
Крок 2: Створення проєкту Spring Boot
- Ініціалізація проєкту: Ви можете ініціалізувати Spring Boot проєкт за допомогою Spring Initializr. Додайте такі залежності:
- Spring Web
- Spring for Apache Kafka
- Додайте залежності у файл
pom.xml: Якщо ви самостійно редагуєте залежності, переконайтеся, що у вашомуpom.xmlнаявні такі бібліотеки:
<dependencies>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.kafka</groupId>
<artifactId>spring-kafka</artifactId>
</dependency>
</dependencies>XMLСтворіть новий конфігураційний файл application.yml або application.properties. Наприклад, у application.yml можна вказати наступні налаштування:
spring:
kafka:
bootstrap-servers:
- localhost:9092 # Адреса Kafka брокера
consumer:
group-id: my-consumer-group # Ідентифікатор групи споживачів
auto-offset-reset: earliest # Установлює точку зчитування (earliest/latest/none)
key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
value-deserializer: org.apache.kafka.common.serialization.StringDeserializer
producer:
key-serializer: org.apache.kafka.common.serialization.StringSerializer
value-serializer: org.apache.kafka.common.serialization.StringSerializer
retries: 2 # Кількість повторних спроб при відмові
ack: all # Очікування підтвердження від усіх реплік
YAMLКрок 3: Реалізація продюсера (producer)
Продюсер відповідає за надсилання даних у Kafka.
Створіть клас KafkaProducer:
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.stereotype.Service;
@Service
public class KafkaProducer {
private final KafkaTemplate<String, String> kafkaTemplate;
public KafkaProducer(KafkaTemplate<String, String> kafkaTemplate) {
this.kafkaTemplate = kafkaTemplate;
}
public void sendMessage(String topic, String message) {
kafkaTemplate.send(topic, message);
System.out.println("Повідомлення надіслано в Kafka: " + message);
}
}JavaКрок 4: Реалізація консюмера (consumer)
Консюмер отримує та обробляє повідомлення, що надходять у Kafka.
Створення клас KafkaConsumer:
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.stereotype.Service;
@Service
public class KafkaConsumer {
@KafkaListener(topics = "test-topic", groupId = "my-group")
public void consumeMessage(ConsumerRecord<String, String> record) {
System.out.println("Отримано повідомлення з Kafka: " + record.value());
}
}JavaКрок 5: Реалізація REST Controller для тестування
Для інтеграції з Kafka можна створити простий REST-контролер.
Створіть клас KafkaController:
import org.springframework.web.bind.annotation.*;
@RestController
@RequestMapping("/kafka")
public class KafkaController {
private final KafkaProducer kafkaProducer;
public KafkaController(KafkaProducer kafkaProducer) {
this.kafkaProducer = kafkaProducer;
}
@PostMapping("/publish")
public String sendMessageToKafka(@RequestParam("message") String message) {
kafkaProducer.sendMessage("test-topic", message);
return "Повідомлення надіслано у Kafka!";
}
}JavaКрок 6: Запуск та тестування
- Запустіть Spring Boot застосунок: Використовуйте команду:bash
mvn spring-boot:run- Відправка повідомлень через REST API: Використовуйте інструмент, як-от Postman або curl, щоб відправити повідомлення до Kafka:
curl -X POST 'http://localhost:8080/kafka/publish?message=Привіт, Kafka!'- Перевірка повідомлень у консюмері: Ви повинні побачити отримане повідомлення у логах консюмера.
І на останок
У цій статті ми розглянули, як налаштувати інтеграцію Apache Kafka зі Spring Boot. Ми створили продюсера та консюмера для відправки та отримання повідомлень, а також налаштували базові конфігурації Kafka в application.yml.
Інтеграція Kafka зі Spring Boot забезпечує гнучкість і простоту в розробці застосунків, які потребують обробки потоків даних у реальному часі. Це потужний інструмент для побудови сучасних масштабованих мікросервісних архітектур.





