Отправить сообщение в RabbitMQ можно за несколько строк. Сложности начинаются после первого сетевого сбоя: connection нужно восстановить, channel открыть заново, exchange и queue повторно объявить, consumer перезапустить, а неподтверждённое сообщение не потерять и не отправить бесконечно много раз.

RabbitMQ не снимает с приложения ответственность за доставку. Broker отвечает за свою часть пути, а Go-клиент должен правильно работать с подтверждениями, повторными доставками и собственным жизненным циклом.

В этой статье сначала разберём модель RabbitMQ. Затем соберём producer и consumer на официальном клиенте rabbitmq/amqp091-go. В последней части посмотрим на пакет GermanGorelkin/pubsub, который прячет часть повторяющегося кода за более узким API.

Как устроен RabbitMQ

Представим сервис заказов. После создания заказа он публикует событие order.created. Сервис доставки и сервис уведомлений должны получить это событие независимо друг от друга.

Путь сообщения выглядит так:

                         binding: order.*
                      +--------------------> queue: delivery
                      |
producer -> exchange: orders
                      |
                      +--------------------> queue: notifications
                         binding: order.created

В этой схеме producer не пишет сообщение прямо в queue. Он публикует его в exchange. Exchange применяет правила маршрутизации из bindings и передаёт копию сообщения в ноль, одну или несколько queues. Consumer читает сообщения уже из своей queue.

Официальное описание называет эту схему моделью AMQP 0-9-1. AMQP расшифровывается как Advanced Message Queuing Protocol - открытый протокол обмена сообщениями. RabbitMQ поддерживает и другие протоколы, но библиотека amqp091-go работает именно с AMQP 0-9-1.

Producer, consumer и broker

Producer, или publisher, создаёт и отправляет сообщения. Consumer получает их и выполняет полезную работу. Между ними находится broker - сервер RabbitMQ, который принимает сообщения, маршрутизирует их и хранит в queues до обработки.

Producer и consumer не обязаны работать одновременно. Если queue сохраняется при перезапуске, consumer может остановиться, а после запуска продолжить работу с накопившимися сообщениями. В этом состоит одно из главных отличий от обычного синхронного HTTP-вызова.

Сообщение состоит из бинарного тела и свойств. В свойствах можно передать тип содержимого, идентификатор сообщения, время создания, срок жизни, correlation ID и пользовательские headers. RabbitMQ не навязывает схему тела: JSON, Protocol Buffers или произвольные байты остаются договорённостью приложений.

Exchange определяет маршрут

Тип exchange задаёт способ выбора queues:

Тип Как маршрутизирует Подходящий пример
direct Сравнивает routing key с ключом binding целиком Команда invoice.generate для конкретного consumer
fanout Отправляет копию во все связанные queues, routing key не учитывает Рассылка события всем заинтересованным сервисам
topic Сопоставляет routing key с шаблоном, где * заменяет одно слово, а # - несколько order.created, order.paid, order.*
headers Проверяет значения headers вместо routing key Маршрутизация по нескольким независимым признакам

Exchange может направить сообщение сразу в несколько queues. Поэтому pub/sub в RabbitMQ строится не на общей queue для всех подписчиков, а на отдельной queue для каждого независимого получателя. Если сервисы доставки и уведомлений читают одну queue, broker будет считать их competing consumers. Каждый сервис увидит только часть событий.

Правило простое: несколько экземпляров одного consumer обычно читают одну queue, а разные бизнес-получатели используют разные queues.

Queue хранит сообщение до обработки

Queue - это упорядоченный буфер сообщений. Несколько consumers одной queue по умолчанию получают сообщения по очереди. Такой режим подходит для фоновых задач: можно добавить workers и увеличить параллелизм, не меняя producer.

Однако порядок обработки не равен порядку публикации при любых условиях. Несколько consumers работают с разной скоростью. Сообщение может вернуться в queue после ошибки, а приоритеты меняют порядок выдачи. Если бизнесу нужен строгий порядок для одного заказа, обычно оставляют один активный consumer или разбивают поток по ключу и сохраняют порядок только внутри каждой группы.

