📨 Очереди сообщений: Погружение в RabbitMQ и Kafka для обработки событий
Синхронные HTTP-вызовы связывают микросервисы в жесткие цепочки. В пиковые нагрузки достаточно сбоя в одном узле (например, в сервисе оплаты), чтобы по цепочке таймаутов упало всё приложение.
Асинхронная коммуникация через брокеры сообщений разрывает эту зависимость: сервис заказов быстро принимает запрос клиента, а фоновые обработчики забирают задачи из очереди по мере готовности.
🐰 RabbitMQ: гибкая маршрутизация
RabbitMQ — классический брокер с развитой маршрутизацией (`direct`, topic, fanout, `headers`) и поддержкой DLQ (dead letter queues). Работает по модели доставки «как минимум один раз» (*at-least-once*).
Пример продюсера (Go, `PublishWithContext` + publisher confirms):
package main
import (
"context"
"log"
"time"
"github.com/rabbitmq/amqp091-go"
)
func main() {
// Подключаемся к RabbitMQ
conn, err := amqp091.Dial("amqp://guest:guest@localhost:5672/")
if err != nil {
log.Fatal(err)
}
defer conn.Close()
ch, err := conn.Channel()
if err != nil {
log.Fatal(err)
}
defer ch.Close()
// Для критичных данных в продакшене используйте реплицируемую quorum queue.
// Этот сокращённый пример объявляет обычную durable-очередь.
q, err := ch.QueueDeclare(
"order_queue", // name
true, // durable
false, // delete when unused
false, // exclusive
false, // no-wait
nil, // arguments
)
if err != nil {
log.Fatal(err)
}
body := `{"order_id": 12345, "user_id": 678, "amount": 5000}`
// Создаём контекст для отправки
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
// Включаем publisher confirms: persistent-флага и durable-очереди недостаточно,
// чтобы продюсер узнал, принял ли брокер сообщение.
if err := ch.Confirm(false); err != nil {
log.Fatal(err)
}
confirmation, err := ch.PublishWithDeferredConfirmWithContext(ctx,
"", // exchange
q.Name, // routing key
false, // маршрут - заранее объявленная очередь
false, // immediate
amqp091.Publishing{
ContentType: "application/json",
Body: []byte(body),
DeliveryMode: amqp091.Persistent,
})
if err != nil {
log.Fatal(err)
}
acked, err := confirmation.WaitContext(ctx)
if err != nil {
log.Fatal(err)
}
if !acked {
log.Fatal("broker rejected the message")
}
log.Printf("[x] Confirmed: %s", body)
}
Пример консьюмера (Go, ручное подтверждение):
package main
import (
"log"
"github.com/rabbitmq/amqp091-go"
)
func main() {
conn, err := amqp091.Dial("amqp://guest:guest@localhost:5672/")
if err != nil {
log.Fatal(err)
}
defer conn.Close()
ch, err := conn.Channel()
if err != nil {
log.Fatal(err)
}
defer ch.Close()
q, err := ch.QueueDeclare(
"order_queue",
true,
false,
false,
false,
nil,
)
if err != nil {
log.Fatal(err)
}
msgs, err := ch.Consume(
q.Name,
"", // consumer tag
false, // auto-ack (false - включаем ручное подтверждение)
false, // exclusive
false, // no-local
false, // no-wait
nil, // args
)
if err != nil {
log.Fatal(err)
}
forever := make(chan struct{})
go func() {
for d := range msgs {
log.Printf("[x] Received: %s", d.Body)
// Обработка сообщения
if err := processOrder(d.Body); err == nil {
d.Ack(false) // Подтверждаем успешную обработку
} else {
log.Printf("[!] Error processing order: %v", err)
// Не делайте бесконечный немедленный requeue: это создаёт горячий цикл.
// В продакшене направляйте сообщение в retry-очередь с задержкой,
// лимитом попыток и последующим переводом в DLQ.
d.Nack(false, false)
}
}
}()
log.Printf("[*] Waiting for messages. To exit press CTRL+C")
<-forever
}
// Заглушка для демонстрации бизнес-логики обработки заказа
func processOrder(body []byte) error {
return nil
}