Мини-проект: LessonCompleted в RabbitMQ и Kafka

Мини-проект: LessonCompleted в RabbitMQ и Kafka

Финальный проект трека. Берём одно событие lesson.completed и доставляем его в два брокера одновременно: в RabbitMQ идёт задача на email-нотификацию, в Kafka - событие для аналитики.

Архитектура

HTTP POST /lessons/{id}/complete
    │
    ▼
CompleteLesson (use case)
    │
    ├── ProgressRepository.MarkCompleted(tx)
    ├── OutboxRepository.Save(tx, RabbitMQ task)  ─────► RabbitMQ
    │                                                         │
    │                                                         ▼
    │                                                    email-worker
    │
    └── KafkaPublisher.Publish(ctx, event)  ──────────► Kafka
                                                            │
                                                            ▼
                                                       analytics-consumer

Два брокера - два разных паттерна:

  • RabbitMQ через Outbox: атомарно с транзакцией БД, гарантия доставки email
  • Kafka напрямую: at-least-once, для аналитики допустим редкий дубль

Структура проекта

lesson-completed/
├── cmd/
│   ├── api/main.go              - HTTP API
│   ├── email-worker/main.go     - RabbitMQ consumer
│   └── analytics/main.go        - Kafka consumer group
├── domain/
│   └── events.go                - LessonCompletedEvent
├── usecases/
│   └── complete_lesson.go       - бизнес-логика
├── adapters/
│   ├── rabbitmq_publisher.go    - Outbox Poller → RabbitMQ
│   ├── kafka_publisher.go       - franz-go producer
│   └── postgres_repo.go         - outbox + progress repo
├── docker-compose.yml
└── Makefile

domain/events.go

package domain

import "time"

type LessonCompletedEvent struct {
    EventID  string    `json:"eventId"`
    UserID   string    `json:"userId"`
    LessonID string    `json:"lessonId"`
    At       time.Time `json:"at"`
}

type EventPublisher interface {
    PublishLessonCompleted(ctx context.Context, e LessonCompletedEvent) error
}

usecases/complete_lesson.go

package usecases

type CompleteLessonUseCase struct {
    progress  ProgressRepository
    outbox    OutboxRepository
    kafka     EventPublisher
    db        *sql.DB
}

func (uc *CompleteLessonUseCase) Execute(ctx context.Context, userID, lessonID string) error {
    event := domain.LessonCompletedEvent{
        EventID:  uuid.New().String(),
        UserID:   userID,
        LessonID: lessonID,
        At:       time.Now().UTC(),
    }

    // 1. Transactionally save progress + RabbitMQ task to outbox
    err := withTx(ctx, uc.db, func(tx *sql.Tx) error {
        if err := uc.progress.MarkCompleted(ctx, tx, userID, lessonID); err != nil {
            return err
        }
        payload, _ := json.Marshal(event)
        return uc.outbox.Save(ctx, tx, OutboxEntry{
            Topic:        "notifications",
            RoutingKey:   "lesson.completed",
            PartitionKey: userID,
            Payload:      payload,
            EventID:      event.EventID,
        })
    })
    if err != nil {
        return err
    }

    // 2. Publish to Kafka (at-least-once, not in transaction)
    return uc.kafka.PublishLessonCompleted(ctx, event)
}

adapters/kafka_publisher.go

package adapters

type KafkaPublisher struct {
    client *kgo.Client
    topic  string
}

func NewKafkaPublisher(brokers []string, topic string) (*KafkaPublisher, error) {
    client, err := kgo.NewClient(
        kgo.SeedBrokers(brokers...),
        kgo.RequiredAcks(kgo.AllISRAcks()),
        kgo.ProducerLinger(5*time.Millisecond),
    )
    if err != nil {
        return nil, err
    }
    return &KafkaPublisher{client: client, topic: topic}, nil
}

func (p *KafkaPublisher) PublishLessonCompleted(ctx context.Context, e domain.LessonCompletedEvent) error {
    payload, _ := json.Marshal(e)
    r := &kgo.Record{
        Topic: p.topic,
        Key:   []byte(e.UserID),
        Value: payload,
        Headers: []kgo.RecordHeader{
            {Key: "event-id", Value: []byte(e.EventID)},
        },
    }
    return p.client.ProduceSync(ctx, r).FirstErr()
}

cmd/email-worker/main.go

func main() {
    ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
    defer stop()

    conn, _ := amqp091.Dial(os.Getenv("RABBITMQ_URL"))
    defer conn.Close()
    ch, _ := conn.Channel()

    ch.QueueDeclare("email-notifications", true, false, false, false, nil)
    ch.QueueBind("email-notifications", "lesson.completed", "lesson-events", false, nil)
    ch.Qos(10, 0, false) // prefetch

    msgs, _ := ch.Consume("email-notifications", "", false, false, false, false, nil)

    for {
        select {
        case <-ctx.Done():
            return
        case msg, ok := <-msgs:
            if !ok {
                return
            }
            if err := sendEmail(msg.Body); err != nil {
                msg.Nack(false, true) // requeue
            } else {
                msg.Ack(false)
            }
        }
    }
}

