УРОК 15 / 30 0%

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. 1
    Зависимости
    
    dependencies {
        implementation("org.springframework.kafka:spring-kafka")
    }
    
  2. 2
    Docker 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
    
  3. 3
    application.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. 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. 5
    config/KafkaTopicConfig.java
    
    package 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. 6
    producer/TaskEventProducer.java
    
    package 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 KafkaTemplate kafkaTemplate;
    
        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. 7
    consumer/TaskEventConsumer.java
    
    package 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. 8
    Использование в TaskService
    
    public 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. 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
    

    ---

🎯

Зачем это на собесе

После урока ты должен уметь ответить на:

Decoupling: TaskFlow не знает, кто слушает события. Завтра добавится notification-service, email-service, analytics — они просто подпишутся на task-events.
Audit log: все изменения пишутся в Kafka → оттуда в data lake → compliance.
Replay: можно перечитать все события с offset=0 (например, для восстановления state).
Exactly-once: enable.idempotence=true + транзакционный producer защищает от дублей.

5 вопросов на углубление

Раскрой вопрос и нажми «🤔 Хочу разобрать подробнее» — он попадёт в страницу ответов.

1
Почему ключ Kafka важен для порядка сообщений?

Разбор внутри: key hashing, partitions и порядок событий по taskId.

2
Что дают acks=all и enable.idempotence=true?

Разбор внутри: надёжная отправка producer-а, retries и delivery timeout.

3
Kafka vs RabbitMQ — когда что использовать?

Разбор внутри: event streaming, очереди задач и replay.

4
Что делает auto-offset-reset?

Разбор внутри: earliest/latest и первый запуск нового consumer group.

5
Когда Kafka не нужна?

Разбор внутри: простые сценарии, лишняя инфраструктура и альтернативы.

Готов идти дальше?

Выбери вопросы, которые тебе интересны, и изучи их. Потом — к следующему уроку.

🚀 Перейти к Уроку 16

Сначала пройди все секции и выбери хотя бы 1 вопрос

⬅️ Назад к Уроку 14
🎉
Новый тир
Новый тир достигнут!