Подробное руководство по интеграции Kafka в приложение Spring Boot
Apache Kafka - это популярная платформа распределенной потоковой передачи данных, которая позволяет создавать масштабируемые, отказоустойчивые и высокопроизводительные приложения. В этой статье мы предоставим Вашему вниманию пошаговое руководство по интеграции Kafka в приложение Spring Boot, дополненное подробными примерами кода и пояснениями. К концу этого руководства Вы поймете, как можно внедрить Kafka в Ваши проекты Spring Boot для того, чтобы повысить производительность и расширить возможности Вашего приложения.
Необходимые условия
- Базовые знания Java и Spring Boot
- Установленный Java Development Kit (JDK) 8 или более поздняя версия
- Подходящая среда разработки, например IntelliJ IDEA или Eclipse.
- Установленный и запущенный на Вашей локальной машине Apache Kafka
Настройка проекта Spring Boot
Для начала создадим новый проект Spring Boot со следующими зависимостями:
- Spring для Apache Kafka
- Spring Boot Starter Web
Вы можете создать проект с помощью Spring Initializr, а можете вручную добавить эти зависимости в файл Maven pom.xml или Gradle build.gradle.
Maven:
<dependencies>
<dependency>
<groupId>org.springframework.kafka</groupId>
<artifactId>spring-kafka</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>
</dependencies>Gradle:
dependencies {
implementation 'org.springframework.kafka:spring-kafka'
implementation 'org.springframework.boot:spring-boot-starter-web'
}
Настройка Kafka в Spring Boot
Далее нужно настроить Kafka в приложении Spring Boot. Начнем с добавления следующих свойств в файл application.properties:
spring.kafka.bootstrap-servers=localhost:9092 spring.kafka.consumer.group-id=my-group-id
Здесь spring.kafka.bootstrap-servers указывает адрес Вашего брокера Kafka, а spring.kafka.consumer.group-id - идентификатор группы консюмеров для Вашего приложения.
Создание продюсера Kafka
Чтобы отправлять сообщения в Kafka, нам нужно создать продюсера Kafka. Сначала создайте новый пакет com.example.kafka.producer, а затем добавьте в него следующий класс KafkaProducerConfig:
package com.example.kafka.producer;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.common.serialization.StringSerializer;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.kafka.core.DefaultKafkaProducerFactory;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.core.ProducerFactory;
import java.util.HashMap;
import java.util.Map;
@Configuration
public class KafkaProducerConfig {
@Bean
public ProducerFactory<String, String> producerFactory() {
Map<String, Object> configProps = new HashMap<>();
configProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
configProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class); configProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class); return new DefaultKafkaProducerFactory<>(configProps);
} @Bean
public KafkaTemplate<String, String> kafkaTemplate() {
return new KafkaTemplate<>(producerFactory());
} }
Затем создайте класс MessageProducer в том же пакете для отправки сообщений в топик Kafka:
package com.example.kafka.producer;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.stereotype.Component;
@Component
public class MessageProducer {
@Autowired
private KafkaTemplate<String, String> kafkaTemplate;
public void sendMessage(String topic, String message) {
kafkaTemplate.send(topic, message); } }
Создание консюмера Kafka
Теперь для получения сообщений из топика Kafka давайте создадим консюмера Kafka. Создайте новый пакет com.example.kafka.consumer и добавьте в него следующий класс KafkaConsumerConfig:
package com.example.kafka.consumer;
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.common.serialization.StringDeserializer;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory;
import org.springframework.kafka.core.ConsumerFactory;
import org.springframework.kafka.core.DefaultKafkaConsumerFactory;
import java.util.HashMap;
import java.util.Map;
@Configuration
public class KafkaConsumerConfig {
@Bean
public ConsumerFactory<String, String> consumerFactory() {
Map<String, Object> configProps = new HashMap<>();
configProps.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
configProps.put(ConsumerConfig.GROUP_ID_CONFIG, "my-group-id");
configProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); configProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); return new DefaultKafkaConsumerFactory<>(configProps);
} @Bean
public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory() {
ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>();
factory.setConsumerFactory(consumerFactory()); return factory;
} }
Создайте класс MessageConsumer в том же пакете для того, чтобы слушать сообщения из топика Kafka:
package com.example.kafka.consumer;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.stereotype.Component;
@Component
public class MessageConsumer {
@KafkaListener(topics = "my-topic", groupId = "my-group-id")
public void listen(String message) {
System.out.println("Received message: " + message);
} }
Тестирование интеграции Kafka
Наконец, давайте протестируем нашу интеграцию с Kafka, отправив и получив сообщения. Создайте новый REST-контроллер в пакете com.example.kafka.controller:
package com.example.kafka.controller;
import com.example.kafka.producer.MessageProducer;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestParam;
import org.springframework.web.bind.annotation.RestController;
@RestController
public class KafkaController {
@Autowired
private MessageProducer messageProducer;
@PostMapping("/send")
public String sendMessage(@RequestParam("message") String message) {
messageProducer.sendMessage("my-topic", message);
return "Message sent: " + message;
} }
Запустите приложение Spring Boot и с помощью инструмента вроде Postman или curl отправьте POST-запрос на http://localhost:8080/send?message=Hello_Kafka. Сообщение будет отправлено в топик Kafka, а консюмер получит и выведет сообщение на консоль.




