golang

Надёжная асинхронная коммуникация: повторы, дубликаты и dead letter queues

  • пятница, 24 июля 2026 г. в 00:00:10
https://habr.com/ru/articles/1062002/

Представим обычную обработку заказа. Сервис заказов публикует событие order.created. Сервис склада получает его и резервирует товар в PostgreSQL. После успешной транзакции обработчик должен отправить RabbitMQ подтверждение (Ack), чтобы broker удалил сообщение из queue.

Но процесс может остановиться после записи в PostgreSQL и до отправки Ack. RabbitMQ не знает, успел ли сервис зарезервировать товар. Broker видит только неподтверждённое сообщение, поэтому доставляет его ещё раз. С точки зрения доставки это правильное поведение. С точки зрения бизнеса один заказ теперь может зарезервировать товар дважды.

Другой сбой возникает раньше: PostgreSQL временно недоступен, и обработчик не может начать работу. Если сразу вернуть сообщение в queue через отрицательное подтверждение Nack(requeue=true), RabbitMQ почти немедленно доставит его снова. Пока база не восстановилась, все попытки будут бесполезными. Нужны задержка и ограничение числа повторов. При этом отложенное сообщение может пропустить вперёд более новые события, поэтому отдельно придётся решить вопрос порядка.

Так одна операция превращается в несколько независимых участков: запись события, публикация, хранение в broker, обработка и подтверждение. Между соседними участками остаются моменты, когда одна сторона уже выполнила действие, а другая ещё не получила подтверждение.

В статье разберём эти моменты по всему пути сообщения. Затем построим практическую схему для RabbitMQ и Go: добавим ограниченные повторы через retry queues, время жизни сообщения (TTL) и dead letter exchange, сделаем обработчик идемпотентным и определим, куда отправлять сообщения, которые не удалось обработать автоматически. В конце сравним этот подход с Kafka, NATS JetStream и Amazon SQS.

Главный вывод: надёжность асинхронной коммуникации означает не отсутствие повторов, а предсказуемое поведение при каждом сбое. Producer не должен терять событие, повторная доставка не должна повторять бизнес-эффект, а retry всегда должен иметь задержку, предел и конечный маршрут.

Где ломается путь сообщения

Рассмотрим полный путь order.created:

orders DB -> outbox -> publisher -> exchange -> queue -> consumer -> inventory DB
                              |                    |
                         confirm/return          Ack/Nack

На этом пути нет одного общего подтверждения. Есть несколько локальных фактов:

  1. PostgreSQL подтвердил транзакцию заказа.

  2. RabbitMQ подтвердил, что принял публикацию.

  3. Exchange нашёл подходящую queue.

  4. Consumer подтвердил доставку через Ack.

  5. PostgreSQL склада зафиксировал резерв.

Сбой возможен между любыми соседними пунктами. Поэтому полезнее спрашивать не «гарантирует ли RabbitMQ доставку», а «кто сейчас владеет сообщением и какое действие можно безопасно повторить».

Проблема

Как проявляется

Что обычно помогает

Dual write

Заказ записан, а событие не опубликовано, или наоборот

Transactional Outbox

Неопределённый результат публикации

Broker принял сообщение, но confirm потерялся

Повтор с тем же message_id, идемпотентный consumer

Повторная доставка

Consumer выполнил изменение и упал до Ack

Inbox или идемпотентная бизнес-операция

Горячий цикл повторов

Неисправное сообщение непрерывно возвращается в queue

Задержка, backoff, лимит попыток

Poison message

Повтор не изменит результат: JSON повреждён или схема неизвестна

Parking lot, или dead letter queue (DLQ), диагностика и ручной redrive - возврат после исправления

Нарушение порядка

Более новое событие обработано раньше повторяемого старого

Партиционирование по ключу, версия агрегата, последовательная обработка

Перегрузка consumer

Queue и возраст сообщений растут быстрее, чем сервис успевает обрабатывать

Prefetch, ограничение конкуренции, backpressure, масштабирование

Несовместимый контракт

Старый consumer не понимает новое поле или смысл события

Версия схемы, совместимые изменения, постепенная миграция

Потеря контекста

Нельзя связать сообщение с запросом и предыдущими попытками

Стабильный ID, correlation ID, trace context, метрики

Одна настройка broker не закрывает всю таблицу. Например, quorum queue защищает данные на стороне RabbitMQ, но ничего не знает о транзакции PostgreSQL внутри consumer.

At most once, at least once и exactly once

У доставки есть три часто используемых описания.

At most once, или не более одного раза. Consumer подтверждает сообщение до обработки либо работает с автоматическим Ack. Повторов нет, но при падении consumer работа может потеряться.

At least once, или как минимум один раз. Consumer сначала выполняет работу и только потом отправляет Ack. Работа не должна потеряться, зато сообщение может прийти повторно.

Exactly once, или ровно один раз. Термин полезен только вместе с границами гарантии. Kafka умеет атомарно записать результат и consumer offset в Kafka-транзакции, когда приложение читает из Kafka и пишет обратно в Kafka. Но запись в произвольную базу или вызов платёжного API снова создают внешнюю границу. Это ограничение прямо описано в документации Kafka о delivery semantics.

RabbitMQ связывает publisher confirms и consumer acknowledgements с передачей ответственности за сообщение. Документация называет надёжность совместной ответственностью broker, producer и consumer. Эти подтверждения не образуют распределённую транзакцию с базой приложения.

На практике рабочая цель звучит так:

at least once delivery + idempotent business effect

Сообщение разрешено доставить несколько раз, но наблюдаемый бизнес-результат должен появиться один раз. Важно именно последнее. Таблица с обработанными ID не поможет, если consumer сначала дважды вызовет внешний API, а только потом посмотрит в таблицу.

Публикация начинается с dual write

Допустим, сервис заказов делает две независимые записи:

if err := orders.Insert(ctx, order); err != nil {
	return err
}

if err := publisher.Publish(ctx, orderCreated); err != nil {
	return err
}