cmd/analytics/main.go

func main() {
    ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
    defer stop()

    client, _ := kgo.NewClient(
        kgo.SeedBrokers(os.Getenv("KAFKA_BROKERS")),
        kgo.ConsumerGroup("analytics"),
        kgo.ConsumeTopics("lesson-events"),
        kgo.DisableAutoCommit(),
    )
    defer client.Close()

    for {
        fetches := client.PollFetches(ctx)
        if ctx.Err() != nil {
            client.CommitMarkedOffsets(context.Background())
            return
        }

        var toCommit []*kgo.Record
        fetches.EachRecord(func(r *kgo.Record) {
            if err := recordAnalytics(r); err != nil {
                log.Printf("analytics error: %v", err)
                return
            }
            toCommit = append(toCommit, r)
        })

        if len(toCommit) > 0 {
            client.CommitRecords(ctx, toCommit...)
        }
    }
}

docker-compose.yml

version: "3.9"
services:
  postgres:
    image: postgres:16-alpine
    environment:
      POSTGRES_DB: app
      POSTGRES_USER: app
      POSTGRES_PASSWORD: secret
    ports: ["5432:5432"]

  rabbitmq:
    image: rabbitmq:3-management
    ports: ["5672:5672", "15672:15672"]
    environment:
      RABBITMQ_DEFAULT_USER: guest
      RABBITMQ_DEFAULT_PASS: guest

  kafka:
    image: confluentinc/cp-kafka:7.6.0
    ports: ["9092:9092"]
    environment:
      KAFKA_NODE_ID: 1
      KAFKA_PROCESS_ROLES: broker,controller
      KAFKA_LISTENERS: PLAINTEXT://0.0.0.0:9092,CONTROLLER://0.0.0.0:9093
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092
      KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER
      KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,CONTROLLER:PLAINTEXT
      KAFKA_CONTROLLER_QUORUM_VOTERS: 1@kafka:9093
      KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
      CLUSTER_ID: "MkU3OEVBNTcwNTJENDM2Qk"

  api:
    build: ./cmd/api
    ports: ["8080:8080"]
    environment:
      DB_URL: postgres://app:secret@postgres:5432/app
      RABBITMQ_URL: amqp://guest:guest@rabbitmq:5672/
      KAFKA_BROKERS: kafka:9092
    depends_on: [postgres, rabbitmq, kafka]

  email-worker:
    build: ./cmd/email-worker
    environment:
      RABBITMQ_URL: amqp://guest:guest@rabbitmq:5672/
    depends_on: [rabbitmq]

  analytics:
    build: ./cmd/analytics
    environment:
      KAFKA_BROKERS: kafka:9092
    depends_on: [kafka]

Проверка

# Запуск
docker compose up -d

# Отправить событие
curl -X POST http://localhost:8080/lessons/go-01/complete \
  -H "X-User-ID: user-42"

# Проверить RabbitMQ
open http://localhost:15672  # guest/guest
# → Queues → email-notifications → Messages

# Проверить Kafka consumer group
docker exec kafka kafka-consumer-groups \
  --bootstrap-server localhost:9092 \
  --describe --group analytics

Типичные ошибки

Не закрывать Kafka producer при shutdown - вызывай client.Flush(ctx) перед client.Close() в HTTP-сервере, иначе последние события теряются.

Публиковать в Kafka до фиксации транзакции - если транзакция откатится после Kafka-публикации, event ушёл без бизнес-данных. В этом проекте Kafka публикуется после успешного commit (строка return uc.kafka.PublishLessonCompleted идёт после withTx).

Не реализовывать idempotent consumer в analytics - Kafka гарантирует at-least-once, дубли возможны. Analytics consumer должен уметь пропускать уже обработанные eventId.

Смотри также

Мини-задание

  • Реализуй структуру проекта из урока: cmd/api, cmd/email-worker, cmd/analytics
  • Напиши CompleteLesson use case с транзакционным outbox для RabbitMQ
  • Реализуй KafkaPublisher с AllISRAcks и Flush на shutdown
  • Подними весь docker-compose, отправь 10 событий, проверь что в RabbitMQ и Kafka всё дошло
  • Симулируй сбой: останови email-worker, отправь событие, перезапусти worker - сообщение должно быть в очереди
  • Добавь idempotency в analytics consumer: пропускать дубли по eventId

Зарегистрируйтесь бесплатно, чтобы пройти квиз, решить задание с автопроверкой и вести прогресс.