Когда сервисы начинают общаться через брокер, возникает вопрос: как правильно организовать очереди, обменники и подписки? Большинство задач укладываются в несколько типовых схем. Разберём каждую на примерах Go и rabbitmq/amqp091-go.
Раздать задачи нескольким обработчикам — Work Queue
Представьте: пользователи загружают фотографии, и каждую надо сжать и нарезать в несколько размеров. Это долго. Делать это прямо в HTTP-запросе нельзя — пользователь будет ждать минуты.
Решение — work queue (очередь задач): сохранить задание в очередь и отдать её одному из пула обработчиков.
Очередь одна, потребителей несколько: брокер раздаёт сообщения по одному, и каждое достаётся ровно кому-то одному. Больше потребителей — быстрее разбирается очередь, но общий порядок между ними уже не сохраняется.
Брокер сам распределяет задачи между обработчиками. Каждое сообщение получает ровно один из них.
ch.QueueDeclare(
"images.to-process",
true, // durable
false, // autoDelete
false, // exclusive
false, // noWait
amqp.Table{"x-queue-type": "quorum"},
)
for i := 0; i < 20; i++ {
go worker(conn)
}
func worker(conn *amqp.Connection) {
ch, err := conn.Channel() // свой канал на воркера: каналы не делят между горутинами
if err != nil {
return
}
defer ch.Close()
ch.Qos(1, 0, false) // одно неподтверждённое сообщение на воркера
msgs, err := ch.Consume("images.to-process", "", false /* autoAck */, false, false, false, nil)
if err != nil {
return
}
for d := range msgs {
if err := processImage(d.Body); err != nil {
d.Nack(false, true) // вернуть в очередь; про отравленные сообщения — ниже, в retry и DLQ
continue
}
d.Ack(false)
}
}
Каждый воркер держит собственный канал: документация amqp091-go прямо просит не делить *amqp.Channel между горутинами, а каналы дёшевы — их открывают по одному на воркера поверх одного соединения. Qos(1) на канале значит, что брокер не выдаст воркеру следующее сообщение, пока тот не подтвердил текущее, — так работа достаётся свободному воркеру, а не первому попавшемуся. Если запустить 5 копий сервиса — получится 100 параллельных обработчиков.
Когда брать: фоновые задачи — обработка файлов, отправка писем, генерация отчётов, любые «положили в очередь — кто-то возьмёт».
Отправить событие всем сервисам сразу — Publish/Subscribe
Другая задача: произошло событие «конфигурация обновлена» — каждый сервис должен обновить свой кеш. Нельзя знать заранее, кто именно подписан и сколько сервисов запущено.
Это publish/subscribe (рассылка всем): издатель отправляет одно сообщение, а копию получают все подписчики одновременно.
Здесь нужен fanout exchange — он копирует каждое сообщение во все привязанные очереди. Каждый сервис объявляет свою очередь и привязывает её к общему обменнику.
ch.ExchangeDeclare("cache.invalidation", "fanout", true, false, false, false, nil)
// своя очередь сервиса: без имени, exclusive, auto-delete
q, err := ch.QueueDeclare("", false, true, true, false, nil)
ch.QueueBind(q.Name, "", "cache.invalidation", false, nil)
msgs, err := ch.Consume(q.Name, "", true, true, false, false, nil)
for d := range msgs {
var event CacheInvalidationEvent
if err := json.Unmarshal(d.Body, &event); err != nil {
continue
}
cache.Evict(event.Key)
}
exclusive + autoDelete — очередь принадлежит одному соединению и удаляется при отключении. При перезапуске сервиса не накапливается мусор в брокере.
Когда брать: инвалидация кешей, broadcast-уведомления всему кластеру, обновления конфигурации.
Направить событие в нужный обработчик — Routing
Иногда нужно не «всем», а «именно тому, кому надо». Например: событие order.created должно идти в сервис выполнения и в аудит, а order.payment-failed — только в алерты.
Это routing (точечная маршрутизация): direct exchange смотрит на ключ маршрутизации сообщения и отправляет его только в очереди с совпадающим binding key.
ch.ExchangeDeclare("orders", "direct", true, false, false, false, nil)
quorum := amqp.Table{"x-queue-type": "quorum"}
ch.QueueDeclare("orders.fulfillment", true, false, false, false, quorum)
ch.QueueDeclare("orders.audit", true, false, false, false, quorum)
ch.QueueDeclare("orders.alerts", true, false, false, false, quorum)
ch.QueueBind("orders.fulfillment", "order.created", "orders", false, nil)
ch.QueueBind("orders.audit", "order.created", "orders", false, nil)
ch.QueueBind("orders.audit", "order.cancelled", "orders", false, nil)
ch.QueueBind("orders.alerts", "order.payment-failed", "orders", false, nil)
order.created→ fulfillment + audit.order.cancelled→ только audit.order.payment-failed→ только alerts.
Когда брать: явное разделение потоков — алерты отдельно от аудита, основной обработчик отдельно от мониторинга.
Подписаться по маске — Topic
Routing хорош для жёстких правил. Но что если сервис хочет подписаться на «все события по заказам»? Или «всё из EU-региона»?
Topic exchange позволяет задавать подписки через шаблоны. Ключи сообщений строятся через точку (order.created.eu), а в подписке можно использовать:
*— ровно одно слово,#— ноль и более слов.
ch.ExchangeDeclare("events", "topic", true, false, false, false, nil)
quorum := amqp.Table{"x-queue-type": "quorum"}
ch.QueueDeclare("audit.orders", true, false, false, false, quorum)
ch.QueueDeclare("dashboard.eu", true, false, false, false, quorum)
ch.QueueDeclare("alerts.critical", true, false, false, false, quorum)
ch.QueueBind("audit.orders", "order.#", "events", false, nil)
ch.QueueBind("dashboard.eu", "*.*.eu", "events", false, nil)
ch.QueueBind("alerts.critical", "payment.failed.#", "events", false, nil)
Сообщение с ключом order.cancelled.eu попадёт в audit.orders (по order.#) и в dashboard.eu (по *.*.eu).
Когда брать: события с иерархической структурой, когда нужно гибко подписываться без переделки топологии при добавлении новых типов событий.
Порядок по ключу при нескольких обработчиках
Work Queue выше раздаёт сообщения нескольким обработчикам, и это ровно то, что нужно, пока сообщения независимы. Как только они связаны (два обновления одного заказа, создание и отмена одной брони), параллельная обработка ломает порядок: второе сообщение попадает к свободному обработчику раньше, чем первый закончил.
В Kafka порядок держит ключ партиции. У RabbitMQ порядок гарантирован только внутри одной очереди с одним потребителем, и из этого растут два приёма.
Single active consumer. Аргумент очереди x-single-active-consumer: true означает: подписчиков может быть много, но сообщения получает один, остальные ждут и мгновенно подхватывают работу, если активный отвалился. Порядок сохранён, отказоустойчивость есть, параллелизма нет. Это правильный выбор, когда очередь и так не перегружена, а порядок важен целиком.
ch.QueueDeclare("bookings", true, false, false, false, amqp.Table{
"x-queue-type": "quorum",
"x-single-active-consumer": true,
})
Consistent hash exchange. Плагин rabbitmq_consistent_hash_exchange добавляет обменник, который раскладывает сообщения по нескольким очередям по хешу ключа маршрутизации или заголовка. Сообщения одного заказа всегда попадают в одну очередь, у каждой очереди свой единственный потребитель, и получается то же, что партиции в Kafka: порядок внутри ключа плюс параллелизм между ключами. Цена ручное управление: очереди создаёте вы, число очередей меняется только с остановкой, а перекос по ключам («один крупный клиент») складывается в одну очередь.
Что не работает: «сделаем один потребитель, но в десять горутин». Они разберут сообщения из общей выдачи и обработают их вперемешку, порядок теряется так же, как при нескольких потребителях. Если нужна и скорость, и порядок, ключ обязан определять, кто обрабатывает; других вариантов нет.
Запрос-ответ через очередь — RPC
Иногда нужен синхронный ответ, но HTTP не подходит: сервис за NAT, нет публичного адреса, или хочется балансировки по пулу обработчиков.
RPC через очередь: клиент отправляет запрос и ждёт ответа. Брокер доставляет запрос одному из обработчиков, тот отвечает в отдельную очередь-ответ. Для сопоставления запроса и ответа используется correlation-id.
В amqp091-go запрос-ответ собирают вручную — через direct reply-to и CorrelationId:
// Клиент: один consumer на reply-to, ответы раздаются по CorrelationId
type PricingClient struct {
ch *amqp.Channel
mu sync.Mutex
pending map[string]chan amqp.Delivery
}
func NewPricingClient(ch *amqp.Channel) (*PricingClient, error) {
replies, err := ch.Consume("amq.rabbitmq.reply-to", "", true, false, false, false, nil)
if err != nil {
return nil, err
}
c := &PricingClient{ch: ch, pending: map[string]chan amqp.Delivery{}}
go func() {
for d := range replies {
c.mu.Lock()
waiter, ok := c.pending[d.CorrelationId]
delete(c.pending, d.CorrelationId)
c.mu.Unlock()
if ok {
waiter <- d
}
}
}()
return c, nil
}
func (c *PricingClient) Quote(ctx context.Context, req QuoteRequest) (PriceQuote, error) {
corrID := uuid.NewString()
waiter := make(chan amqp.Delivery, 1)
body, _ := json.Marshal(req)
c.mu.Lock() // канал один на клиента: регистрация и публикация под одним замком
c.pending[corrID] = waiter
err := c.ch.PublishWithContext(ctx, "pricing.exchange", "pricing.quote", false, false,
amqp.Publishing{CorrelationId: corrID, ReplyTo: "amq.rabbitmq.reply-to", Body: body})
c.mu.Unlock()
if err != nil {
return PriceQuote{}, err
}
select {
case d := <-waiter:
var quote PriceQuote
err := json.Unmarshal(d.Body, "e)
return quote, err
case <-ctx.Done():
c.mu.Lock()
delete(c.pending, corrID)
c.mu.Unlock()
return PriceQuote{}, ctx.Err()
}
}
// Сервер
for d := range requests {
var req QuoteRequest
if err := json.Unmarshal(d.Body, &req); err != nil {
d.Nack(false, false)
continue
}
body, _ := json.Marshal(ComputePriceQuote(req))
ch.PublishWithContext(ctx, "", d.ReplyTo, false, false, // ответ уходит в reply-to
amqp.Publishing{CorrelationId: d.CorrelationId, Body: body})
d.Ack(false)
}
Direct reply-to (amq.rabbitmq.reply-to) — псевдоочередь RabbitMQ: брокер доставляет ответ прямо в соединение клиента, без создания временной очереди. Сервер публикует ответ в d.ReplyTo с тем же CorrelationId. Consumer на reply-to один на канал и живёт всё время работы клиента: открывать его на каждый запрос нельзя — ответы соседних запросов достанутся чужому циклу и пропадут. Поэтому клиент держит карту ожидающих по CorrelationId, а единственная горутина-диспетчер раздаёт ответы.
Когда брать: нужен синхронный вызов, но HTTP не работает (NAT, firewall, нет публичного адреса); нужна балансировка запросов по пулу обработчиков.
Когда не брать: если HTTP/gRPC просто работает — RPC через брокер сложнее в отладке и дороже.
Что делать с повторной доставкой — Idempotent Consumer
AMQP гарантирует at-least-once доставку: одно и то же сообщение может прийти дважды. Это случается, когда брокер не получил подтверждение обработки (например, из-за сетевой проблемы) и переотправляет сообщение.
Обработчик обязан уметь работать с повторами, не ломая бизнес-логику.
Ключ идемпотентности в базе данных
Самый надёжный способ — запоминать уже обработанные сообщения:
func (p *PaymentProcessor) Process(ctx context.Context, event PaymentEvent) error {
tx, err := p.db.BeginTx(ctx, nil)
if err != nil {
return err
}
defer tx.Rollback()
res, err := tx.ExecContext(ctx,
`INSERT INTO processed_events (idempotency_key) VALUES ($1)
ON CONFLICT DO NOTHING`, event.IdempotencyKey)
if err != nil {
return err
}
if n, _ := res.RowsAffected(); n == 0 {
return nil // уже обработано — просто подтверждаем получение
}
if _, err := tx.ExecContext(ctx,
`UPDATE account SET balance = balance - $1 WHERE id = $2`,
event.Amount, event.AccountID); err != nil {
return err
}
return tx.Commit()
}
Таблица processed_events с уникальным индексом на idempotency_key. Если два одинаковых сообщения придут одновременно — база данных поймает дубль через нарушение уникального ограничения.
Проверка состояния объекта
Если событие переводит объект в новое состояние, достаточно проверить текущее:
func (s *OrderService) OnOrderConfirmed(ctx context.Context, event OrderConfirmedEvent) error {
tx, err := s.db.BeginTx(ctx, nil)
if err != nil {
return err
}
defer tx.Rollback()
var status string
if err := tx.QueryRowContext(ctx,
`SELECT status FROM orders WHERE id = $1 FOR UPDATE`,
event.OrderID).Scan(&status); err != nil {
return err
}
if status == "CONFIRMED" {
return nil // уже в нужном состоянии
}
if _, err := tx.ExecContext(ctx,
`UPDATE orders SET status = 'CONFIRMED' WHERE id = $1`,
event.OrderID); err != nil {
return err
}
return tx.Commit()
}
Не требует отдельной таблицы — состояние уже хранится в бизнес-объекте.
Retry с задержкой и Dead Letter Queue
Что если обработчик упал не из-за бага, а из-за временной недоступности внешнего сервиса? Нужно попробовать снова, но не немедленно.
Delayed retry собирается через x-message-ttl и Dead Letter Exchange:
ch.QueueDeclare("orders.retry", true, false, false, false, amqp.Table{
"x-queue-type": "quorum",
"x-message-ttl": int32(30_000), // ждём 30 секунд
"x-dead-letter-exchange": "orders",
"x-dead-letter-routing-key": "order.created",
})
Поток: обработчик отклонил сообщение → оно попадает в retry очередь → через 30 секунд по истечении TTL уходит через DLX обратно в основную очередь → новая попытка.
Количество попыток считается через заголовок x-death.count — его нужно проверять вручную, встроенного ограничителя нет.
Ещё одна ловушка — очередь с TTL отдаёт сообщения строго по порядку: сообщение с задержкой 5 секунд, вставшее за десятиминутным, дождётся его. Поэтому под каждую задержку заводят свою retry-очередь (orders.retry.5s, orders.retry.30s, orders.retry.10m) либо ставят плагин rabbitmq_delayed_message_exchange.
Плагин включается на брокере (rabbitmq-plugins enable rabbitmq_delayed_message_exchange) и добавляет особый тип обменника, у которого задержка задаётся у каждого сообщения заголовком:
ch.ExchangeDeclare("orders.delayed", "x-delayed-message", true, false, false, false,
amqp.Table{"x-delayed-type": "direct"})
ch.PublishWithContext(ctx, "orders.delayed", "order.created", false, false, amqp.Publishing{
Headers: amqp.Table{"x-delay": int32(30_000)}, // задержка этого сообщения, мс
ContentType: "application/json",
DeliveryMode: amqp.Persistent,
Body: body,
})
Одна очередь, произвольные задержки, никакой головной блокировки. Цена тоже есть: плагин держит отложенные сообщения в своей встроенной базе на том узле, куда они опубликованы, это не очередь, и она не реплицируется. Падение узла до срока доставки теряет отложенные сообщения, а большое их число заметно ест память. Для повторов и коротких задержек плагин удобен, а для «напомнить через месяц» берут не брокер, а таблицу в базе и задачу по расписанию.
Сообщения, которые не получилось обработать после всех попыток, уходят в Dead Letter Queue (DLQ) — отдельную очередь для разбора вручную или через алерты.
Гарантированная публикация — Outbox
Бывает задача: сохранить заказ в базу данных и опубликовать событие — атомарно. Если сначала сохранить, потом опубликовать, то сервис может упасть между двумя операциями. Событие потеряется.
Outbox pattern: событие сохраняется в ту же транзакцию, что и бизнес-данные. Отдельный процесс читает таблицу и публикует в AMQP.
func (s *OrderService) Confirm(ctx context.Context, orderID string) error {
tx, err := s.db.BeginTx(ctx, nil)
if err != nil {
return err
}
defer tx.Rollback()
if _, err := tx.ExecContext(ctx,
`UPDATE orders SET status = 'CONFIRMED' WHERE id = $1`, orderID); err != nil {
return err
}
payload, _ := json.Marshal(OrderConfirmedEvent{OrderID: orderID})
if _, err := tx.ExecContext(ctx,
`INSERT INTO outbox (id, routing_key, exchange, payload)
VALUES ($1, $2, $3, $4)`,
uuid.NewString(), "order.confirmed", "orders", payload); err != nil {
return err
}
return tx.Commit()
}
// Фоновая публикация — goroutine с тикером; канал переведён в режим подтверждений
func (s *OrderService) PublishOutbox(ctx context.Context) {
if err := s.ch.Confirm(false); err != nil {
return
}
ticker := time.NewTicker(500 * time.Millisecond)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return
case <-ticker.C:
batch, err := s.outbox.FetchUnpublished(ctx, 100)
if err != nil {
continue
}
for _, e := range batch {
confirm, err := s.ch.PublishWithDeferredConfirmWithContext(ctx, e.Exchange, e.RoutingKey, false, false,
amqp.Publishing{DeliveryMode: amqp.Persistent, Body: e.Payload})
if err != nil {
break
}
if acked, err := confirm.WaitContext(ctx); err != nil || !acked {
break // брокер не подтвердил: строка остаётся неопубликованной, повторим на следующем тике
}
s.outbox.MarkPublished(ctx, e.ID)
}
}
}
}
Две детали, без которых «гарантированная» публикация ничего не гарантирует. PublishWithContext возвращает nil, как только кадр ушёл в сокет, — брокер мог его ещё не принять; поэтому канал переводят в режим publisher confirms (ch.Confirm) и ждут подтверждения через DeferredConfirmation, и только после него помечают строку. И DeliveryMode: amqp.Persistent — без него сообщение живёт в памяти брокера и пропадёт при его перезапуске даже из durable-очереди.
Либо оба изменения зафиксированы, либо ни одного. Дубли возможны (публикация прошла, но пометить как отправленное не успело) — поэтому получатель всё равно должен быть идемпотентным.
Шпаргалка по выбору
| Задача | Паттерн | Тип обменника |
|---|---|---|
| Распределить нагрузку между обработчиками | Work Queue | direct (default) |
| Broadcast события всем сервисам | Publish/Subscribe | fanout |
| Разные события в разные очереди | Routing | direct |
| Подписка по маске на иерархические события | Topic | topic |
| Синхронный вызов через очередь | RPC | direct + reply-to |
| Защита от повторной доставки | Idempotent Consumer | любой |
| Retry с задержкой | Delayed Retry | direct + DLX |
| Атомарная публикация вместе с записью в БД | Outbox | direct |
Глубже: чего amqp091-go не делает за васрасширенное
Java-клиент RabbitMQ переподключается сам; amqp091-go — нет. Оборвалось соединение — Connection и все его каналы становятся невалидными навсегда, а цикл for d := range msgs тихо завершается, потому что канал доставок закрыт. Поэтому в сервисе всегда есть цикл переподключения: подписка на conn.NotifyClose(make(chan *amqp.Error, 1)), при закрытии — пауза с нарастающей задержкой, новый Dial, новые каналы и повторное объявление обменников, очередей и привязок (объявления идемпотентны, повторять их безопасно). Вторая деталь из документации Consume: доставки нужно читать из канала постоянно — непрочитанные блокируют все методы на том же соединении, так что «получил и отложил на потом» здесь не работает.
Коротко
- Work Queue — одна очередь, несколько обработчиков, каждое сообщение получает ровно один. Для фоновых задач.
- Publish/Subscribe — fanout exchange копирует сообщение во все привязанные очереди. Для broadcast-событий.
- Routing — direct exchange смотрит на ключ маршрутизации. Для точного разделения потоков.
- Topic — как routing, но с шаблонами
*и#. Для иерархических событий с гибкой подпиской. - RPC через очередь — запрос-ответ через брокер с
reply-toиcorrelation-id. Для вызовов без HTTP. - Idempotent Consumer — at-least-once означает возможные дубли. Защита: ключ идемпотентности в БД или проверка состояния объекта.
- Delayed Retry — TTL + DLX: сообщение «паркуется» на время, потом возвращается.
- Outbox — событие сохраняется в той же транзакции, что и данные. Атомарность без двухфазного коммита.
- amqp091-go не переподключается сам и просит не делить канал между горутинами: цикл переподключения по
NotifyClose, канал на воркера, publisher confirms иDeliveryMode: Persistentтам, где терять нельзя.
Что почитать дальше
- Протокол AMQP — модель exchange/binding/queue изнутри.
- Фоновые горутины и outbox-relay — как relay из этой статьи останавливается на выкате, не бросив пачку.
- RabbitMQ в production — Quorum Queues, кластеризация, мониторинг.
- AMQP vs Kafka — какой брокер и когда брать.