Между вызовами процесс может упасть. Перестановка строк не решает проблему: тогда событие способно уйти до появления заказа.

Transactional Outbox переносит событие в ту же транзакцию, где создаётся заказ:

BEGIN;

INSERT INTO orders (id, status)
VALUES ('42', 'created');

INSERT INTO outbox (id, event_type, payload, created_at)
VALUES (
    '01K0ORDERCREATED',
    'order.created',
    '{"order_id":"42"}',
    now()
);

COMMIT;

Отдельный relay читает outbox, публикует событие и после publisher confirm отмечает запись отправленной. Если relay упадёт после публикации, но до отметки в таблице, он отправит событие повторно. Поэтому ID создают при записи в outbox и сохраняют при каждом повторе. Новый ID превратил бы тот же факт в новое событие и сделал дедупликацию невозможной.

Publisher confirm сообщает, что RabbitMQ принял ответственность за публикацию. Для persistent-сообщения, попавшего в durable queue, это означает сохранение на диске. Quorum queue подтверждает публикацию после того, как сообщение приняло большинство реплик. Если exchange направил сообщение в несколько queues, confirm придёт после того, как его примут все выбранные queues. Но потерянный ответ снова оставляет producer в неопределённости: broker мог принять сообщение, а confirm мог не дойти. RabbitMQ рекомендует повторно публиковать неподтверждённые сообщения и быть готовым к дубликатам.

Есть ещё одна деталь. Confirm не означает, что exchange нашёл queue. Сообщение без подходящего маршрута, или unroutable message, RabbitMQ тоже может подтвердить. Для обязательных событий producer публикует с mandatory=true и читает basic.return. Если отсутствие подписчиков допустимо по смыслу события, потеря такого сообщения может быть нормальной политикой.

Когда отправлять Ack

Consumer должен отправить Ack только после того, как действие стало устойчивым к падению процесса. Для записи в PostgreSQL это означает успешный COMMIT.

Рассмотрим четыре варианта:

Момент падения

Что увидит RabbitMQ

Что должен сделать consumer после запуска

До начала транзакции

Ack не получен, сообщение вернётся

Выполнить обработку

Во время транзакции

Транзакция откатится, сообщение вернётся

Выполнить обработку

После COMMIT, до Ack

Сообщение вернётся, хотя результат уже записан

Распознать дубликат и ничего не менять

После Ack

RabbitMQ удалит сообщение

Результат уже обязан быть сохранён

Автоматический Ack сдвигает границу слишком рано. Ack до COMMIT делает окно потери. Ack после COMMIT оставляет окно дубликата, которое можно закрыть идемпотентностью.

Флаг Redelivered полезен для диагностики, но не для корректности. Он говорит, что RabbitMQ уже пытался выдать сообщение, а не что предыдущий consumer успел изменить данные. Consumer может впервые увидеть сообщение с Redelivered=true, если предыдущая доставка не дошла до приложения перед разрывом соединения. RabbitMQ прямо называет этот флаг подсказкой, а не доказательством выполненной работы.

Идемпотентный consumer через Inbox

Идемпотентная операция даёт одинаковый наблюдаемый результат при одном и нескольких одинаковых вызовах. Присваивание status = 'paid' часто идемпотентно. Команда balance = balance - 100 или отправка письма - нет.

Для произвольного изменения удобно использовать Inbox. Consumer хранит ID обработанных событий в собственной базе. Уникальное ограничение и бизнес-изменение входят в одну транзакцию:

CREATE TABLE processed_messages (
    consumer_name text        NOT NULL,
    message_id    text        NOT NULL,
    processed_at  timestamptz NOT NULL DEFAULT now(),
    PRIMARY KEY (consumer_name, message_id)
);

Ниже упрощённый метод резервирования товара. Если сообщение уже обработано этим consumer, INSERT ... ON CONFLICT не вернёт строку и повтор завершится без второго резерва.

type OrderCreated struct {
	OrderID string `json:"order_id"`
	SKU     string `json:"sku"`
	Count   int    `json:"count"`
}

func reserve(
	ctx context.Context,
	db *sql.DB,
	messageID string,
	event OrderCreated,
) error {
	tx, err := db.BeginTx(ctx, nil)
	if err != nil {
		return fmt.Errorf("begin transaction: %w", err)
	}
	defer func() {
		// После успешного Commit вернётся sql.ErrTxDone.
		_ = tx.Rollback()
	}()

	var inserted int
	err = tx.QueryRowContext(ctx, `
		INSERT INTO processed_messages (consumer_name, message_id)
		VALUES ('inventory.reserve', $1)
		ON CONFLICT DO NOTHING
		RETURNING 1
	`, messageID).Scan(&inserted)
	switch {
	case errors.Is(err, sql.ErrNoRows):
		return nil // событие уже применили, defer закроет транзакцию
	case err != nil:
		return fmt.Errorf("claim message: %w", err)
	}

	result, err := tx.ExecContext(ctx, `
		UPDATE stock
		SET reserved = reserved + $1
		WHERE sku = $2
		  AND available - reserved >= $1
	`, event.Count, event.SKU)
	if err != nil {
		return fmt.Errorf("reserve stock: %w", err)
	}

	rows, err := result.RowsAffected()
	if err != nil {
		return fmt.Errorf("read affected rows: %w", err)
	}
	if rows == 0 {
		// Отсутствие товара - бизнес-результат, а не сбой инфраструктуры.
		// Сохраняем исходящее событие в той же транзакции.
		_, err := tx.ExecContext(ctx, `
			INSERT INTO inventory_outbox (id, event_type, payload)
			VALUES (
				$1,
				'inventory.reservation_rejected',
				jsonb_build_object('order_id', $2, 'sku', $3)
			)
			ON CONFLICT DO NOTHING
		`, messageID+":reservation-rejected", event.OrderID, event.SKU)
		if err != nil {
			return fmt.Errorf("save reservation rejection: %w", err)
		}
	}

	if err := tx.Commit(); err != nil {
		return fmt.Errorf("commit reservation: %w", err)
	}
	return nil
}

