← назад к разделу

Когда сервисы начинают общаться через брокер, возникает вопрос: как правильно организовать очереди, обменники и подписки? Большинство задач укладываются в несколько типовых схем. Разберём каждую на примерах Go и rabbitmq/amqp091-go.

Обязательно

Раздать задачи нескольким обработчикам — Work Queue

Представьте: пользователи загружают фотографии, и каждую надо сжать и нарезать в несколько размеров. Это долго. Делать это прямо в HTTP-запросе нельзя — пользователь будет ждать минуты.

Решение — work queue (очередь задач): сохранить задание в очередь и отдать её одному из пула обработчиков.

producer одна очередь consumer 1 consumer 2 consumer 3

Очередь одна, потребителей несколько: брокер раздаёт сообщения по одному, и каждое достаётся ровно кому-то одному. Больше потребителей — быстрее разбирается очередь, но общий порядок между ними уже не сохраняется.

Брокер сам распределяет задачи между обработчиками. Каждое сообщение получает ровно один из них.

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, &quote)
		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 Queuedirect (default)
Broadcast события всем сервисамPublish/Subscribefanout
Разные события в разные очередиRoutingdirect
Подписка по маске на иерархические событияTopictopic
Синхронный вызов через очередьRPCdirect + reply-to
Защита от повторной доставкиIdempotent Consumerлюбой
Retry с задержкойDelayed Retrydirect + DLX
Атомарная публикация вместе с записью в БДOutboxdirect
Дополнительно: при первом чтении можно пропустить

Глубже: чего 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 там, где терять нельзя.

Что почитать дальше