Мини-проект: 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.
Смотри также
- Hex Architecture - Use Cases - EventPublisher как порт
- Event-Driven - Outbox + Inbox - теория паттерна
- Message Brokers - RabbitMQ Producer - детали AMQP publisher
Мини-задание
- Реализуй структуру проекта из урока:
cmd/api,cmd/email-worker,cmd/analytics - Напиши
CompleteLessonuse case с транзакционным outbox для RabbitMQ - Реализуй
KafkaPublisherсAllISRAcksиFlushна shutdown - Подними весь docker-compose, отправь 10 событий, проверь что в RabbitMQ и Kafka всё дошло
- Симулируй сбой: останови email-worker, отправь событие, перезапусти worker - сообщение должно быть в очереди
- Добавь idempotency в analytics consumer: пропускать дубли по
eventId