Порядок важен: уникальный ID и резерв фиксируются одним COMMIT. Если сделать две отдельные транзакции, падение между ними либо потеряет работу, либо снова допустит дубликат. Отложенный Rollback безопасен и после успешного Commit: завершённая транзакция вернёт sql.ErrTxDone, а код намеренно проигнорирует эту ошибку.

Ошибку Commit нельзя автоматически назвать временной. Отдельная проблема в том, что при обрыве соединения результат иногда остаётся неизвестным приложению. Например, PostgreSQL успел зафиксировать транзакцию, а ответ не дошёл до consumer. В показанном цикле такая ошибка пойдёт по ветке retry. RabbitMQ доставит сообщение ещё раз, а запись в Inbox не даст повторно применить уже зафиксированное изменение.

Inbox растёт, поэтому для него нужна политика хранения. Удалять записи можно только после максимального срока, в течение которого старое событие способно вернуться из broker, архива или ручного redrive. Если события разрешено переигрывать через год, очистка Inbox через неделю вернёт побочные эффекты вместе с replay.

Ветка с отсутствующим товаром заодно показывает, что consumer часто становится producer. Запись в inventory_outbox входит в ту же транзакцию, а отдельный relay опубликует inventory.reservation_rejected. Таблица должна иметь уникальный ключ id, поэтому повтор исходного события не создаст два исходящих.

Для внешнего HTTP API локальная таблица не создаёт атомарность. Передавайте message_id как idempotency key, если API поддерживает такой контракт. Если не поддерживает, сохраняйте намерение локально и выполняйте вызов отдельным процессом либо проектируйте компенсирующую операцию. Волшебного Ack для чужой системы здесь не появляется.

Какие ошибки нужно повторять

Не каждая ошибка временна. Retry policy начинается с классификации, а не с формулы backoff.

Временная ошибка: connection reset, недоступная база, 429 Too Many Requests, deadlock или короткий timeout. Повтор может сработать без изменения сообщения и кода.

Постоянная ошибка: повреждённый JSON, неизвестная версия схемы, отсутствующее обязательное поле, запрещённый переход состояния. Повтор тех же байтов ничего не исправит.

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

Неизвестная ошибка: новый дефект кода или неожиданный ответ зависимости. Бесконечно считать её временной опасно. Разумный вариант - несколько ограниченных повторов, затем parking lot и сигнал дежурному.

Иногда класс зависит от времени. customer not found может означать постоянную ошибку, а может появиться из-за того, что событие клиента ещё не дошло до проекции данных для чтения, или read model. Тогда retry policy должна содержать не только число попыток, но и бизнес-дедлайн: например, ждать не более 15 минут после occurred_at.

Хороший retry budget задаёт:

  • максимальное число попыток;

  • максимальный возраст сообщения;

  • задержки между попытками;

  • случайный разброс, или jitter, чтобы consumers не проснулись одновременно;

  • конечный маршрут после исчерпания бюджета.

Последовательность 5 секунд -> 30 секунд -> 5 минут понятнее бесконечного экспоненциального повтора. Конкретные числа зависят от времени восстановления зависимости и допустимой задержки бизнеса. Если база обычно восстанавливается за 20 минут, три повтора за 5 минут 35 секунд всё равно лишь быстрее доставят весь поток в DLQ.

Почему Nack(requeue=true) опасен

При requeue=true сообщение снова становится доступным для доставки. Если все consumers получают одну и ту же временную ошибку, они создают requeue/redelivery loop. Документация RabbitMQ предупреждает, что такой цикл расходует CPU и сетевой трафик.

if err := handle(delivery.Body); err != nil {
	// Плохая политика: задержки и лимита нет.
	return delivery.Nack(false, true)
}

Не стоит заменять broker обычным time.Sleep внутри consumer:

time.Sleep(5 * time.Minute)
return delivery.Nack(false, true)

Пять минут сообщение остаётся unacked, занимает окно prefetch и зависит от жизни процесса. При падении задержка обнулится. Тысяча таких сообщений потребует тысячу ожидающих обработчиков или остановит получение новой работы.

Задержку лучше хранить на стороне broker. Тогда consumer освобождает свой слот, а сообщение переживает перезапуск приложения.

Retry queues, TTL и dead letter exchange в RabbitMQ

В AMQP (Advanced Message Queuing Protocol) 0-9-1 нет обычной команды «доставь это сообщение через 30 секунд». Переносимая схема строится из retry queues. У этих queues нет consumers. Сообщение ждёт свой TTL, истекает и через dead letter exchange возвращается в рабочую queue.

Сделаем три корзины задержек:

                                  transient error
                                       |
                                       v
orders -> inventory.orders -> consumer -> inventory.orders.retry (direct)
              ^                         | 5s  -> retry.5s  -- TTL 5s  --+
              |                         | 30s -> retry.30s -- TTL 30s --+-> orders
              |                         | 5m  -> retry.5m  -- TTL 5m  --+
              |
              +---------------------- expired

permanent/exhausted
        |
        v
inventory.orders.dlx -> inventory.orders.parking

Здесь есть три разных понятия.

TTL, или time to live, определяет, сколько сообщение может находиться в queue. После истечения message-ttl RabbitMQ не отдаёт его consumer и отправляет в настроенный dead letter exchange (DLX) либо удаляет.

Dead letter exchange - это exchange, а не хранилище. RabbitMQ публикует туда сообщения после basic.reject или basic.nack с requeue=false, истечения TTL, превышения длины queue или delivery limit quorum queue. Полный список причин есть в документации DLX. Чтобы сообщения где-то сохранились, к exchange нужно привязать queue.

Parking lot, часто называемый DLQ, - обычная queue для сообщений, которые автоматическая обработка прекратила. Для неё задают срок хранения, или retention, алерт и контролируемый redrive - возврат сообщений в рабочую queue после устранения причины ошибки.

Почему нужны отдельные queues для задержек

