RabbitMQ — популярный брокер сообщений. Работать с ним из Java можно напрямую через его Java-клиент, но это много низкоуровневого кода: управление соединениями, потоками, ack/nack вручную. Spring AMQP (библиотека spring-rabbit) берёт этот код на себя и даёт удобный Spring-способ отправлять и получать сообщения.
Отправить и получить — самая простая часть. Главное начинается там, где обработчик упал: по умолчанию сообщение возвращается в очередь и тут же приходит снова. У классической очереди так может продолжаться бесконечно; у надёжной (quorum) с RabbitMQ 4.0 брокер оборвёт круг сам после двадцатой доставки — и выбросит сообщение, если DLX не настроен. Весь рабочий рецепт ниже — про то, как разорвать круг раньше и ничего при этом не потерять.
Без настройки упавший обработчик возвращает сообщение в ту же очередь, и оно приходит снова: у классической очереди круг не кончается, у надёжной брокер обрывает его на двадцатой доставке и может выбросить сообщение. Рабочий набор: setDefaultRequeueRejected(false), retry на 4 попытки (полосы в одном масштабе: 1с, 2с, 4с — семь секунд до отказа) и DLX. После четвёртой сообщение отклоняется без возврата и уходит в orders.dlq с заголовком x-death, а основная очередь берёт следующее.
Подключение
Одна зависимость в Gradle добавляет всё нужное:
// build.gradle.kts
dependencies {
implementation("org.springframework.boot:spring-boot-starter-amqp")
}
Базовая конфигурация в application.yml:
spring:
rabbitmq:
host: rabbit
port: 5672
username: billing-service
password: ${RABBIT_PASSWORD}
virtual-host: /prod
Spring Boot автоматически создаст соединение и настроит все нужные бины.
Как отправить сообщение: RabbitTemplate
Раньше, чтобы опубликовать сообщение, нужно было вручную открывать канал, сериализовать объект, вызывать методы протокола AMQP. Spring AMQP делает это одним вызовом.
Главный инструмент для отправки — RabbitTemplate. Его создаёт Spring Boot автоматически, достаточно внедрить как зависимость:
@Component
@RequiredArgsConstructor
public class OrderEventPublisher {
private final RabbitTemplate rabbit;
public void publish(OrderCreatedEvent event) {
rabbit.convertAndSend("orders", "order.created", event);
}
}
convertAndSend(exchange, routingKey, payload) — основной метод. Три параметра:
"orders"— имя exchange (куда шлём);"order.created"— routing key (по нему exchange решает, в какую очередь направить);event— объект, который станет телом сообщения.
Зачем менять сериализацию
По умолчанию Spring AMQP сериализует объект через Java-сериализацию. Это работает, но у такого подхода есть проблема: другой сервис (или другой язык) не сможет прочитать сообщение. Формат Java-сериализации понимает только Java.
Решение — переключиться на JSON. Для этого нужно объявить конвертер и указать его шаблону:
@Configuration
public class RabbitConfig {
@Bean
MessageConverter jacksonConverter(ObjectMapper mapper) {
return new Jackson2JsonMessageConverter(mapper);
}
@Bean
RabbitTemplate rabbitTemplate(ConnectionFactory cf, MessageConverter converter) {
var t = new RabbitTemplate(cf);
t.setMessageConverter(converter);
return t;
}
}
После этого payload уходит в виде JSON. Spring автоматически добавит заголовок content-type: application/json — получатель поймёт, как десериализовать.
Конвертер решает формат тела: по умолчанию это Java-сериализация, понятная только Java, а Jackson кладёт JSON и проставляет content-type и __TypeId__, по которому получатель на Java собирает нужный класс.
Одна оговорка про имя класса. Jackson2JsonMessageConverter — это spring-amqp 3.x, то есть Spring Boot 3. В spring-amqp 4.0 библиотека перешла на Jackson 3: конвертер там называется JacksonJsonMessageConverter, а прежнее имя объявлено устаревшим и тянет за собой Jackson 2, которого в зависимостях свежего Boot уже нет. Так что сверяйтесь с версией: имя класса — единственное, что меняется, смысл прежний.
Половина гарантии на стороне отправителя: confirms и returns
Всё, что описано выше, защищает сообщение после того, как оно попало в брокер. Но rabbitTemplate.convertAndSend(...) возвращает управление сразу и ничего не говорит о судьбе сообщения: брокер мог быть недоступен, диск мог не успеть, ключ мог не совпасть ни с одной привязкой. Метод при этом не бросит исключение — он просто вернётся.
Закрывают это две настройки, и включают их вместе.
spring:
rabbitmq:
publisher-confirm-type: correlated
publisher-returns: true
Подтверждения (confirms) отвечают на вопрос «брокер принял сообщение». Режим correlated присылает подтверждение вместе с объектом связи, по которому видно, о каком именно сообщении речь:
rabbit.setConfirmCallback((correlation, ack, cause) -> {
if (!ack) {
log.error("Брокер не принял сообщение {}: {}", correlation.getId(), cause);
outbox.markFailed(correlation.getId());
}
});
rabbit.convertAndSend("orders.exchange", "order.created", event,
new CorrelationData(event.id().toString()));
Обратный вызов приходит асинхронно, уже после возврата из convertAndSend, и это главное, что нужно про него понять: писать код в расчёте на «отправили и сразу узнали» нельзя. Идентификатор в CorrelationData поэтому должен быть тем, по которому вы найдёте своё сообщение потом, — идентификатор записи в таблице исходящих, а не случайная строка. Если нужен синхронный ответ, у шаблона есть waitForConfirms внутри invoke, но он останавливает поток и годится для пакетной отправки, а не для обработки запроса.
Возвраты (returns) отвечают на второй вопрос — «нашлась ли очередь». Брокер принял сообщение, не нашёл ни одной подходящей привязки и вернул его обратно. Чтобы это случилось, нужен флаг mandatory и обработчик возврата:
rabbit.setMandatory(true);
rabbit.setReturnsCallback(returned ->
log.error("Некуда доставить: ключ {}, причина {}",
returned.getRoutingKey(), returned.getReplyText()));
Без него сообщение с неверным ключом исчезает молча, а подтверждение при этом приходит положительное — брокер же его принял. Разбор обоих механизмов на уровне протокола — в статье про AMQP.
Что делать в этих обработчиках. Писать в лог недостаточно: запись в логе не отправит сообщение заново. Рабочая схема — таблица исходящих сообщений: перед отправкой запись создаётся в той же транзакции, что бизнес-операция, при положительном подтверждении помечается отправленной, при отрицательном или возврате остаётся неотправленной, и фоновый процесс пробует снова. Тогда недоступность брокера превращается в задержку, а не в потерю; сам шаблон разбирает статья про шаблоны обмена.
Как получить сообщение: @RabbitListener
Подписаться на очередь проще всего через аннотацию @RabbitListener:
@Component
public class OrderEventListener {
@RabbitListener(queues = "orders.fulfillment", concurrency = "3-10")
public void handle(OrderCreatedEvent event) {
// обрабатываем событие
}
}
Spring запустит пул потоков, которые будут читать из очереди и вызывать метод. Параметр concurrency = "3-10" задаёт диапазон: минимум 3 потока постоянно, и до 10 при увеличении нагрузки.
Объект OrderCreatedEvent Spring десериализует автоматически из JSON, если настроен MessageConverter.
Как объявить очередь, exchange и binding
Обычно очередь создаётся вручную в RabbitMQ или через Infrastructure as Code. Но Spring AMQP умеет создавать их при старте приложения — это удобно для разработки и простых сетапов:
@Configuration
public class RabbitTopology {
@Bean
Queue ordersFulfillment() {
return QueueBuilder
.durable("orders.fulfillment")
.quorum()
.withArgument("x-dead-letter-exchange", "orders.dlx")
.build();
}
@Bean
DirectExchange ordersExchange() {
return new DirectExchange("orders", true, false);
}
@Bean
Binding binding(Queue ordersFulfillment, DirectExchange ordersExchange) {
return BindingBuilder
.bind(ordersFulfillment)
.to(ordersExchange)
.with("order.created");
}
}
Spring видит эти бины и при старте проверит, существуют ли они в брокере. Если нет — создаст.
Заголовки в сообщениях
Технические метаданные (идентификатор запроса, версия события) удобно передавать через заголовки сообщения, а не смешивать с бизнес-данными в теле:
// При отправке
public void publish(OrderEvent event, String correlationId) {
var message = MessageBuilder
.withBody(jsonOf(event))
.setContentType("application/json")
.setHeader("X-Correlation-ID", correlationId)
.setHeader("X-Event-Version", "2")
.build();
rabbit.send("orders", "order.created", message);
}
// При получении
@RabbitListener(queues = "orders.fulfillment")
public void handle(
@Payload OrderEvent event,
@Header("X-Correlation-ID") String correlationId) {
MDC.put("correlationId", correlationId);
// ...
}
Что происходит при ошибке
По умолчанию, если ваш метод-listener бросает исключение, Spring возвращает сообщение обратно в очередь — и оно тут же приходит снова. Сообщение обрабатывается, падает, возвращается, снова обрабатывается.
У классической очереди это буквально бесконечный цикл. У надёжной очереди, которую мы объявили выше через .quorum(), с RabbitMQ 4.0 есть встроенный ограничитель x-delivery-limit со значением 20 по умолчанию: на двадцать первой доставке брокер сам снимет сообщение — отправит в Dead Letter Exchange, если тот настроен, и молча выбросит, если нет. Оба исхода плохие: двадцать бесполезных прогонов подряд, а в конце ещё и потеря данных.
Чтобы этого избежать, настраивают два механизма вместе: retry (повторные попытки) и Dead Letter Exchange (место, куда уходят сообщения после всех неудачных попыток).
Повторные попытки
Если ошибка временная (база данных на секунду недоступна, внешний сервис ответил 503), имеет смысл попробовать несколько раз с паузой между попытками:
@Bean
RetryOperationsInterceptor retryInterceptor() {
return RetryInterceptorBuilder.stateless()
.maxAttempts(4)
.backOffOptions(1000L, 2.0, 30000L) // 1с, потом 2с, потом 4с; верхняя граница — 30с
.recoverer(new RejectAndDontRequeueRecoverer())
.build();
}
@Bean
SimpleRabbitListenerContainerFactory rabbitListenerContainerFactory(
SimpleRabbitListenerContainerFactoryConfigurer configurer,
ConnectionFactory cf, MessageConverter converter) {
var factory = new SimpleRabbitListenerContainerFactory();
configurer.configure(factory, cf);
factory.setMessageConverter(converter);
factory.setAdviceChain(retryInterceptor());
factory.setDefaultRequeueRejected(false);
return factory;
}
Что здесь происходит:
- Spring сделает 4 попытки с экспоненциальной паузой (1с, 2с, 4с).
- После 4-й неудачной попытки
RejectAndDontRequeueRecovererскажет брокеру «не возвращай это сообщение» — и оно уйдёт в Dead Letter Exchange. setDefaultRequeueRejected(false)— на случай исключений, которые retry не перехватил.configurer.configure(factory, cf)— строчка, которую легко пропустить и потом долго искать. Свой бин с этим именем заменяет собой тот, что собрал Spring Boot, и без вызова конфигуратора все настройкиspring.rabbitmq.listener.*— prefetch, concurrency, режим подтверждения — молча перестают действовать. Конфигуратор сначала накатывает их на фабрику, а уже потом мы дописываем своё.
Dead Letter Exchange
Dead Letter Exchange (DLX) — это специальный exchange в RabbitMQ, куда попадают «мёртвые» сообщения: те, которые не удалось обработать. Оттуда их можно проанализировать, отправить уведомление или попробовать обработать по-другому.
Настройка DLX на уровне очереди:
@Bean
Queue ordersMain() {
return QueueBuilder
.durable("orders.fulfillment")
.quorum()
.withArgument("x-dead-letter-exchange", "orders.dlx")
.withArgument("x-dead-letter-routing-key", "order.failed")
.build();
}
@Bean DirectExchange ordersDlx() {
return new DirectExchange("orders.dlx", true, false);
}
@Bean
Queue ordersDlq() {
return QueueBuilder.durable("orders.dlq").quorum().build();
}
@Bean
Binding dlqBinding(Queue ordersDlq, DirectExchange ordersDlx) {
return BindingBuilder.bind(ordersDlq).to(ordersDlx).with("order.failed");
}
Поток при ошибке: основная очередь → исключение → reject без возврата → DLX → очередь orders.dlq.
В очереди orders.dlq обычно сидит отдельный listener, который логирует проблему, сохраняет сообщение в базу для ручного разбора или отправляет алерт:
@RabbitListener(queues = "orders.dlq")
public void inspect(
OrderCreatedEvent event,
@Header("x-death") List<Map<String, Object>> deaths) {
long attempts = (Long) deaths.get(0).get("count"); // именно count, а не размер списка
log.error("Не удалось обработать orderId={}, попыток={}", event.orderId(), attempts);
}
Заголовок x-death добавляет брокер. Внутри — список записей, по одной на каждую пару «очередь и причина отклонения», а число попыток лежит внутри записи, в поле count. Причина в поле reason бывает не только rejected: expired — вышел срок жизни сообщения, maxlen — очередь упёрлась в лимит длины, delivery_limit — сообщение сняли по x-delivery-limit. Код разбора DLQ, который умеет только rejected, на остальных случаях покажет неверную картину. Длина списка — это не количество попыток: сообщение, вернувшееся из одной очереди сто раз, даст список из одного элемента со счётчиком 100.
То же самое настройками, без своего бина
Перехватчик выше даёт полный контроль, и в большинстве проектов он не нужен: те же повторы включаются свойствами, и тогда своего бина фабрики не требуется вовсе.
spring:
rabbitmq:
listener:
simple:
prefetch: 20
concurrency: 4
max-concurrency: 8
default-requeue-rejected: false
retry:
enabled: true
max-attempts: 4
initial-interval: 1s
multiplier: 2
max-interval: 30s
Что здесь важно понимать про каждую строку.
prefetch — сколько неподтверждённых сообщений брокер отдаёт одному потребителю; по умолчанию 250, и это много для медленного обработчика. Значение меньше числа потоков оставляет потоки без работы, значение в сотни при обработке в секунды означает, что при падении потребителя в повторную доставку уйдёт вся пачка. Разумно начинать с числа потоков плюс небольшой запас.
concurrency — сколько потоков слушают очередь. Вместе с prefetch это и есть весь параллелизм потребителя.
default-requeue-rejected: false — то же, что setDefaultRequeueRejected(false) в коде: после исчерпания попыток сообщение не возвращается в очередь, а уходит в DLX. Без этой строки повторы закончатся бесконечным кругом.
retry.* — те же попытки и та же выдержка, что в перехватчике. Восстановитель по умолчанию при этом RejectAndDontRequeueRecoverer, то есть поведение совпадает с примером выше.
Когда всё-таки нужен свой бин: если восстановитель другой (например, тот, что перекладывает сообщение в DLQ с причиной — о нём ниже), если нужна своя логика «какие исключения повторять, а какие нет», или если у разных очередей разные политики повторов и нужны несколько фабрик.
Причина отказа в самом сообщении: RepublishMessageRecoverer
С настройками выше сообщение после всех попыток уходит в DLX, и там оно лежит ровно в том виде, в каком пришло: тело, заголовки, x-death с причиной вроде rejected. Чего в нём нет — того, что случилось у вас: ни класса исключения, ни сообщения, ни стектрейса. Разбирать такую очередь приходится сопоставлением времени с логами, а если сообщений в ней сотня, это работа на день.
Лечится это заменой восстановителя. RepublishMessageRecoverer не отвергает сообщение, а публикует его сам — в указанный обменник с указанным ключом, добавив заголовки с причиной:
@Bean
MessageRecoverer messageRecoverer(RabbitTemplate rabbit) {
return new RepublishMessageRecoverer(rabbit, "orders.dlx", "orders.failed")
.errorRoutingKeyPrefix("");
}
В отправленном сообщении появятся заголовки x-exception-message, x-exception-stacktrace, x-original-exchange, x-original-routingKey. Теперь очередь недоставленных читается сама: видно, какое исключение, на какой строке и что было исходным ключом.
Три оговорки. Публикует он обычным шаблоном, поэтому на его отправку тоже действуют подтверждения и возвраты из раздела выше. Стектрейс — большой текст, и заголовки у AMQP не бесплатны: при потоке отказов очередь недоставленных растёт быстрее ожидаемого, поэтому у неё тоже должен быть предел длины. И третье: когда восстановитель задан бином MessageRecoverer, автоконфигурация Spring Boot подхватывает его для повторов из свойств — то есть свой бин фабрики по-прежнему не нужен.
Режимы подтверждения
RabbitMQ работает на принципе подтверждений: брокер держит сообщение до тех пор, пока consumer не скажет «получил» (ack) или «не смог» (nack). Это защита от потери сообщений.
Spring AMQP отвечает за подтверждения сам: в режиме AUTO (он по умолчанию) ack уходит после успешного возврата из метода, а при исключении — nack. Для большинства обработчиков этого достаточно.
MANUAL нужен, когда подтвердить можно только позже, чем метод вернулся, — например, после ответа внешней системы. Тогда ack и nack вызывают руками через объект Channel, а забытое подтверждение оставляет сообщение висеть невыданным до разрыва соединения.
NONE отключает подтверждения вовсе: брокер считает сообщение доставленным в момент отправки, и падение обработчика его теряет. Режим для данных, которые не жалко.
Стандартный рецепт: AUTO + setDefaultRequeueRejected(false) + DLX. Этого достаточно для большинства случаев.
Ручное подтверждение в коде и повторная доставка
MANUAL описан выше строкой, и стоит показать, как это выглядит: метод получает канал и метку доставки и подтверждает сам.
@RabbitListener(queues = "orders.created", ackMode = "MANUAL")
public void handle(OrderCreatedEvent event, Channel channel,
@Header(AmqpHeaders.DELIVERY_TAG) long tag) throws IOException {
try {
orders.apply(event);
channel.basicAck(tag, false);
} catch (DuplicateEventException e) {
channel.basicAck(tag, false);
} catch (TemporaryFailureException e) {
channel.basicNack(tag, false, true);
} catch (Exception e) {
channel.basicNack(tag, false, false);
}
}
Второй аргумент у обоих вызовов — multiple: с true подтверждаются все сообщения до этой метки включительно, и это опасно, если обработка идёт в несколько потоков. Третий у basicNack — requeue: true возвращает сообщение в очередь (то есть попробуем ещё раз), false отправляет в DLX. Забытое подтверждение не даёт ошибки: сообщение просто висит невыданным, пока не разорвётся соединение, а потом возвращается в очередь — отсюда странные жалобы «сообщения обрабатываются по второму разу при перезапуске».
И главное, что идёт в комплекте с at-least-once: обработчик обязан быть идемпотентным. Сообщение придёт дважды не потому, что кто-то ошибся, а потому, что таков протокол: подтверждение не дошло, соединение оборвалось, сработал повтор. Минимальная защита — таблица обработанных событий:
@Transactional
public void apply(OrderCreatedEvent event) {
if (!processed.tryInsert(event.id())) {
return;
}
orders.create(event);
}
Вставка идентификатора и бизнес-действие — в одной транзакции, иначе между ними приложение упадёт, и защита не сработает. Где хранить такие ключи, сколько и как чистить, разбирает статья про шаблоны обмена. Естественный ключ здесь предпочтительнее случайного: идентификатор события из его тела, а не метка доставки — она у повторной доставки другая.
Целостный пример
Вот как выглядит минимальный рабочий сетап: топология, listener и обработчик мёртвых сообщений вместе:
@Configuration
public class OrdersTopology {
@Bean DirectExchange ordersEx() { return new DirectExchange("orders", true, false); }
@Bean DirectExchange ordersDlx() { return new DirectExchange("orders.dlx", true, false); }
@Bean
Queue ordersFulfillment() {
return QueueBuilder.durable("orders.fulfillment").quorum()
.withArgument("x-dead-letter-exchange", "orders.dlx")
.withArgument("x-dead-letter-routing-key", "order.failed")
.build();
}
@Bean Queue ordersDlq() { return QueueBuilder.durable("orders.dlq").quorum().build(); }
@Bean
Binding bindFulfillment(Queue ordersFulfillment, DirectExchange ordersEx) {
return BindingBuilder.bind(ordersFulfillment).to(ordersEx).with("order.created");
}
@Bean
Binding bindDlq(Queue ordersDlq, DirectExchange ordersDlx) {
return BindingBuilder.bind(ordersDlq).to(ordersDlx).with("order.failed");
}
}
@Component
@RequiredArgsConstructor
class OrderListener {
private final OrderHandler handler;
@RabbitListener(queues = "orders.fulfillment", concurrency = "3-10")
public void on(OrderCreatedEvent event) {
handler.process(event);
}
}
@Component
class DlqInspector {
@RabbitListener(queues = "orders.dlq")
public void inspect(
OrderCreatedEvent event,
@Header("x-death") List<Map<String, Object>> deaths) {
log.error("Не удалось обработать: orderId={}", event.orderId());
}
}
Когда брокер недоступен
Про это обычно узнают в первый же сбой, поэтому лучше знать заранее: что делает приложение, пока брокера нет.
Соединение восстанавливается само. Клиентская библиотека RabbitMQ по умолчанию включает автоматическое восстановление: потеряв соединение, она пробует переподключиться (раз в пять секунд по умолчанию), а после успеха заново открывает каналы и восстанавливает подписки. Spring AMQP этим и пользуется, поэтому в логах видно череду «потеряно соединение — восстановлено», а обработчики продолжают работать.
Топология объявляется заново. Объявления, сделанные бинами Queue, Exchange и Binding, при переподключении повторяются: RabbitAdmin объявляет их снова. Это важное свойство, ради которого топологию и описывают бинами, а не командами в коде: брокер, поднятый с нуля (после аварии или в тесте), получит очереди и привязки сам. Обратная сторона — расхождение свойств: если очередь уже существует с другими параметрами, объявление не пройдёт, канал закроется с ошибкой PRECONDITION_FAILED, и это самая частая ошибка при смене настроек очереди. Очередь в таком случае удаляют и создают заново, что для боевого контура означает окно с остановкой потребителей.
Публикация не ждёт. RabbitTemplate при недоступном брокере не копит сообщения бесконечно: попытка отправки ждёт соединения и, не получив его, бросает AmqpConnectException (или AmqpTimeoutException). Отдельный случай — брокер доступен, но заблокировал публикацию: при нехватке памяти или места на диске RabbitMQ применяет обратное давление и перестаёт принимать сообщения. Тогда отправка зависает до таймаута, а не отвечает ошибкой сразу. Готовятся к обоим случаям одинаково: таблица исходящих сообщений и фоновая отправка из неё, чтобы недоступность брокера превращалась в задержку. Проверить это на стенде просто — остановить контейнер брокера и посмотреть, что делает ваш обработчик запроса.
Потребитель при этом не теряет ничего. Неподтверждённые сообщения брокер вернёт в очередь после разрыва соединения, и они придут снова — что возвращает нас к идемпотентности из раздела выше.
Глубже: ошибка до обработчика: MessageConversionException и фатальные исключениярасширенное
Всё в разделе про ошибки касается исключений из вашего метода. Есть класс ошибок, до метода не доходящих, и именно он даёт настоящий бесконечный цикл на классической очереди: сообщение не удаётся превратить в аргумент. Тело не JSON, JSON не сходится с классом (строка в числовом поле), заголовок __TypeId__ указывает на класс, которого нет, @Valid на аргументе отверг тело.
Spring AMQP знает про это и по умолчанию ведёт себя правильно: ConditionalRejectingErrorHandler с DefaultExceptionStrategy считает MessageConversionException, MethodArgumentNotValidException, MethodArgumentTypeMismatchException, NoSuchMethodException и ClassCastException фатальными и отклоняет сообщение без возврата в очередь. Дальше два исхода: у очереди есть Dead Letter Exchange, и сообщение уезжает туда с x-death; DLX нет, и брокер молча удаляет сообщение. Второй исход это потеря без единой строки в логе брокера, только WARN в логе приложения, поэтому DLX объявляют у каждой очереди, куда пишут не вы.
Что здесь легко сломать. Свой ErrorHandler на контейнере, который «логирует и бросает дальше» или возвращает сообщение в очередь, отменяет умное умолчание, и битое сообщение начинает крутиться. Повторы через RetryOperationsInterceptor фатальные исключения не повторяют, и это правильно: разбор не починится от паузы. Список фатальных расширяют, а не сужают: DefaultExceptionStrategy наследуют и добавляют свои исключения валидации, чтобы «сумма отрицательная» тоже уезжала в DLQ с первой попытки, а не после четырёх.
Как это разбирать потом: в DLQ лежит исходное тело как байты, и читать его нужно не тем же слушателем с JSON-конвертером (он споткнётся снова), а отдельным, со строковым конвертером, который пишет тело и заголовки в журнал и, если исправление возможно, переотправляет. Тот же приём для Kafka называется ErrorHandlingDeserializer, о чём статья про Kafka в production.
Глубже: тест пути до DLQ на настоящем RabbitMQрасширенное
Настройка повторов и DLX это код, который не выполняется неделями, а потом обязан сработать с первого раза. Проверяют его тестом на настоящем брокере: RabbitMQ в Testcontainers, объявления очередей из приложения, отравленное сообщение и ожидание в очереди недоставленных.
@SpringBootTest(properties = "app.retry.max-attempts=2")
@Testcontainers
class OrderDlqTest {
@Container
@ServiceConnection
static RabbitMQContainer rabbit = new RabbitMQContainer("rabbitmq:4.0-management");
@Autowired RabbitTemplate template;
@Test
void invalidOrderLandsInDlq() {
template.convertAndSend("orders", "order.created", new OrderCreated("order-1", new BigDecimal("-1")));
Message dead = template.receive("orders.dlq", 10_000);
assertThat(dead).isNotNull();
List<Map<String, ?>> deaths = dead.getMessageProperties().getXDeathHeader();
assertThat(deaths.get(0).get("reason")).isEqualTo("rejected");
}
@Test
void garbageLandsInDlqWithoutRetries() {
template.send("orders", "order.created", new Message("не json".getBytes(UTF_8)));
assertThat(template.receive("orders.dlq", 10_000)).isNotNull();
}
}
Первый тест проверяет бизнес-ошибку: сообщение прошло повторы (их число уменьшено свойством, чтобы тест не ждал минуту) и легло в DLQ с причиной rejected. Второй проверяет фатальную ошибку из предыдущего раздела: мусор вместо JSON попадает в DLQ сразу. receive с таймаутом ждёт появления сообщения, и десять секунд это запас, а не норма: на здоровом стенде путь занимает миллисекунды. Читать DLQ через template.receive безопасно, потому что слушателя на неё в приложении нет, а если есть, тест смотрит через RabbitAdmin.getQueueInfo на число сообщений.
Встроенного брокера, как EmbeddedKafka, у RabbitMQ нет, и заглушки на RabbitTemplate путь до DLQ не проверяют: он целиком живёт в настройках брокера и контейнера. Один такой тест на приложение окупается первым же инцидентом с «сообщение исчезло».
Глубже: трасса через брокер и метрики слушателярасширенное
Запрос пользователя закончился публикацией события, обработка случилась в другом сервисе через секунду, и в трассе это выглядит как два несвязанных куска. Связывает их заголовок traceparent: Spring Boot 3 с Micrometer Tracing пишет его при отправке и продолжает трассу в слушателе, если включить spring.rabbitmq.template.observation-enabled и spring.rabbitmq.listener.observation-enabled. Спан слушателя становится потомком спана публикации, и в трассе видно время, которое сообщение пролежало в очереди, как разрыв между ними. Тот же заголовок переживает DLQ и повторы, поэтому по трассе видно и путь через них.
Метрики с двух сторон. Со стороны приложения таймеры spring.rabbitmq.listener и spring.rabbitmq.template с тегом результата: 99-й перцентиль слушателя растёт раньше, чем очередь. Со стороны брокера, через плагин Prometheus, rabbitmq_queue_messages_ready (ждут), rabbitmq_queue_messages_unacked (в работе, застрявшие видны здесь), rabbitmq_queue_consumers (ноль это инцидент), и число сообщений в очередях недоставленных, каждое из которых для платёжных очередей это тревога. Возраст самого старого сообщения, который важнее длины очереди, отдаёт head_message_timestamp через management API для классических очередей; для остальных его считают сами по отметке времени в сообщении при получении, как now - timestamp. По каким из этих чисел звать дежурного, разбирает статья про RabbitMQ в production.
Коротко
- Spring AMQP (
spring-boot-starter-amqp) оборачивает RabbitMQ-клиент и убирает низкоуровневый код. RabbitTemplate — для отправки; основной методconvertAndSend(exchange, routingKey, payload). @RabbitListener— для получения; параметрconcurrencyзадаёт размер пула потоков. По умолчанию сериализация — Java, для межсервисного взаимодействия нуженJackson2JsonMessageConverter(в spring-amqp 4.0 он называетсяJacksonJsonMessageConverter).- При ошибке без настройки сообщение крутится по кругу: у классической очереди бесконечно, у надёжной брокер оборвёт круг на двадцатой доставке и без DLX выбросит сообщение. Нужны
setDefaultRequeueRejected(false)и DLX. - Повторы (
RetryOperationsInterceptorили свойстваspring.rabbitmq.listener.simple.retry.*) делают несколько попыток с паузой; после исчерпания — сообщение уходит в Dead Letter Exchange. DLX/DLQ настраивается черезQueueBuilder.withArgument("x-dead-letter-exchange", ...). - Стандартный режим подтверждений:
AUTO+requeue=false+ DLX. Ошибка разбора до метода фатальна по умолчанию: без возврата в очередь, в DLX или в никуда, если DLX нет; свойErrorHandlerс возвратом в очередь включает вечный цикл; DLQ читают строковым конвертером. - Путь до DLQ проверяют одним тестом на RabbitMQ в Testcontainers: бизнес-ошибка после повторов и мусор сразу,
receiveс таймаутом и причина вx-death. Трасса идёт черезtraceparentпри включённых observation у template и listener; смотрят таймер слушателя,messages_ready,messages_unacked, число потребителей и возраст головы очереди. - Гарантия отправителя — это
publisher-confirm-type: correlatedплюсpublisher-returnsсmandatory: первое отвечает «брокер принял», второе «нашлась очередь», и оба обратных вызова асинхронные. - Повторы,
prefetch(по умолчанию 250) иdefault-requeue-rejected: falseнастраиваются свойствамиspring.rabbitmq.listener.simple.*— свой бин фабрики нужен только для нестандартного восстановителя. RepublishMessageRecovererкладёт в DLQ причину и стектрейс заголовками — без него очередь недоставленных разбирают по логам. At-least-once требует идемпотентного обработчика: вставка идентификатора события и бизнес-действие в одной транзакции, ключ из тела события, а не метка доставки.- При недоступном брокере соединение и топология восстанавливаются сами (и ловят
PRECONDITION_FAILEDпри смене свойств очереди), а публикация падает по таймауту — поэтому отправку страхуют таблицей исходящих.
Что почитать дальше
- Протокол AMQP — модель доставки: exchange, queue, binding, ack.
- RabbitMQ в production — Quorum Queues, кластеризация, мониторинг.
- Messaging-паттерны через AMQP — work queue, pub/sub, RPC, идемпотентный consumer.