Kafka — Producer + Consumer
Цель урока
Сделать полноценный CRUD для задач с правильными HTTP-статусами (201/204/404), обработкой ошибок, и подготовить Task к подключению JPA.
Теория · для собеса
1 Что такое Kafka ▸
Apache Kafka — распределённый event streaming platform. Используется как message broker для event-driven архитектур.
2 Ключевые понятия ▸
- Broker — Kafka-сервер (обычно кластер из 3+)
- Topic — категория событий (
task-events,user-events) - Partition — топик делится на партиции (для параллелизма)
- Offset — позиция сообщения в партиции (монотонно растёт)
- Producer — пишет в топик
- Consumer — читает из топика (в рамках consumer group)
- Consumer Group — группа инстансов, читающих один топик (каждое сообщение — одному инстансу)
3 Когда использовать Kafka ▸
- Event-driven microservices — «при создании заказа → отправь в Kafka, платежный сервис слушает»
- Audit log — все события пишутся в Kafka, потом в data lake
- Async processing — producer не ждёт consumer (decoupling)
- Streaming — Kafka Streams / ksqlDB для real-time аналитики
⚠️ Не для: простых RPC (используй REST/gRPC), долговременного хранения (Kafka — не БД, retention по умолчанию 7 дней).
---
Практика: улучшаем CRUD
-
1Зависимости
dependencies { implementation("org.springframework.kafka:spring-kafka") } -
2Docker Compose с Kafka
# docker-compose.yml version: '3.8' services: postgres: image: postgres:16-alpine environment: POSTGRES_DB: taskflow_dev POSTGRES_USER: taskflow POSTGRES_PASSWORD: dev_password ports: ["5432:5432"] volumes: [pgdata:/var/lib/postgresql/data] kafka: image: bitnami/kafka:3.7 environment: KAFKA_CFG_NODE_ID: 0 KAFKA_CFG_PROCESS_ROLES: controller,broker KAFKA_CFG_LISTENERS: PLAINTEXT://:9092,CONTROLLER://:9093 KAFKA_CFG_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092 KAFKA_CFG_CONTROLLER_LISTENER_NAMES: CONTROLLER KAFKA_CFG_CONTROLLER_QUORUM_VOTERS: 0@kafka:9093 ALLOW_PLAINTEXT_LISTENER: "yes" ports: ["9092:9092"] volumes: [kafkadata:/bitnami/kafka/data] volumes: pgdata: {} kafkadata: {}docker compose up -d kafka -
3application.yml
spring: kafka: bootstrap-servers: localhost:9092 producer: key-serializer: org.apache.kafka.common.serialization.StringSerializer value-serializer: org.springframework.kafka.support.serializer.JsonSerializer acks: all # ждём подтверждения от всех реплик properties: enable.idempotence: true # защита от дублей consumer: group-id: taskflow-consumer key-deserializer: org.apache.kafka.common.serialization.StringDeserializer value-deserializer: org.springframework.kafka.support.serializer.JsonDeserializer auto-offset-reset: earliest properties: spring.json.trusted.packages: "com.taskflow.event" -
4
event/TaskEvent.java(базовый класс)package com.taskflow.event; import java.time.Instant; public sealed interface TaskEvent permits TaskCreatedEvent, TaskUpdatedEvent, TaskDeletedEvent { Long taskId(); String eventType(); Instant occurredAt(); }package com.taskflow.event; import java.time.Instant; public record TaskCreatedEvent( Long taskId, String title, String ownerUsername, Instant occurredAt ) implements TaskEvent { @Override public String eventType() { return "TaskCreated"; } }package com.taskflow.event; import java.time.Instant; public record TaskUpdatedEvent( Long taskId, String title, boolean done, Instant occurredAt ) implements TaskEvent { @Override public String eventType() { return "TaskUpdated"; } }package com.taskflow.event; import java.time.Instant; public record TaskDeletedEvent( Long taskId, Instant occurredAt ) implements TaskEvent { @Override public String eventType() { return "TaskDeleted"; } } -
5
config/KafkaTopicConfig.javapackage com.taskflow.config; import org.apache.kafka.clients.admin.NewTopic; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; @Configuration public class KafkaTopicConfig { public static final String TASK_EVENTS_TOPIC = "task-events"; @Bean public NewTopic taskEventsTopic() { return new NewTopic(TASK_EVENTS_TOPIC, 3, (short) 1); // 3 партиции } } -
6
producer/TaskEventProducer.javapackage com.taskflow.producer; import com.taskflow.config.KafkaTopicConfig; import com.taskflow.event.TaskEvent; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.kafka.core.KafkaTemplate; import org.springframework.stereotype.Component; @Component public class TaskEventProducer { private static final Logger log = LoggerFactory.getLogger(TaskEventProducer.class); private final KafkaTemplatekafkaTemplate; public TaskEventProducer(KafkaTemplate kafkaTemplate) { this.kafkaTemplate = kafkaTemplate; } public void publish(TaskEvent event) { // Ключ = taskId → все события одной задачи идут в одну партицию (порядок!) String key = String.valueOf(event.taskId()); kafkaTemplate.send(KafkaTopicConfig.TASK_EVENTS_TOPIC, key, event) .whenComplete((result, ex) -> { if (ex != null) { log.error("Failed to publish event: {}", event, ex); } else { log.info("Published event: topic={}, partition={}, offset={}, event={}", result.getRecordMetadata().topic(), result.getRecordMetadata().partition(), result.getRecordMetadata().offset(), event); } }); } } -
7
consumer/TaskEventConsumer.javapackage com.taskflow.consumer; import com.taskflow.event.TaskCreatedEvent; import com.taskflow.event.TaskDeletedEvent; import com.taskflow.event.TaskEvent; import com.taskflow.event.TaskUpdatedEvent; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.kafka.annotation.KafkaListener; import org.springframework.stereotype.Component; @Component public class TaskEventConsumer { private static final Logger log = LoggerFactory.getLogger(TaskEventConsumer.class); @KafkaListener( topics = "task-events", groupId = "taskflow-audit-consumer" ) public void onTaskEvent(TaskEvent event) { switch (event) { case TaskCreatedEvent e -> log.info("AUDIT: Task created: id={}, title={}, owner={}", e.taskId(), e.title(), e.ownerUsername()); case TaskUpdatedEvent e -> log.info("AUDIT: Task updated: id={}, title={}, done={}", e.taskId(), e.title(), e.done()); case TaskDeletedEvent e -> log.info("AUDIT: Task deleted: id={}", e.taskId()); } } } -
8Использование в
TaskServicepublic TaskResponse create(CreateTaskRequest req) { User me = currentUser.get(); Task task = mapper.toEntity(req); task.setOwner(me); Task saved = repo.save(task); // Отправляем событие eventProducer.publish(new TaskCreatedEvent( saved.getId(), saved.getTitle(), me.getUsername(), Instant.now() )); return mapper.toResponse(saved); } -
9Тест
# Создаём задачу $ curl -X POST http://localhost:8080/api/tasks \ -H "Authorization: Bearer $TOKEN" \ -H "Content-Type: application/json" \ -d '{"title":"Test Kafka"}' # В логах Spring Boot увидишь: Producer: Published event: topic=task-events, partition=0, offset=0, event=TaskCreatedEvent[...] Consumer: AUDIT: Task created: id=1, title=Test Kafka, owner=alice---
Зачем это на собесе
После урока ты должен уметь ответить на:
task-events.enable.idempotence=true + транзакционный producer защищает от дублей.5 вопросов на углубление
Раскрой вопрос и нажми «🤔 Хочу разобрать подробнее» — он попадёт в страницу ответов.
Разбор внутри: key hashing, partitions и порядок событий по taskId.
Разбор внутри: надёжная отправка producer-а, retries и delivery timeout.
Разбор внутри: event streaming, очереди задач и replay.
Разбор внутри: earliest/latest и первый запуск нового consumer group.
Разбор внутри: простые сценарии, лишняя инфраструктура и альтернативы.
Готов идти дальше?
Выбери вопросы, которые тебе интересны, и изучи их. Потом — к следующему уроку.
🚀 Перейти к Уроку 16Сначала пройди все секции и выбери хотя бы 1 вопрос
⬅️ Назад к Уроку 14