Можно задать свой Expiration каждому сообщению и сложить все повторы в одну queue. Но сообщение с коротким TTL способно оказаться за сообщением с длинным TTL. RabbitMQ удалит короткое сообщение только тогда, когда оно достигнет головы queue. До этого уже истёкшее сообщение продолжит занимать место и учитываться в статистике. Это ограничение отдельно описано в документации message TTL.

Небольшой набор queues с одинаковым TTL внутри каждой избегает этой проблемы. Задержка получается приблизительной: retry.30s означает «не раньше чем примерно через 30 секунд», а не точный планировщик. Для обработки заказа такая семантика обычно подходит.

Объявляем топологию из Go

Тип queue задаётся при объявлении. Для важных данных используем quorum queues. Они реплицируют состояние через протокол консенсуса Raft и подтверждают persistent-публикацию после принятия большинством реплик. В локальном примере можно запустить один узел, но production-смысл quorum queue появляется с нечётным числом реплик, обычно с тремя.

func declareTopology(ch *amqp.Channel) error {
	for _, exchange := range []struct {
		name string
		kind string
	}{
		{"orders", "topic"},
		{"inventory.orders.retry", "direct"},
		{"inventory.orders.dlx", "direct"},
	} {
		if err := ch.ExchangeDeclare(
			exchange.name,
			exchange.kind,
			true,  // durable
			false, // auto-delete
			false, // internal
			false, // no-wait
			nil,
		); err != nil {
			return fmt.Errorf("declare exchange %s: %w", exchange.name, err)
		}
	}

	queueArgs := amqp.Table{
		amqp.QueueTypeArg: amqp.QueueTypeQuorum,
	}
	queues := []string{
		"inventory.orders",
		"inventory.orders.retry.5s",
		"inventory.orders.retry.30s",
		"inventory.orders.retry.5m",
		"inventory.orders.parking",
	}
	for _, name := range queues {
		if _, err := ch.QueueDeclare(
			name,
			true,  // durable
			false, // auto-delete
			false, // exclusive
			false, // no-wait
			queueArgs,
		); err != nil {
			return fmt.Errorf("declare queue %s: %w", name, err)
		}
	}

	bindings := []struct {
		queue, key, exchange string
	}{
		{"inventory.orders", "order.created", "orders"},
		{"inventory.orders.retry.5s", "5s", "inventory.orders.retry"},
		{"inventory.orders.retry.30s", "30s", "inventory.orders.retry"},
		{"inventory.orders.retry.5m", "5m", "inventory.orders.retry"},
		{"inventory.orders.parking", "inventory.orders.failed", "inventory.orders.dlx"},
	}
	for _, binding := range bindings {
		if err := ch.QueueBind(
			binding.queue,
			binding.key,
			binding.exchange,
			false,
			nil,
		); err != nil {
			return fmt.Errorf("bind queue %s: %w", binding.queue, err)
		}
	}

	return nil
}

TTL и dead lettering лучше задавать RabbitMQ policies, а не жёсткими x-arguments в приложении. Policy можно изменить без удаления queue и нового развёртывания сервиса. RabbitMQ прямо рекомендует policies для DLX.

Для основной queue зададим parking lot и ограничение неожиданных повторных доставок. Для retry queues зададим TTL и возврат в exchange orders:

rabbitmqctl set_policy inventory-orders-main \
  '^inventory\.orders$' \
  '{"dead-letter-exchange":"inventory.orders.dlx","dead-letter-routing-key":"inventory.orders.failed","dead-letter-strategy":"at-least-once","overflow":"reject-publish","delivery-limit":8}' \
  --apply-to quorum_queues

rabbitmqctl set_policy inventory-orders-retry-5s \
  '^inventory\.orders\.retry\.5s$' \
  '{"message-ttl":5000,"dead-letter-exchange":"orders","dead-letter-routing-key":"order.created","dead-letter-strategy":"at-least-once","overflow":"reject-publish"}' \
  --apply-to quorum_queues

rabbitmqctl set_policy inventory-orders-retry-30s \
  '^inventory\.orders\.retry\.30s$' \
  '{"message-ttl":30000,"dead-letter-exchange":"orders","dead-letter-routing-key":"order.created","dead-letter-strategy":"at-least-once","overflow":"reject-publish"}' \
  --apply-to quorum_queues

rabbitmqctl set_policy inventory-orders-retry-5m \
  '^inventory\.orders\.retry\.5m$' \
  '{"message-ttl":300000,"dead-letter-exchange":"orders","dead-letter-routing-key":"order.created","dead-letter-strategy":"at-least-once","overflow":"reject-publish"}' \
  --apply-to quorum_queues

Зачем здесь dead-letter-strategy=at-least-once? Обычный DLX переносит сообщение внутренней публикацией без publisher confirms и удаляет его из исходной queue. Если целевая queue недоступна, сообщение можно потерять. Quorum queues поддерживают at-least-once dead lettering: исходная queue хранит сообщение, пока целевая сторона не подтвердит приём. Если подтверждение потеряется, целевая queue может получить дубликат, поэтому идемпотентность всё ещё нужна. Такой перенос требует больше CPU и памяти, а также overflow=reject-publish. На кластерах, обновлённых со старых версий RabbitMQ, также нужно проверить feature flag stream_queue.

RabbitMQ рекомендует дополнительно ограничить исходные queues через max-length или max-length-bytes. Если parking queue или рабочая queue долго не принимает сообщения, безопасный DLX продолжит хранить их в источнике. Конкретный лимит зависит от размера сообщений, доступного диска и допустимого backlog, поэтому в примере он не задан произвольной константой.

На classic queues и при стандартной стратегии DLX остаётся at most once. Это важная оговорка: наличие DLX само по себе ещё не означает надёжный перенос.

Публикуем повтор до Ack исходного сообщения

Consumer должен выбрать следующую задержку. Для этого добавим собственный header x-retry-attempt. Не стоит считать попытки только по x-death: этот header меняется при dead lettering, зависит от queue и причины, а обычная повторная публикация приложения создаёт другую историю.