Queue, exchange и binding вместе образуют топологию RabbitMQ. Приложение может объявлять её при запуске. Повторное объявление с теми же параметрами безопасно. Если объявить существующую queue с другими параметрами, RabbitMQ закроет channel с ошибкой. Такая проверка конфигурации полезна, хотя при первом развёртывании иногда выглядит неожиданно.

Connection, channel и vhost

Connection - это долгоживущее TCP-соединение с RabbitMQ. Его установка включает аутентификацию и согласование параметров, поэтому открывать новый connection для каждого сообщения дорого.

Внутри одного connection работают channels. Это лёгкие логические соединения, которые разделяют один TCP-сокет. Объявление топологии, публикация и чтение выполняются через channel. Ошибка протокола может закрыть только один channel, тогда как разрыв TCP-соединения закроет все channels.

В Go важно учитывать конкурентный доступ. Документация amqp091-go не рекомендует вызывать методы одного Channel из нескольких goroutines. На практике публикацию и обработку сообщений удобнее разделить по channels. Вызовы общего publisher channel можно сериализовать или выполнять из одной выделенной goroutine.

Virtual host, или vhost, создаёт изолированное пространство имён для пользователей, exchanges, queues и прав доступа. Это не отдельный процесс RabbitMQ, а логическая граница внутри broker. Обычно окружения или независимые приложения получают разные vhosts и отдельные учётные данные.

Что RabbitMQ гарантирует, а что нет

Фраза «сообщение положили в надёжную queue» скрывает несколько независимых настроек.

Чтобы после перезапуска broker сохранились и сообщение, и маршрут для новых публикаций, обычно нужны:

  1. exchange с флагом durable, который сохраняется после перезапуска;
  2. durable queue и её привязка;
  3. сообщение с DeliveryMode: Persistent.

Но даже с этой комбинацией producer не знает, принял ли broker конкретное сообщение. TCP-соединение могло разорваться в момент отправки. Для этой границы существуют publisher confirms: broker отвечает, когда принимает ответственность за сообщение. Для persistent-сообщения в durable queue подтверждение связано с записью на диск, а для реплицируемой quorum queue - с принятием сообщения большинством реплик. Подробные условия описаны в документации RabbitMQ.

На другой стороне работает consumer acknowledgement. Consumer отправляет Ack только после успешного выполнения своей работы. Если connection закроется до подтверждения, RabbitMQ вернёт сообщение в queue и позднее доставит его снова.

Publisher confirm и consumer acknowledgement не образуют одну сквозную транзакцию. Первое подтверждает участок producer -> broker, второе - участок broker -> consumer. Между ними остаётся бизнес-операция: например, запись в PostgreSQL.

На практике это означает доставку at least once, или как минимум один раз: сообщение не должно потеряться, но может прийти повторно. Поэтому consumer должен быть идемпотентным. Один из вариантов - передавать MessageId и сохранять идентификаторы уже обработанных сообщений в базе вместе с бизнес-изменением.

Есть ещё одна граница. Publisher confirm означает, что exchange принял публикацию, но не обязательно нашёл для неё queue. Если producer использует mandatory=false, RabbitMQ может отбросить сообщение без подходящего маршрута и всё равно подтвердить публикацию. Для важных данных стоит использовать mandatory=true, обрабатывать basic.return и следить за числом unroutable messages, то есть сообщений без подходящей queue. Это отдельно подчёркнуто в руководстве для publishers.

Повторная доставка требует политики

После временной ошибки хочется вызвать Nack с requeue=true. Если причина не исчезает, consumer сразу получает то же сообщение, снова падает и создаёт горячий цикл. Один неисправимый JSON способен надолго занять worker - весьма эффективный способ превратить queue в карусель.

Обычно ошибки разделяют как минимум на две группы:

  • временные: недоступна база, сработал rate limit, истёк сетевой timeout;
  • постоянные: неверная схема, отсутствует обязательное поле, операция запрещена бизнес-правилом.

