Введение в Apache Kafka со Spring Boot
Apache Kafka - это распределенная платформа потоковой передачи данных, предназначенная для обработки больших объемов данных в режиме реального времени. Архитектура Kafka состоит из четырех основных компонентов: продюсеров, брокеров, консюмеров и ZooKeeper.
- Продюсеры: Продюсеры - это клиенты, которые публикуют сообщения в топики Kafka. Ими могут быть любые приложения или системы, генерирующие данные, например датчики, веб-серверы или базы данных. Когда продюсер публикует сообщение в топике, он отправляет его брокеру Kafka.
- Брокеры: Брокеры - это серверы, которые образуют кластер Kafka. Они отвечают за прием сообщений от производителей, хранение сообщений в распределенном журнале фиксации и передачу сообщений консюмерам. Каждый брокер в кластере Kafka идентифицируется уникальным целым числом, называемым идентификатором брокера.
- Консюмеры: Консюмеры - это клиенты, которые читают сообщения из топиков Kafka. Ими могут быть любые приложения или системы, обрабатывающие данные, например, аналитические системы, поисковые системы или алгоритмы машинного обучения. Когда консюмер подписывается на топик, он сразу же получает сообщения из одного или нескольких разделов этого топика.
- ZooKeeper: ZooKeeper - это распределенный координационный сервис, который предназначен для управления состоянием кластера, выборов лидеров и выполнения различных административных задач. Он поддерживает список всех брокеров в кластере Kafka, а также их текущее состояние и метаданные.
Модель обмена сообщениями Kafka основана на топиках и разделах. Топик - это логическая категория или имя ленты, в которую публикуются сообщения. Каждый топик делится на один или несколько разделов, которые представляют собой упорядоченные, неизменяемые последовательности сообщений, которые могут храниться на разных брокерах в кластере Kafka.
Когда продюсер публикует сообщение в тему, он указывает имя топика и ключ, который используется для определения раздела, в который записывается сообщение. Раздел выбирается на основе хэша ключа, поэтому сообщения с одинаковым ключом всегда записываются в один и тот же раздел. Каждый раздел реплицируется на нескольких брокерах для того, чтобы обеспечить высокую доступность и долговечность данных.
Консюмеры могут подписаться на один или несколько топиков и читать сообщения из одного или нескольких разделов. Когда консюмер читает сообщения из раздела, он сохраняет свое собственное смещение, которое является позицией последнего использованного сообщения в разделе. Это позволяет консюмерам читать сообщения в своем собственном темпе и в случае сбоев возобновлять чтение с того места, на котором они остановились. Архитектура Kafka представляет собой масштабируемую, отказоустойчивую и высокопроизводительную платформу для создания конвейеров данных, работающих в режиме реального времени, и потоковых приложений. Ее распределенный характер позволяет оперативно обрабатывать большие объемы данных и обеспечивать доступ к ним с низкой задержкой в различных приложениях и системах.
Разделы, Фактор репликации и период удержания
Разделы, коэффициент репликации и период удержания - важные понятия, которые определяют то, как данные хранятся, реплицируются и сохраняются в кластере Kafka.
- Разделы: Топик Kafka делится на один или несколько разделов, которые по сути являются упорядоченными и неизменяемыми последовательностями записей. Каждый раздел может храниться на разных брокерах Kafka, что позволяет Kafka распределять нагрузку по чтению и записи данных между несколькими узлами. Это позволяет Kafka масштабироваться горизонтально и оперативно обрабатывать большие объемы данных.
- Фактор репликации: Для обеспечения высокой доступности и отказоустойчивости Kafka реплицирует каждый раздел между несколькими брокерами. Количество реплик определяется коэффициентом репликации, который задает количество копий каждого раздела, которые необходимо поддерживать. Например, если коэффициент репликации равен 3, то каждый раздел будет реплицирован три раза через три разных брокера. Это означает, что даже если один или два брокера выйдут из строя, данные все равно будут доступны.
- Период удержания: Топики Kafka могут быть настроены с учетом периода удержания или периода хранения, который определяет, как долго Kafka должен хранить сообщения в топике перед их удалением. Это может быть полезно для отслеживания последних данных или для реализации политики хранения данных. Kafka поддерживает два типа периода удержания: по времени и по размеру. Политики, основанные на времени, удаляют сообщения, которые находились в топике в течение определенного времени, а политики, основанные на размере, удаляют сообщения на основе общего объема данных, хранящихся в топике.
Настраивая разделы, коэффициент репликации и период удержания данных, пользователи Kafka могут создавать топики, отвечающие высоким требованиям к масштабируемости, отказоустойчивости и хранению данных.
Начало работы с Apache Kafka
Вот несколько шагов, которые помогут Вам быстро начать работу с Apache Kafka:
1. Скачайте и установите Apache Kafka: Вы можете скачать последнюю версию Apache Kafka с официального сайта. После загрузки следуйте инструкциям по установке.
2. Запустите сервер Kafka: После установки Kafka Вы сможете запустить сервер Kafka, выполнив следующую команду:
bin/kafka-server-start.sh config/server.properties
Данная команда запустит сервер Kafka на Вашем локальном устройстве.
3. Создайте топик: Топик - это название категории или ленты, в которую публикуются сообщения. Вы можете создать свой собственный топик, выполнив следующую команду:
bin/kafka-topics.sh --create --topic my-topic --bootstrap-server localhost:9092
Эта команда создаст топик под названием my-topic на сервере Kafka, работающем с портом 9092.
4. Отправьте сообщения в опик: Вы можете отправлять сообщения в топик с помощью Kafka Producer API. Пример:
import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.clients.producer.KafkaProducer;
import java.util.Properties;
public class KafkaProducerExample {
public static void main(String[] args) {
Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092");
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
Producer<String, String> producer = new KafkaProducer<>(props);
String message = "Hello, Kafka!";
ProducerRecord<String, String> record = new ProducerRecord<>("my-topic", message);
producer.send(record);
producer.close();
}
}
В этом примере мы создаем продюсера Kafka, который отправляет сообщение «Hello, Kafka!» в топик «my-topic».
5. Возьмите из топика сообщение и прочитайте его: Вы можете потреблять сообщения из топика, используя API консюмера Kafka. Пример:
import org.apache.kafka.clients.consumer.Consumer;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.common.serialization.StringDeserializer;
import java.util.Collections;
import java.util.Properties;
public class KafkaConsumerExample {
public static void main(String[] args) {
Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092");
props.put("key.deserializer", StringDeserializer.class.getName());
props.put("value.deserializer", StringDeserializer.class.getName());
props.put("group.id", "my-group");
Consumer<String, String> consumer = new KafkaConsumer<>(props);
consumer.subscribe(Collections.singleton("my-topic"));
while (true) {
ConsumerRecords<String, String> records = consumer.poll(1000);
records.forEach(record -> {
System.out.println("Received message: " + record.value());
});
}
}
}
В этои примере мы создали консмера, которые читает сообщения из топика “my-topic” и выводит их на специальную консоль.
Пример:
Вот пример конечной точки Kafka Spring Boot, которая передает событие в топик Kafka, как только а в базу данных добавляется новая запись о сотруднике:
Во-первых, Вам нужно добавить зависимости Kafka и Spring Data JPA в Ваш проект Spring Boot. Вы можете добавить их в файл pom.xml, если Вы используете Maven, или в файл build.gradle, если Вы используете Gradle
<!-- Kafka dependency -->
<dependency>
<groupId>org.springframework.kafka</groupId>
<artifactId>spring-kafka</artifactId>
<version>2.7.4</version>
</dependency>
<!-- Spring Data JPA dependency -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-data-jpa</artifactId>
</dependency>
Создайте класс конфигурации Kafka, который определит продюсера.
@Configuration
@EnableKafka
public class KafkaConfig {
@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());
}
}
Создайте службу Kafka producer, которая будет отправлять сообщения в топик Kafka.
@Service
public class KafkaProducerService {
private final KafkaTemplate<String, String> kafkaTemplate;
@Autowired
public KafkaProducerService(KafkaTemplate<String, String> kafkaTemplate) {
this.kafkaTemplate = kafkaTemplate;
}
public void sendMessage(String message) {
kafkaTemplate.send("employee-topic", message);
}
}
Создайте экземпляр Spring Data JPA для записи о сотруднике.
@Entity
public class Employee {
@Id
@GeneratedValue(strategy = GenerationType.IDENTITY)
private Long id;
private String name;
private String department;
// getters and setters
}
Создайте JPA-репозиторий Spring Data для сотрудника.
@Repository
public interface EmployeeRepository extends JpaRepository<Employee, Long> {
}
Создайте контроллер Spring MVC, который будет обрабатывать HTTP POST-запрос для добавления новой записи о сотруднике в базу данных
@RestController
public class EmployeeController {
private final EmployeeRepository employeeRepository;
private final KafkaProducerService kafkaProducerService;
@Autowired
public EmployeeController(EmployeeRepository employeeRepository, KafkaProducerService kafkaProducerService) {
this.employeeRepository = employeeRepository;
this.kafkaProducerService = kafkaProducerService;
}
@PostMapping("/employees")
public ResponseEntity<Employee> createEmployee(@RequestBody Employee employee) {
Employee savedEmployee = employeeRepository.save(employee);
kafkaProducerService.sendMessage("New employee added: " + savedEmployee.getName());
return ResponseEntity.ok(savedEmployee);
}
}
В приведенном выше примере, когда новая запись сотрудника добавляется в базу данных с помощью HTTP POST-запроса, метод createEmployee отправляет сообщение в топик Kafka с помощью kafkaProducerService. Сообщение содержит имя нового добавленного сотрудника. Вы можете настроить формат сообщения и имя топика в соответствии с Вашими требованиями.
Чтобы отправить только что добавленную запись о сотруднике в топик employee-topic, измените KafkaTemplate
@Service
public class KafkaProducerService {
private static final String TOPIC_NAME = "employee-topic";
private static final int NUM_PARTITIONS = 3;
private static final short REPLICATION_FACTOR = 1;
private static final long RETENTION_PERIOD_MS = 86400000L; // 1 day
private final KafkaTemplate<Long, Employee> kafkaTemplate;
@Autowired
public KafkaProducerService(KafkaTemplate<Long, Employee> kafkaTemplate) {
this.kafkaTemplate = kafkaTemplate;
}
public void sendEmployee(Employee employee) {
kafkaTemplate.executeInTransaction(operations -> {
operations.send(TOPIC_NAME, employee.getId(), employee);
return null;
});
AdminClient adminClient = AdminClient.create(kafkaTemplate.getProducerFactory().getConfigurationProperties());
NewTopic newTopic = new NewTopic(TOPIC_NAME, NUM_PARTITIONS, REPLICATION_FACTOR);
Map<String, String> config = new HashMap<>();
config.put(TopicConfig.RETENTION_MS_CONFIG, String.valueOf(RETENTION_PERIOD_MS));
newTopic.configs(config);
try {
adminClient.createTopics(Collections.singleton(newTopic)).all().get();
} catch (InterruptedException | ExecutionException e) {
throw new RuntimeException(e);
} finally {
adminClient.close();
}
}
}
В этой измененной версии KafkaProducerService для отправки записи Employee в топик Kafka «employee-topic» мы используем метод executeInTransaction , который гарантирует то, что продюсер Kafka отправит сообщение атомарно, то есть сообщение будет либо успешно отправлено, либо не отправлено вовсе.
Мы также используем класс AdminClient из клиентской библиотеки Kafka для настройки свойств топика Kafka. В этом примере мы устанавливаем количество разделов на 3, коэффициент репликации на 1 и период удержания на 1 день. Вы можете настроить эти свойства в соответствии с Вашими требованиями. Обратите внимание, что класс AdminClient должен быть закрыт после использования для того, чтобы освободить ресурсы, используемые клиентом Kafka.
Консюмер Kafka
Пример консюмера Kafka в Spring Boot, который читает сообщения из топика «employee-topic».:
@Service
public class KafkaConsumerService {
private static final String TOPIC_NAME = "employee-topic";
@KafkaListener(topics = TOPIC_NAME, groupId = "employee-consumer-group")
public void receiveEmployee(Employee employee) {
// do something with the received employee record
System.out.println("Received employee: " + employee);
}
}
В этом примере мы используем аннотацию @KafkaListener из библиотеки Spring Kafka для настройки метода, который будет вызываться каждый раз, когда будет получено новое сообщение из топика «employee-topic». Атрибут groupId указывает уникальный идентификатор группы консюмеров, к которой принадлежит данный потребитель. Обратите внимание, что несколько консюмеров могут принадлежать к одной группе, и Kafka будет автоматически распределять разделы топиками между ее членами.
При получении нового сообщения будет вызван метод receiveEmployee с объектом Employee, который был десериализован из сообщения Kafka. Вы можете модифицировать этот метод для выполнения любой необходимой бизнес-логики с полученными данными.
Чтобы включить консюмера Kafka в Ваше приложение Spring Boot, Вам нужно настроить свойства Kafka в файле application.properties или application.yml. Пример конфигурации:
spring.kafka.consumer.bootstrap-servers=<kafka-broker-url> spring.kafka.consumer.group-id=employee-consumer-group spring.kafka.consumer.auto-offset-reset=earliest
В этой конфигурации мы устанавливаем свойство bootstrap servers на URL брокера (брокеров) Kafka, к которым должен подключаться консюмер, а также свойство group ID на то же значение, которое мы использовали в аннотации @KafkaListener. Свойство auto-offset-reset определяет поведение консюмера, когда он запускается и не имеет действительного смещения для раздела. В этом примере мы установим значение «earliest» для того, чтобы консюмер читал все сообщения в топике с самого начала.