Порядок действий после временной ошибки такой:

  1. Скопировать сообщение с тем же MessageId.

  2. Увеличить x-retry-attempt.

  3. Опубликовать копию в нужную retry queue.

  4. Дождаться publisher confirm и проверить basic.return.

  5. Только после этого отправить Ack исходной доставке.

Если consumer отправит Ack раньше публикации, сбой потеряет сообщение. Если процесс упадёт после подтверждённой публикации, но до Ack, появятся две копии. Убрать обе опасности одной локальной транзакцией нельзя, поэтому выбираем дубликат и полагаемся на Inbox.

Ниже publisher сериализует вызовы на отдельном AMQP channel. Это упрощает связь basic.return с единственной публикацией в полёте.

type ConfirmingPublisher struct {
	mu      sync.Mutex
	ch      *amqp.Channel
	returns <-chan amqp.Return
}

func NewConfirmingPublisher(ch *amqp.Channel) (*ConfirmingPublisher, error) {
	if err := ch.Confirm(false); err != nil {
		return nil, fmt.Errorf("enable publisher confirms: %w", err)
	}

	returns := ch.NotifyReturn(make(chan amqp.Return, 1))
	return &ConfirmingPublisher{ch: ch, returns: returns}, nil
}

func (p *ConfirmingPublisher) Close() error {
	p.mu.Lock()
	defer p.mu.Unlock()
	return p.ch.Close()
}

func (p *ConfirmingPublisher) Publish(
	ctx context.Context,
	exchange string,
	key string,
	msg amqp.Publishing,
) error {
	p.mu.Lock()
	defer p.mu.Unlock()

	confirmation, err := p.ch.PublishWithDeferredConfirmWithContext(
		ctx,
		exchange,
		key,
		true,  // mandatory
		false, // immediate RabbitMQ не поддерживает
		msg,
	)
	if err != nil {
		return fmt.Errorf("publish: %w", err)
	}

	acked, err := confirmation.WaitContext(ctx)
	if err != nil {
		return fmt.Errorf("wait for confirm: %w", err)
	}
	if !acked {
		return errors.New("broker nacked publication")
	}

	// RabbitMQ отправляет basic.return до confirm. Channel выделен этому
	// publisher, а публикации сериализованы, поэтому return относится к msg.
	select {
	case returned, ok := <-p.returns:
		if !ok {
			return errors.New("return notifications channel closed")
		}
		return fmt.Errorf(
			"unroutable message %s: %s",
			returned.MessageId,
			returned.ReplyText,
		)
	default:
		return nil
	}
}

Production-обёртка также должна завершать ожидающие публикации при NotifyClose, восстанавливать channel и учитывать все неподтверждённые сообщения. Подробно жизненный цикл соединения разобран в посте «RabbitMQ и Go».

Теперь выберем retry queue и сохраним исходные headers:

var retryPlan = []string{"5s", "30s", "5m"}

func scheduleRetry(
	ctx context.Context,
	publisher *ConfirmingPublisher,
	delivery amqp.Delivery,
	cause error,
) (scheduled bool, err error) {
	attempt := headerInt(delivery.Headers["x-retry-attempt"])
	if attempt >= len(retryPlan) {
		return false, nil
	}

	headers := cloneHeaders(delivery.Headers)
	headers["x-retry-attempt"] = int32(attempt + 1)
	headers["x-last-error"] = truncate(cause.Error(), 512)

	err = publisher.Publish(ctx, "inventory.orders.retry", retryPlan[attempt], amqp.Publishing{
		Headers:         headers,
		ContentType:     delivery.ContentType,
		ContentEncoding: delivery.ContentEncoding,
		DeliveryMode:    amqp.Persistent,
		CorrelationId:   delivery.CorrelationId,
		MessageId:       delivery.MessageId, // ID не меняется при retry
		Timestamp:       delivery.Timestamp,
		Type:            delivery.Type,
		Body:            delivery.Body,
	})
	if err != nil {
		return false, err
	}
	return true, nil
}

func cloneHeaders(src amqp.Table) amqp.Table {
	dst := make(amqp.Table, len(src)+2)
	for key, value := range src {
		dst[key] = value
	}
	return dst
}

func headerInt(value any) int {
	var result int
	switch value := value.(type) {
	case int8:
		result = int(value)
	case int16:
		result = int(value)
	case int32:
		result = int(value)
	case int64:
		result = int(value)
	default:
		return 0
	}
	if result < 0 {
		return 0
	}
	return result
}

func truncate(value string, limit int) string {
	runes := []rune(value)
	if len(runes) <= limit {
		return value
	}
	return string(runes[:limit])
}

Мы намеренно не копируем Expiration. Если producer задал исходному событию TTL, его нужно трактовать как бизнес-дедлайн отдельно. Слепое копирование может привести к тому, что повтор истечёт раньше выбранной задержки.

Пример ограничивает число попыток, но не проверяет максимальный возраст события и использует фиксированные задержки без jitter. В production consumer должен прочитать occurred_at из envelope, сравнить его с бизнес-дедлайном и при необходимости сразу завершить retry. Jitter можно добавить выбором одной из соседних корзин задержки. Использовать произвольный per-message TTL в общей retry queue опасно из-за описанной выше блокировки в голове queue.

Consumer с тремя исходами

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

type PermanentError struct {
	Err error
}

func (e *PermanentError) Error() string { return e.Err.Error() }
func (e *PermanentError) Unwrap() error { return e.Err }

func process(
	ctx context.Context,
	db *sql.DB,
	delivery amqp.Delivery,
) error {
	if delivery.MessageId == "" {
		return &PermanentError{Err: errors.New("message_id is required")}
	}

	var event OrderCreated
	if err := json.Unmarshal(delivery.Body, &event); err != nil {
		return &PermanentError{Err: fmt.Errorf("decode order.created: %w", err)}
	}
	if event.OrderID == "" || event.SKU == "" || event.Count <= 0 {
		return &PermanentError{Err: errors.New("invalid order.created")}
	}

	return reserve(ctx, db, delivery.MessageId, event)
}