Для временных ошибок задают ограниченное число повторов с задержкой. В RabbitMQ такую схему часто строят через retry queues и TTL (time to live) - срок хранения сообщения. Сообщения, которые не удалось обработать, можно направить в dead letter exchange (DLX) - отдельный exchange для ошибок. Постоянные ошибки отправляют в отдельную queue для разбора или подтверждают после сохранения диагностической информации. Бесконечный немедленный requeue почти никогда не заменяет продуманную политику повторов.

Когда RabbitMQ подходит

RabbitMQ особенно полезен в четырёх сценариях.

Фоновые задачи

HTTP-сервис быстро принимает запрос и ставит тяжёлую работу в queue. Несколько workers разбирают изображения, формируют отчёты или отправляют письма. Queue сглаживает всплеск нагрузки, а prefetch ограничивает число сообщений, которые broker отдаёт занятому consumer без подтверждения.

События между сервисами

Один сервис публикует факт order.created, а несколько систем реагируют независимо. Exchange типа topic или fanout позволяет добавлять новых получателей, не меняя producer. Такая схема ослабляет временную связанность, но не отменяет версионирование контрактов и наблюдаемость.

Буфер перед медленной системой

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

Запрос-ответ (request/reply) с особыми требованиями

RabbitMQ поддерживает request/reply через CorrelationId и ReplyTo; этому посвящён отдельный tutorial. Подход бывает полезен внутри уже существующей инфраструктуры обмена сообщениями. Для обычного синхронного запроса HTTP или gRPC чаще проще: у них понятнее ограничение времени, балансировка и диагностика полного пути.

RabbitMQ не всегда лучший выбор. Обычные queues удаляют подтверждённые сообщения и плохо подходят на роль долговременного журнала с replay. Для повторного чтения истории, очень больших queues и высокой потоковой нагрузки у RabbitMQ есть отдельная модель Streams, а в других системах может лучше подойти Kafka. Если два компонента всегда должны ответить друг другу синхронно, broker способен добавить больше состояний отказа, чем пользы.

Плюсы и минусы RabbitMQ

Плюсы Минусы
Гибкая маршрутизация через exchanges и bindings Топологию и правила доставки нужно проектировать и сопровождать
Publisher confirms и consumer acknowledgements At least once требует идемпотентных consumers
Сглаживание всплесков и горизонтальное масштабирование workers Растущая queue увеличивает задержку и расход диска
Quorum queues, TLS-шифрование, права доступа, веб-интерфейс управления и метрики Кластер остаётся отдельной распределённой системой, которую нужно обновлять и наблюдать
Несколько протоколов и зрелая экосистема клиентов Обычная queue не даёт удобного replay после подтверждения
TTL, dead lettering, priorities и другие политики Большое число возможностей увеличивает число неверных комбинаций

RabbitMQ в Go через amqp091-go

amqp091-go - официальный Go-клиент, который поддерживает команда RabbitMQ. Его API близко отображает модель AMQP 0-9-1: приложение явно открывает connection и channel, объявляет топологию, включает publisher confirms, настраивает Quality of Service (QoS) и отправляет consumer acknowledgements.

Установить пакет можно обычной командой:

go get github.com/rabbitmq/amqp091-go

Для начала создадим channel для событий заказов:

func prepareChannel(conn *amqp.Connection) (*amqp.Channel, error) {
	ch, err := conn.Channel()
	if err != nil {
		return nil, fmt.Errorf("open channel: %w", err)
	}

	closeOnError := func(err error) (*amqp.Channel, error) {
		_ = ch.Close()
		return nil, err
	}

	if err := ch.ExchangeDeclare(
		"orders", // name
		"topic",  // kind
		true,     // durable
		false,    // auto-delete
		false,    // internal
		false,    // no-wait
		nil,      // arguments
	); err != nil {
		return closeOnError(fmt.Errorf("declare exchange: %w", err))
	}

	if _, err := ch.QueueDeclare(
		"delivery.orders", // name
		true,              // durable
		false,             // auto-delete
		false,             // exclusive
		false,             // no-wait
		nil,               // arguments
	); err != nil {
		return closeOnError(fmt.Errorf("declare queue: %w", err))
	}

	if err := ch.QueueBind(
		"delivery.orders",
		"order.*",
		"orders",
		false,
		nil,
	); err != nil {
		return closeOnError(fmt.Errorf("bind queue: %w", err))
	}

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

	if err := ch.Confirm(false); err != nil {
		return closeOnError(fmt.Errorf("enable confirms: %w", err))
	}

	return ch, nil
}