Для краткости пример считает все ошибки reserve временными. В production нужно классифицировать и ошибки базы: deadlock, serialization failure и обрыв соединения обычно требуют retry, а нарушение constraint может означать дефект данных или кода, который повтор не исправит.

Цикл чтения применяет политику:

func consume(
	ctx context.Context,
	consumerChannel *amqp.Channel,
	db *sql.DB,
	retryPublisher *ConfirmingPublisher,
) error {
	// Функция владеет consumerChannel. При любом выходе незавершённые
	// доставки вернутся в queue.
	defer func() {
		_ = consumerChannel.Close()
	}()

	if err := consumerChannel.Qos(16, 0, false); err != nil {
		return fmt.Errorf("set qos: %w", err)
	}

	deliveries, err := consumerChannel.Consume(
		"inventory.orders",
		"inventory-worker",
		false, // auto-ack
		false, // exclusive
		false, // no-local
		false, // no-wait
		nil,
	)
	if err != nil {
		return fmt.Errorf("consume: %w", err)
	}

	for {
		select {
		case <-ctx.Done():
			return ctx.Err()

		case delivery, ok := <-deliveries:
			if !ok {
				return errors.New("delivery channel closed")
			}

			handleCtx, cancel := context.WithTimeout(ctx, 10*time.Second)
			err := process(handleCtx, db, delivery)
			cancel()

			if err == nil {
				if err := delivery.Ack(false); err != nil {
					return fmt.Errorf("ack: %w", err)
				}
				continue
			}

			var permanent *PermanentError
			if errors.As(err, &permanent) {
				// requeue=false отправит сообщение в inventory.orders.dlx.
				if nackErr := delivery.Nack(false, false); nackErr != nil {
					return fmt.Errorf("park message: %w", nackErr)
				}
				continue
			}

			retryCtx, retryCancel := context.WithTimeout(ctx, 5*time.Second)
			scheduled, retryErr := scheduleRetry(
				retryCtx,
				retryPublisher,
				delivery,
				err,
			)
			retryCancel()
			if retryErr != nil {
				// Не подтверждаем исходную доставку и больше не используем
				// publisher channel с неопределённым состоянием.
				_ = retryPublisher.Close()
				return fmt.Errorf("schedule retry: %w", retryErr)
			}

			if !scheduled {
				if nackErr := delivery.Nack(false, false); nackErr != nil {
					return fmt.Errorf("park exhausted message: %w", nackErr)
				}
				continue
			}

			if err := delivery.Ack(false); err != nil {
				return fmt.Errorf("ack scheduled retry: %w", err)
			}
		}
	}
}

Код показывает границы ответственности, но не весь production runtime. Нужны reconnect с ограниченным backoff, повторное объявление топологии, graceful shutdown и метрики. При ошибке публикации внешний supervisor должен создать заново и consumer channel, и publisher channel. Один amqp.Channel не стоит использовать конкурентно из нескольких goroutines без собственной сериализации.

Для параллельной обработки проще запустить несколько workers, каждый со своим consumer channel, чем раздать amqp.Delivery произвольным goroutines и затем конкурентно отправлять подтверждения через общий channel. Prefetch задаёт верхнюю границу unacked-сообщений. Значение 0 означает отсутствие лимита, поэтому начинать с него рискованно. RabbitMQ объясняет связь prefetch, памяти и throughput в руководстве по acknowledgements.

Что изменилось в RabbitMQ 4.3

Начиная с RabbitMQ 4.3, quorum queues поддерживают встроенный delayed retry. Broker удерживает возвращённое сообщение до следующей доставки и рассчитывает линейную задержку:

delay = min(min_delay * delivery_count, max_delay)

RabbitMQ применяет к queue не более одной обычной policy. Если создать вторую policy с большим приоритетом, она не объединится с предыдущей, а вытеснит её. Поэтому обновим уже заданную inventory-orders-main и сохраним в одном JSON все параметры:

rabbitmqctl set_policy inventory-orders-main \
  '^inventory\.orders$' \
  '{"dead-letter-exchange":"inventory.orders.dlx","dead-letter-routing-key":"inventory.orders.failed","dead-letter-strategy":"at-least-once","overflow":"reject-publish","delivery-limit":8,"delayed-retry-type":"all","delayed-retry-min":1000,"delayed-retry-max":30000}' \
  --apply-to quorum_queues

После этого consumer может вернуть сообщение в ту же queue, а broker применит задержку. Если нужны растущая задержка и учёт delivery limit, в amqp091-go следует использовать Reject(true), который отправляет basic.reject:

if err := delivery.Reject(true); err != nil {
	return fmt.Errorf("reject for delayed retry: %w", err)
}

Отдельные retry queues в этом варианте не нужны. Решение проще, если линейный backoff подходит и вся инфраструктура уже работает на 4.3.

Есть важная деталь протокола. В RabbitMQ 4.3 basic.nack увеличивает x-acquired-count, но не x-delivery-count. basic.reject и потеря connection считаются неуспешной доставкой и увеличивают оба значения. Поэтому Nack(requeue=true) при линейном расчёте способен каждый раз получать минимальную задержку. Если политика опирается на delivery count и лимит попыток, проверьте поведение конкретного метода клиента в таблице poison message handling.

У встроенного механизма и TTL buckets разные компромиссы:

Подход

Плюсы

Ограничения

Retry queues + TTL + DLX

Работает на разных версиях, произвольные ступени задержек, явный маршрут

Больше queues и bindings, перенос через DLX, application republish создаёт окно дубликата

RabbitMQ 4.3 delayed retry

Нет отдельных queues и перепубликации, broker видит delayed messages

Только quorum queues 4.3+, линейный backoff, тонкости nack и reject

rabbitmq_delayed_message_exchange

Произвольная задержка в header

Community plugin больше не поддерживается и не входит в RabbitMQ

Последний вариант долго был популярен, но RabbitMQ сейчас помечает delayed message exchange как неподдерживаемый community plugin.