Код многословен, но каждый аргумент имеет протокольный смысл. Можно включить auto-delete для временной топологии, передать x-queue-type: quorum, задать dead lettering или настроить prefetch под конкретную нагрузку. Низкоуровневый API почти ничего не решает за приложение.

Публикация с подтверждением

Обычный успешный возврат из PublishWithContext означает, что клиент записал фреймы AMQP в connection. Для важных сообщений этого недостаточно. Publisher confirms позволяют связать публикацию с последующим ответом от broker:

func publishOrderCreated(
	ctx context.Context,
	ch *amqp.Channel,
	messageID string,
	body []byte,
) error {
	confirmation, err := ch.PublishWithDeferredConfirmWithContext(
		ctx,
		"orders",
		"order.created",
		true,  // mandatory: вернуть сообщение, если нет подходящей queue
		false, // immediate не поддерживается RabbitMQ
		amqp.Publishing{
			ContentType:  "application/json",
			DeliveryMode: amqp.Persistent,
			MessageId:    messageID,
			Timestamp:    time.Now(),
			Body:         body,
		},
	)
	if err != nil {
		return fmt.Errorf("publish order.created: %w", err)
	}

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

	return nil
}

Этот пример всё ещё неполон: при mandatory=true приложению нужно постоянно читать channel из NotifyReturn. Confirm и return приходят как разные асинхронные события. Кроме того, после timeout нельзя достоверно сказать, принял ли broker сообщение. Повторная публикация может создать дубликат, поэтому messageID и идемпотентность нужны даже при publisher confirms.

Для высокой пропускной способности не стоит ждать confirm после каждого сообщения. Клиент позволяет отправить несколько сообщений, не дожидаясь ответа на каждое, и сопоставить confirms по порядковому номеру delivery tag. Это быстрее, но требует отдельного учёта неподтверждённых сообщений и аккуратной синхронизации.

Consumer и ручной Ack

При чтении важных сообщений обычно отключают auto-ack:

func consume(
	ctx context.Context,
	ch *amqp.Channel,
	handle func(context.Context, []byte) error,
) error {
	deliveries, err := ch.Consume(
		"delivery.orders",
		"delivery-worker",
		false, // auto-ack
		false, // exclusive
		false, // no-local
		false, // no-wait
		nil,
	)
	if err != nil {
		return fmt.Errorf("start consumer: %w", err)
	}

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

			if err := handle(ctx, delivery.Body); err != nil {
				// requeue=false: RabbitMQ отправит сообщение в настроенный
				// dead letter exchange или удалит его.
				if nackErr := delivery.Nack(false, false); nackErr != nil {
					return fmt.Errorf("nack message: %w", nackErr)
				}
				continue
			}

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

В production handle должен различать временные и постоянные ошибки. После временной ошибки сообщение можно отправить в топологию повторов с задержкой. После постоянной - в dead letter queue. Один Nack на все случаи редко выражает бизнес-требования полностью.

Особенности amqp091-go в production

Прямой клиент даёт полный контроль, но оставляет приложению несколько обязанностей:

  • не создавать connection для каждой публикации;
  • выбрать схему channels и не использовать один channel конкурентно без собственной сериализации;
  • постоянно читать NotifyClose, NotifyReturn, NotifyPublish и другие включённые channels уведомлений;
  • повторно открыть channel после протокольной ошибки;
  • сопоставить подтверждения публикации с конкретными сообщениями;
  • перезапустить consumers и восстановить топологию после сбоя;
  • завершить consumers, дождаться текущей работы и закрыть channels и connection;
  • измерять число reconnect, неподтверждённых публикаций, redelivery и ошибок consumers.

Здесь есть важная поправка для свежих версий. В amqp091-go 1.12 появился recovery для connections и channels, а версия 1.13 добавила recovery для exchanges, queues, bindings и consumers. Функция включается через Config.Recovery:

conn, err := amqp.DialConfig(addr, amqp.Config{
	Recovery: &amqp.Recovery{},
})

На момент написания статьи документация помечает эту возможность как экспериментальную. Recovery нужно включить явно, наблюдать его состояния и проверить на собственных сценариях: обрыв TCP, перезапуск узла, удаление queue, закрытие channel из-за несовпавшей декларации и остановка процесса во время recovery.

Сильная сторона amqp091-go - точное управление AMQP и доступ к новым возможностям RabbitMQ. Обратная сторона того же решения - протокольные детали проникают в инфраструктурный код сервиса. Если нескольким приложениям нужна одна и та же ограниченная модель работы, появляется смысл вынести её в общий компонент.

Более узкий API через pubsub

Пакет GermanGorelkin/pubsub построен поверх amqp091-go. Он вводит одну сущность Session, которая владеет connection, channel, объявлениями ресурсов и consumers.

Установить пакет можно так:

go get github.com/germangorelkin/pubsub@v1.0.0

Версия 1.0.0 требует Go 1.24 или новее.

Тот же базовый сценарий выглядит компактнее:

ctx, stop := signal.NotifyContext(
	context.Background(),
	os.Interrupt,
	syscall.SIGTERM,
)
defer stop()

session := pubsub.New(
	os.Getenv("AMQP_URL"),
	pubsub.WithDeclare(
		pubsub.Exchange{
			Name:           "orders",
			Kind:           "topic",
			Durable:        true,
			IsUsageDefault: true,
		},
		pubsub.Queue{
			Name:           "delivery.orders",
			Durable:        true,
			IsUsageDefault: true,
		},
		pubsub.Bind{
			QueueName:      "delivery.orders",
			ExchangeName:   "orders",
			Key:            "order.created",
			IsUsageDefault: true,
		},
	),
)

readyCtx, cancelReady := context.WithTimeout(ctx, 15*time.Second)
if err := session.WaitReady(readyCtx); err != nil {
	cancelReady()
	log.Fatalf("RabbitMQ is not ready: %v", err)
}
cancelReady()

if err := session.Subscribe(func(delivery pubsub.Delivery) error {
	log.Printf("received %s: %s", delivery.RoutingKey, delivery.Body)
	return nil
}); err != nil {
	log.Fatalf("subscribe: %v", err)
}

publishCtx, cancelPublish := context.WithTimeout(ctx, 5*time.Second)
err := session.PublishWithContext(
	publishCtx,
	[]byte(`{"order_id":"42"}`),
)
cancelPublish()
if err != nil {
	log.Printf("publish: %v", err)
}

<-ctx.Done()

if err := session.Close(); err != nil {
	log.Printf("close RabbitMQ session: %v", err)
}

В этом примере нет явных Dial, Channel, Confirm, NotifyClose, Ack и Nack. Эти операции не исчезли, а стали частью поведения Session.

Что делает Session

После New пакет запускает фоновые goroutines. Они устанавливают connection, открывают channel и запускают зарегистрированных consumers. Если TCP-соединение теряется, Session подключается снова. Если протокольная ошибка закрывает channel, пакет создаёт новый, не разрывая исправный connection.

При каждой инициализации channel пакет снова объявляет зарегистрированные exchanges, queues и bindings. Consumers хранятся внутри Session и запускаются заново, когда channel готов. Это закрывает неприятный класс ошибок: connection уже восстановился, но приложение больше не читает сообщения.

Ресурсы можно объявлять и после создания Session. Пакет запоминает их и снова объявляет после следующего подключения. Через WithLogger можно подключить любой логгер с методом Printf. Параметры WithReconnectDelay, WithReInitDelay и WithResendDelay задают задержки между попытками подключиться, открыть channel и отправить сообщение.

WaitReady(ctx) отделяет создание объекта от готовности инфраструктуры. New возвращает Session сразу и не сообщает ошибкой, что broker недоступен. Поэтому сервис, который не может работать без RabbitMQ, должен явно вызвать WaitReady. Если RabbitMQ не критичен при старте, приложение может запуститься и дождаться готового connection позже.

Подтверждённая публикация как поведение по умолчанию

PublishWithContext ждёт publisher confirm. При ошибке отправки, Nack или закрытии confirm channel пакет делает retry после настраиваемой задержки. Контекст ограничивает ожидание готовности, ожидание своей очереди на публикацию и ожидание confirm.

Такой API упрощает прикладной код: nil означает, что broker подтвердил приём сообщения, а timeout или остановка Session возвращаются как ошибка. Для некритичных метрик и логов есть UnsafePublishToWithContext, который не ждёт confirm.

Есть и более узкая граница гарантии. Текущая версия pubsub публикует с mandatory=false и не открывает приложению basic.return. Поэтому успешный PublishWithContext подтверждает приём публикации exchange, но не доказывает, что сообщение попало хотя бы в одну queue. Простую статическую топологию можно контролировать тестами и мониторингом. Если приложение должно обработать каждое сообщение без подходящего маршрута, нужен прямой API amqp091-go.

За простоту приходится платить пропускной способностью. Текущая реализация выполняет публикации с confirms последовательно и ждёт ответа на каждую. Это разумное поведение по умолчанию для небольшого потока важных событий, но не для producer, которому нужны тысячи параллельных публикаций и batch confirms.

Ack и Nack через результат handler

Handler получает упрощённый pubsub.Delivery с Exchange, RoutingKey и Body. Если функция возвращает nil, пакет вызывает Ack. Если она возвращает ошибку или panic, пакет вызывает Nack с requeue=true.

err := session.SubscribeTo("delivery.orders", func(d pubsub.Delivery) error {
	var event OrderCreated
	if err := json.Unmarshal(d.Body, &event); err != nil {
		// Повреждённое сообщение не исправится после повторной доставки.
		// Сохраняем его для разбора и возвращаем nil, чтобы отправить Ack.
		if saveErr := saveRejectedEvent(d.Body, err); saveErr != nil {
			return fmt.Errorf("save rejected order.created: %w", saveErr)
		}
		return nil
	}

	if err := createDelivery(event.OrderID); err != nil {
		return fmt.Errorf("create delivery: %w", err)
	}

	return nil
})

Контракт легко запомнить, но у него есть важное следствие: неисправимое сообщение будет возвращаться в queue снова и снова. В handler нужно заранее решить, какие ошибки действительно стоит повторять. Если нужна отдельная команда Nack(false, false), dead lettering по числу попыток или разные стратегии для разных ошибок, низкоуровневый amqp091-go даёт больше контроля.

pubsub.Delivery не открывает MessageId, headers и признак Redelivered. Поэтому идентификатор для защиты от дубликатов следует включить в схему тела, например как поле event_id, а затем сохранять вместе с бизнес-изменением.

Пакет также поддерживает auto-ack для данных, которые допустимо потерять. В этом режиме RabbitMQ считает сообщение обработанным сразу после доставки. Ошибка или panic в handler уже не вернут его в queue.

Управление consumers и завершение работы

AddConsumer регистрирует consumer с явным или автоматически созданным тегом. RemoveConsumer останавливает его через AMQP basic.cancel. Это удобно, если подписки зависят от конфигурации, которая меняется во время работы, или от жизненного цикла компонента.

Close прекращает фоновые циклы, ждёт consumers, которые уже выполняют работу, затем закрывает channel и connection. Это лучше, чем дать процессу оборвать TCP-сокет. Однако у Close нет context или timeout. Handler, который завис навсегда, задержит остановку всего процесса. Timeout по-прежнему нужно задавать внутри бизнес-операций.

Где абстракция помогает, а где мешает

pubsub не пытается заменить весь API RabbitMQ. Он выбирает конкретный профиль использования:

  • одна управляемая Session;
  • автоматический reconnect и повторное объявление простой топологии;
  • persistent-сообщения;
  • publisher confirm для обычной публикации;
  • prefetch, равный 1;
  • последовательная обработка внутри одного consumer;
  • Ack при nil и Nack(requeue=true) при ошибке;
  • несколько настраиваемых задержек для reconnect, повторного открытия channel и отправки.

Такой профиль хорошо подходит небольшим и средним сервисам с несколькими JSON-событиями, умеренной нагрузкой и одинаковыми правилами доставки. Команда получает меньше кода для управления жизненным циклом и один узнаваемый способ подключать RabbitMQ в разных приложениях.

Граница становится заметна, когда нужны аргументы декларации queue, настраиваемый prefetch, headers, MessageId, TTL, приоритеты, mandatory=true, обработка returns, разные виды Nack, несколько channels или конвейер publisher confirms. Текущий API pubsub не открывает эти детали. В таком сервисе попытка обойти абстракцию быстро станет сложнее прямого использования amqp091-go.

Отдельно стоит учесть свежие изменения официального клиента. В go.mod пакета pubsub 1.0.0 указана зависимость amqp091-go 1.10.0, а recovery реализован внутри Session. В amqp091-go 1.12 и 1.13 похожие базовые возможности появились в самом клиенте. Поэтому выбор уже нельзя свести к вопросу «кто умеет делать reconnect». Нужно сравнивать полный контракт: как публикуются сообщения, как consumer управляет Ack и Nack, какие свойства доступны и сколько протокольных деталей должно видеть приложение.

Задача amqp091-go pubsub
Уровень API Почти прямое отображение AMQP 0-9-1 Узкий API с заранее выбранным сценарием использования
Automatic recovery Есть в свежих версиях, пока помечено как экспериментальное Собственные reconnect, повторное открытие channel и перезапуск consumers
Объявление топологии Полный набор параметров и arguments Простые Exchange, Queue и Bind
Publisher confirms Несколько API, включая асинхронные и deferred confirms Ожидание подтверждения и повтор внутри Publish*
Подтверждение обработки Приложение явно выбирает Ack, Nack или Reject nil вызывает Ack, ошибка вызывает Nack с requeue
Свойства сообщения Headers, ID, timestamp, TTL, correlation ID и другие Тело, exchange и routing key
Пропускная способность Можно построить конвейер и групповые подтверждения Подтверждаемые публикации выполняются последовательно
Объём инфраструктурного кода Больше Меньше
Контроль над протоколом Максимальный Намеренно ограниченный

Практический выбор

Начинать стоит не с названия библиотеки, а с требуемого поведения.

Выбирайте прямой amqp091-go, если сервису нужны сложная топология, максимальная пропускная способность, несколько независимых channels, богатые свойства сообщений или собственная политика повторов. Цена - больше кода и тестов на сбои.

Рассмотрите pubsub, если сервис использует простую топологию из exchange, queue и binding, публикует важные сообщения с умеренной скоростью и может выразить результат обработки как «успех» или «повторить». В этом случае пакет убирает достаточно кода управления connection и consumers, чтобы бизнес-логика стала заметнее.

Независимо от клиента, перед запуском в production полезно проверить один и тот же набор сценариев:

  1. broker недоступен во время старта;
  2. TCP-соединение обрывается во время публикации;
  3. channel закрывается из-за ошибочной декларации;
  4. сообщение не подходит ни под один binding;
  5. consumer падает до Ack;
  6. одно сообщение доставляется дважды;
  7. consumer постоянно возвращает ошибку;
  8. процесс получает SIGTERM во время обработки;
  9. queue растёт быстрее, чем consumers успевают её разбирать.

Главный вывод не зависит от выбранной обёртки: надёжность RabbitMQ складывается из поведения broker, producer, consumer и бизнес-хранилища. Клиентская библиотека может сделать правильный путь короче, но не может сама выбрать политику повторов, идемпотентность и допустимую потерю данных.

Дополнительные материалы


Комментарии в Telegram-группе!