Начиная с RabbitMQ 4.0, delivery limit для quorum queues по умолчанию равен 20. После превышения сообщение удаляется или попадает в DLX, если он настроен. Это защита от poison messages, а не замена retry policy: версия 4.3 не считает обычный basic.nack неуспешной доставкой. Явный application header и собственный бюджет всё ещё проще связать с бизнес-дедлайном.

Dead letter queue требует процесса, а не только binding

Parking lot бесполезен, если команда узнаёт о нём через месяц. Для него нужен небольшой операционный процесс.

Сохраняйте рядом с сообщением:

  • исходный message_id, тип и версию схемы;

  • число попыток и время первой ошибки;

  • последнюю классификацию и короткое описание ошибки;

  • correlation ID и trace ID;

  • исходные routing key и producer;

  • версию consumer, на которой произошёл отказ.

RabbitMQ добавляет header x-death при dead lettering. В нём есть queue, причина, число событий и время. История сжимается по паре {queue, reason}, поэтому header не растёт на каждый цикл. Возможные причины включают rejected, expired, maxlen и delivery_limit. Структура описана в документации dead-lettered effects.

Нужны как минимум два алерта:

  1. В parking lot появилось новое сообщение.

  2. Само dead lettering не работает или целевая queue недоступна.

Redrive не должен быть автоматической петлёй parking -> main -> parking. Сначала исправляют данные, контракт или код. Затем оператор выбирает сообщения по понятному критерию, сохраняет их ID и возвращает ограниченной скоростью. Иначе одна команда redrive способна повторить исходный инцидент на всём backlog.

Для parking lot задают retention и доступ отдельно. Сообщения нередко содержат персональные данные, токены или бизнес-информацию. DLQ не должна становиться бессрочным архивом, который читает любой пользователь management UI.

Повторы нарушают порядок

Пусть для заказа 42 пришли два события:

v7: order.created
v8: order.cancelled

Снача consumer временно не смог обработать v7 и отправил его в retry.30s. Затем consumer успешно обработал оставшееся в рабочей queue событие v8. Через 30 секунд вернулось v7. FIFO (first in, first out - «первым пришёл, первым вышел») на исходной queue не спасает, потому что событие покинуло её и затем было опубликовано заново.

Несколько consumers, redelivery и priorities также меняют наблюдаемый порядок. RabbitMQ перечисляет эти случаи в руководстве по ordering.

Есть несколько стратегий, и ни одна не бесплатна.

Версия агрегата

Producer добавляет aggregate_id=42 и aggregate_version=7. Consumer хранит последнюю применённую версию.

  • версия равна ожидаемой - применить;

  • версия уже применена - подтвердить как дубликат;

  • версия больше ожидаемой - отложить и дождаться пропуска;

  • версия меньше текущей - не откатывать состояние назад.

Подход хорошо работает для событий состояния, но требует стратегии восстановления пропущенного события.

Партиционирование по бизнес-ключу

Все события одного заказа направляются в одну из нескольких queues по хешу order_id, а внутри каждой группы, или shard, работает один активный consumer. RabbitMQ предлагает для этого modulus hash exchange и Single Active Consumer. Параллелизм сохраняется между заказами, но долгая обработка одного заказа способна задержать весь его shard.

Остановить партицию

В Kafka порядок гарантируется внутри partition. Consumer может не продвигать offset, пока проблемная запись не обработана. Порядок сохраняется, но одна ошибка блокирует все следующие записи partition. Retry topic снимает блокировку и одновременно нарушает исходный порядок. Выбор между latency всего потока и порядком нужно сделать явно.

Правило простое: retry - это изменение порядка, пока архитектура отдельно не доказывает обратное.

Backpressure: queue не создаёт пропускную способность

Queue хорошо поглощает короткий всплеск. Если producer стабильно создаёт 10 000 сообщений в секунду, а consumers обрабатывают 8 000, backlog будет расти на 2 000 сообщений в секунду. Broker лишь отсрочит момент, когда закончится диск или нарушится допустимое время обработки.

Наблюдайте не только длину queue. Сто сообщений по 10 минут важнее десяти тысяч сообщений, появившихся секунду назад. Полезные сигналы:

  • возраст самого старого готового сообщения;

  • входящая и подтверждаемая скорость;

  • число ready и unacked messages;

  • длительность обработки и доля ошибок;

  • число сообщений в каждой retry queue;

  • количество повторов на одно сообщение;

  • скорость поступления в parking lot;

  • publisher confirm latency и unroutable returns;

  • reconnect и consumer cancellations.

Prefetch ограничивает число сообщений в работе у consumer. Он не должен быть ни бесконечным, ни автоматически равным числу goroutines. Для тяжёлых операций разумно начать с небольшого значения около реальной конкуренции, измерить throughput и память, а затем увеличивать.

Автомасштабирование по глубине queue без ограничения тоже опасно. Если downstream database уже перегружена, запуск ещё ста consumers ускорит только её падение. Ограничение конкуренции, circuit breaker для временной остановки вызовов и rate limit для ограничения их частоты иногда полезнее, чем запуск новой replica.

Повторный поток нужно учитывать как отдельную нагрузку. После восстановления зависимости тысячи сообщений могут одновременно вернуться из одной TTL queue. Несколько ступеней, jitter и контролируемый redrive уменьшают такой retry storm.

Контракт сообщения тоже может сломаться

Producer и consumer развёртываются независимо, поэтому контракт должен переживать смешанные версии.

Полезный envelope выглядит примерно так:

{
  "id": "01K0ORDERCREATED",
  "type": "order.created",
  "schema_version": 2,
  "occurred_at": "2026-07-21T09:00:00Z",
  "producer": "orders",
  "correlation_id": "01K0REQUEST",
  "causation_id": "01K0COMMAND",
  "aggregate": {
    "type": "order",
    "id": "42",
    "version": 7
  },
  "data": {
    "sku": "book-17",
    "count": 1
  }
}

id идентифицирует факт и не меняется при повторе. correlation_id объединяет сообщения одного бизнес-процесса. causation_id указывает, какое сообщение или команда вызвали это событие. Версия агрегата помогает с порядком, а версия схемы - с эволюцией контракта.

Совместимое изменение обычно добавляет необязательное поле, если старые consumers игнорируют неизвестные поля, а значение по умолчанию безопасно. Удаление поля, смена его смысла или переиспользование старого event type под новый контракт ломают consumers. При несовместимой миграции безопаснее некоторое время публиковать новую версию отдельно, перевести consumers и только потом остановить старую.

Не превращайте каждую ошибку десериализации в retry. Пока код не обновился, те же байты останутся теми же байтами. Сообщение должно быстро попасть в parking lot вместе с версией схемы и понятной диагностикой.

Как проверить схему сбоями

Happy path подтверждает только happy path. Retry topology нужно проверять в тех местах, ради которых она появилась.

Consumer упал после COMMIT, но до Ack

Добавьте тестовую точку остановки сразу после reserve и завершите процесс через SIGKILL. После запуска RabbitMQ доставит сообщение снова. Значение stock.reserved должно измениться один раз, а Inbox должен содержать один message_id.

PostgreSQL недоступен

Остановите базу дольше полного retry budget. Сообщение должно последовательно пройти retry.5s, retry.30s и retry.5m, а consumer не должен создавать горячий requeue-цикл. В отдельном запуске восстановите базу до исчерпания попыток и убедитесь, что сообщение обработалось.

Повреждённый JSON

Опубликуйте {not-json} с валидным MessageId. Сообщение должно сразу попасть в inventory.orders.parking, без трёх бессмысленных попыток.

Retry опубликован, но исходный Ack потерян

Остановите consumer после publisher confirm и до delivery.Ack. В рабочей queue останется исходное сообщение, а в retry queue уже будет копия. Обе доставки должны привести к одному резерву.

Binding удалён

Удалите binding retry exchange. Публикация с mandatory=true должна получить basic.return. Consumer не должен подтверждать исходное сообщение. Заодно проверьте, что supervisor не перезапускает consumer без задержки и не создаёт reconnect storm.

Нарушение порядка

Заставьте v7 уйти в retry, а v8 обработайте сразу. Тест должен показать выбранную политику: блокировку shard, ожидание пропущенной версии или безопасное игнорирование устаревшего события. Если тест просто «иногда проходит», политики пока нет.

Parking lot недоступен

Удалите binding или остановите quorum целевой queue. При at-least-once dead lettering исходная quorum queue должна удерживать сообщение и повторять внутреннюю публикацию. Следите за ростом исходной queue: безопасный DLX способен создать backpressure, если цель долго недоступна.

Как выглядят те же задачи в других brokers

Названия механизмов меняются, но границы остаются.

Broker

Повтор и подтверждение

Poison message

Порядок и replay

RabbitMQ

Manual Ack/Nack, delayed retry в quorum queues 4.3 или retry queues с TTL и DLX

Delivery limit и DLX, отдельная parking queue

Порядок queue ослабляют несколько consumers и retries; Streams подходят для replay

Kafka

Consumer управляет offset; retry обычно означает остановку partition или публикацию в retry topic

Отдельный dead-letter topic создаёт приложение или framework

Replay встроен в модель log; порядок только внутри partition

NATS JetStream

AckWait, BackOff, NakWithDelay, MaxDeliver

После MaxDeliver сообщение остаётся в stream, advisory позволяет построить DLQ

Consumer задаёт позицию в stream; порядок зависит от consumer и subjects

Amazon SQS

Visibility timeout, после которого неудалённое сообщение видно снова

Redrive policy с maxReceiveCount перемещает сообщение в DLQ

Standard queue допускает дубликаты и best-effort ordering; FIFO упорядочивает внутри MessageGroupId

NATS JetStream позволяет задать последовательность BackOff и максимальное число доставок в конфигурации consumer. Обычный Nak() вызывает немедленную доставку, а для задержки нужен NakWithDelay(). После MaxDeliver сообщение не исчезает из stream; server публикует advisory, по которому приложение может реализовать DLQ. Эти детали собраны в документации JetStream consumers.

В Amazon SQS consumer не отправляет Ack, а удаляет обработанное сообщение. До удаления оно скрыто на время visibility timeout. Если timeout короче обработки, второй consumer способен получить ту же работу; если слишком длинный, повтор после падения задержится. AWS также подчёркивает, что visibility timeout не отменяет at-least-once delivery. Redrive policy перемещает сообщение в DLQ после maxReceiveCount.

Kafka ближе к журналу, чем к очереди заданий. Consumer может вернуться к старому offset и переиграть историю. Это удобно для исправления ошибочного обработчика, но требует долговечной идемпотентности. События с одним ключом направляют в одну partition, чтобы сохранить их порядок. Если проблемную запись вынести в retry topic, следующие записи продолжат обработку, но строгий порядок по ключу исчезнет.

Выбор broker меняет удобство конкретного механизма. Он не отменяет стабильные ID, dual write, неопределённый результат сетевого вызова и необходимость определить границу бизнес-эффекта.

Итог

Асинхронная коммуникация не устраняет сбои, а позволяет разместить между компонентами устойчивый буфер и выбрать реакцию на каждый сбой. За это приходится явно проектировать владение сообщением.

Надёжная схема выглядит довольно прозаично:

  • producer пишет бизнес-данные и Outbox одной транзакцией;

  • relay публикует событие со стабильным ID и ждёт confirm;

  • consumer меняет свою базу и Inbox одной транзакцией;

  • Ack отправляется после COMMIT;

  • временные ошибки получают ограниченные повторы с задержкой;

  • постоянные ошибки попадают в наблюдаемый parking lot;

  • порядок задаётся по бизнес-ключу, а не предполагается из слова FIFO;

  • backlog, возраст сообщений и redrive имеют операционные ограничения.

Broker закрывает важную часть пути. Остальная надёжность появляется из протокола между приложением, broker и базой. Обычно это не одна эффектная гарантия, а несколько скучных локальных правил. В распределённых системах скучные правила имеют приятное свойство: их можно проверить.