Когда сервисы начинают общаться через брокер, возникает вопрос: как правильно организовать очереди, обменники и подписки? Большинство задач укладываются в несколько типовых схем. Разберём каждую на примерах Spring AMQP.
Разница между ними не в коде отправителя, а в том, сколько получателей увидит одну публикацию.
Отправитель во всех четырёх тактах делает одно и то же — публикует. Сколько получателей увидит публикацию, решает топология: work queue делит 3 задания на 3 обработчика, fanout превращает одно событие в 3 копии, direct по точному ключу даёт 4 доставки из 3 событий, topic по маске — 2 из одного.
Раздать задачи нескольким обработчикам — Work Queue
Представьте: пользователи загружают фотографии, и каждую надо сжать и нарезать в несколько размеров. Это долго. Делать это прямо в HTTP-запросе нельзя — пользователь будет ждать минуты.
Решение — work queue (очередь задач): сохранить задание в очередь и отдать её одному из пула обработчиков.
Очередь одна, потребителей несколько: брокер раздаёт сообщения по одному, и каждое достаётся ровно кому-то одному. Больше потребителей — быстрее разбирается очередь, но общий порядок между ними уже не сохраняется.
Брокер сам распределяет задачи между обработчиками. Каждое сообщение получает ровно один из них.
@Configuration
class ImageProcessingTopology {
@Bean
Queue imagesQueue() {
return QueueBuilder.durable("images.to-process").quorum().build();
}
}
@Component
class ImageProcessor {
@RabbitListener(queues = "images.to-process", concurrency = "5-20")
public void process(ImageJob job) {
// обработка изображения
}
}
concurrency = "5-20" означает: минимум 5 потоков, максимум 20. Если запустить 5 копий сервиса — получится 25–100 параллельных обработчиков.
Когда брать: фоновые задачи — обработка файлов, отправка писем, генерация отчётов, любые «положили в очередь — кто-то возьмёт».
Отправить событие всем сервисам сразу — Publish/Subscribe
Другая задача: произошло событие «конфигурация обновлена» — каждый сервис должен обновить свой кеш. Нельзя знать заранее, кто именно подписан и сколько сервисов запущено.
Это publish/subscribe (рассылка всем): издатель отправляет одно сообщение, а копию получают все подписчики одновременно.
Здесь нужен fanout exchange — он копирует каждое сообщение во все привязанные очереди. Каждый сервис объявляет свою очередь и привязывает её к общему обменнику.
@Configuration
class CacheInvalidationTopology {
@Bean FanoutExchange cacheInvalidation() {
return new FanoutExchange("cache.invalidation", true, false);
}
@Bean Queue serviceACache() {
return QueueBuilder.nonDurable().exclusive().autoDelete().build();
}
@Bean Binding bindA(Queue serviceACache, FanoutExchange cacheInvalidation) {
return BindingBuilder.bind(serviceACache).to(cacheInvalidation);
}
}
@Component
class ServiceACacheListener {
@RabbitListener(queues = "#{serviceACache.name}")
public void invalidate(CacheInvalidationEvent event) {
cache.evict(event.key());
}
}
exclusive + autoDelete — очередь принадлежит одному соединению и удаляется при отключении. При перезапуске сервиса не накапливается мусор в брокере.
Когда брать: инвалидация кешей, broadcast-уведомления всему кластеру, обновления конфигурации.
Направить событие в нужный обработчик — Routing
Иногда нужно не «всем», а «именно тому, кому надо». Например: событие order.created должно идти в сервис выполнения и в аудит, а order.payment-failed — только в алерты.
Это routing (точечная маршрутизация): direct exchange смотрит на ключ маршрутизации сообщения и отправляет его только в очереди с совпадающим binding key.
@Configuration
class OrderRoutingTopology {
@Bean DirectExchange orders() { return new DirectExchange("orders", true, false); }
@Bean Queue fulfillment() { return QueueBuilder.durable("orders.fulfillment").quorum().build(); }
@Bean Queue audit() { return QueueBuilder.durable("orders.audit").quorum().build(); }
@Bean Queue alerts() { return QueueBuilder.durable("orders.alerts").quorum().build(); }
@Bean Binding b1(Queue fulfillment, DirectExchange orders) {
return BindingBuilder.bind(fulfillment).to(orders).with("order.created");
}
@Bean Binding b2(Queue audit, DirectExchange orders) {
return BindingBuilder.bind(audit).to(orders).with("order.created");
}
@Bean Binding b3(Queue audit, DirectExchange orders) {
return BindingBuilder.bind(audit).to(orders).with("order.cancelled");
}
@Bean Binding b4(Queue alerts, DirectExchange orders) {
return BindingBuilder.bind(alerts).to(orders).with("order.payment-failed");
}
}
order.created→ fulfillment + audit.order.cancelled→ только audit.order.payment-failed→ только alerts.
Когда брать: явное разделение потоков — алерты отдельно от аудита, основной обработчик отдельно от мониторинга.
Подписаться по маске — Topic
Routing хорош для жёстких правил. Но что если сервис хочет подписаться на «все события по заказам»? Или «всё из EU-региона»?
Topic exchange позволяет задавать подписки через шаблоны. Ключи сообщений строятся через точку (order.created.eu), а в подписке можно использовать:
*— ровно одно слово,#— ноль и более слов.
@Configuration
class TopicRoutingTopology {
@Bean TopicExchange events() { return new TopicExchange("events", true, false); }
@Bean Queue auditAllOrders() { return QueueBuilder.durable("audit.orders").quorum().build(); }
@Bean Queue euDashboard() { return QueueBuilder.durable("dashboard.eu").quorum().build(); }
@Bean Queue alerts() { return QueueBuilder.durable("alerts.critical").quorum().build(); }
@Bean Binding b1(Queue auditAllOrders, TopicExchange events) {
return BindingBuilder.bind(auditAllOrders).to(events).with("order.#");
}
@Bean Binding b2(Queue euDashboard, TopicExchange events) {
return BindingBuilder.bind(euDashboard).to(events).with("*.*.eu");
}
@Bean Binding b3(Queue alerts, TopicExchange events) {
return BindingBuilder.bind(alerts).to(events).with("payment.failed.#");
}
}
Сообщение с ключом order.cancelled.eu попадёт в auditAllOrders (по order.#) и в euDashboard (по *.*.eu).
Когда брать: события с иерархической структурой, когда нужно гибко подписываться без переделки топологии при добавлении новых типов событий.
Порядок по ключу при нескольких обработчиках
Work Queue выше раздаёт сообщения нескольким обработчикам, и это ровно то, что нужно, пока сообщения независимы. Как только они связаны — два обновления одного заказа, создание и отмена одной брони, — параллельная обработка ломает порядок: второе сообщение попадает к свободному обработчику раньше, чем первый закончил.
В Kafka порядок держит ключ партиции. У RabbitMQ порядок гарантирован только внутри одной очереди с одним потребителем, и из этого растут два приёма.
Single active consumer. Аргумент очереди x-single-active-consumer: true означает: подписчиков может быть много, но сообщения получает один — остальные ждут и мгновенно подхватывают работу, если активный отвалился. Порядок сохранён, отказоустойчивость есть, параллелизма нет. Это правильный выбор, когда очередь и так не перегружена, а порядок важен целиком.
@Bean Queue bookings() {
return QueueBuilder.durable("bookings")
.singleActiveConsumer()
.build();
}
Consistent hash exchange. Плагин rabbitmq_consistent_hash_exchange добавляет обменник, который раскладывает сообщения по нескольким очередям по хешу ключа маршрутизации (или заголовка). Сообщения одного заказа всегда попадают в одну очередь, у каждой очереди свой единственный потребитель — получается то же, что партиции в Kafka: порядок внутри ключа плюс параллелизм между ключами. Цена — ручное управление: очереди создаёте вы, число очередей меняется только с остановкой, а перекос по ключам («один крупный клиент») складывается в одну очередь.
Что не работает: «сделаем один потребитель, но в десять потоков». Потоки разберут сообщения из общей выдачи и обработают их вперемешку — порядок теряется так же, как при нескольких потребителях. Если нужна и скорость, и порядок, ключ обязан определять, кто обрабатывает; других вариантов нет.
Запрос-ответ через очередь — RPC
Иногда нужен синхронный ответ, но HTTP не подходит: сервис за NAT, нет публичного адреса, или хочется балансировки по пулу обработчиков.
RPC через очередь: клиент отправляет запрос и ждёт ответа. Брокер доставляет запрос одному из обработчиков, тот отвечает в отдельную очередь-ответ. Для сопоставления запроса и ответа используется correlation-id.
В Spring AMQP это скрыто за sendAndReceive:
// Клиент
@Component
@RequiredArgsConstructor
class PricingClient {
private final RabbitTemplate rabbit;
public PriceQuote quote(QuoteRequest request) {
return (PriceQuote) rabbit.convertSendAndReceive(
"pricing.exchange", "pricing.quote", request);
}
}
// Сервер
@Component
class PricingServer {
@RabbitListener(queues = "pricing.quote")
public PriceQuote handle(QuoteRequest request) {
return PriceQuote.compute(request); // возврат автоматически уходит в reply-to
}
}
Очередь-ответ Spring AMQP при этом не заводит: начиная с версии 1.4.1 RabbitTemplate отвечает через прямой ответ — псевдо-очередь amq.rabbitmq.reply-to, а настоящую временную очередь создаёт только как запасной путь для старых брокеров. Он же проставляет reply-to и correlation-id и ждёт ответа. Возврат из @RabbitListener автоматически публикуется обратно.
Ждёт, но не бесконечно: по умолчанию convertSendAndReceive стоит на месте пять секунд (reply-timeout), а потом возвращает null. Пример выше этого не проверяет — при таймауте метод quote молча отдаст null, и падение вылезет где-то дальше по коду, далеко от настоящей причины. Таймаут надо либо обрабатывать на месте, либо превращать в понятную ошибку.
Что ломается, кроме таймаута. У RPC через брокер отказы неочевидные, и их стоит перечислить.
Ответ пришёл после таймаута. Клиент уже бросил ошибку и ушёл, а обработчик закончил работу и опубликовал ответ. С прямым ответом (amq.rabbitmq.reply-to) такой ответ просто выбрасывается — но работа-то выполнена. Значит, обработчик RPC должен быть идемпотентным так же, как обработчик события: повтор запроса после таймаута не должен списать деньги второй раз. Для чего-либо изменяющего данные RPC вообще плохой выбор, и лучше разделить на команду и уведомление.
Ответ потерян. Прямой ответ живёт на конкретном соединении клиента: разорвалось соединение — отвечать некуда, и ответ пропадёт. Перезапуск экземпляра клиента во время ожидания означает потерянный ответ, даже если запрос обработан.
Потребители кончились. Если обработчиков нет вовсе, запрос ляжет в очередь и будет лежать до таймаута клиента. Внешне это выглядит как «сервис тормозит», а не «сервис недоступен». Полезно ставить очереди x-message-ttl чуть меньше таймаута клиента: тогда просроченные запросы не будут обрабатываться после того, как их перестали ждать.
Одновременных запросов слишком много. Каждый ждущий вызов занимает поток, и при медленном обработчике пул потоков приложения выедается ожиданием. Здесь и появляется главный аргумент против RPC через брокер: у HTTP-клиента есть пул соединений с понятным пределом, у ожидания в очереди предела нет. Ограничитель ставят сами: семафор на число одновременных вызовов или отдельный пул.
Проверять это стоит на стенде: остановить обработчик и посмотреть, что делает клиентский поток и что видит пользователь.
Когда брать: нужен синхронный вызов, но HTTP не работает (NAT, firewall, нет публичного адреса); нужна балансировка запросов по пулу обработчиков.
Когда не брать: если HTTP/gRPC просто работает — RPC через брокер сложнее в отладке и дороже.
Что делать с повторной доставкой — Idempotent Consumer
С ручным подтверждением AMQP даёт at-least-once доставку: одно и то же сообщение может прийти дважды. Это случается, когда брокер не получил подтверждение обработки (например, из-за сетевой проблемы) и переотправляет сообщение. Гарантию даёт именно ack: если включить автоподтверждение, брокер забудет о сообщении сразу после отправки, и вместо повторов вы получите потери — at-most-once.
Обработчик обязан уметь работать с повторами, не ломая бизнес-логику.
Ключ идемпотентности в базе данных
Самый надёжный способ — запоминать уже обработанные сообщения:
@RabbitListener(queues = "payments")
@Transactional
public void process(PaymentEvent event) {
// INSERT ... ON CONFLICT DO NOTHING: вернёт 1, если ключ вставлен, и 0, если он уже был
if (processedEventsRepo.insertIfAbsent(event.idempotencyKey()) == 0) {
return; // такой ключ уже лежит в таблице — это повтор, просто подтверждаем получение
}
accountRepo.debit(event.accountId(), event.amount());
}
Таблица processed_events с уникальным индексом на idempotency_key. Важен порядок: сначала вставка ключа, потом сама работа. Наоборот — сперва «а не обработано ли уже?», потом вставка — не годится: два потребителя, которым достались копии одного сообщения, успевают оба увидеть «нет» и оба списывают деньги. От такой гонки спасает не проверка в коде, а уникальный индекс: первую вставку база пропустит, вторую отобьёт, а ON CONFLICT DO NOTHING превратит отказ в ноль строк вместо исключения посреди транзакции.
Проверка состояния объекта
Если событие переводит объект в новое состояние, достаточно проверить текущее:
@Transactional
public void onOrderConfirmed(OrderConfirmedEvent event) {
var order = orderRepo.findById(event.orderId()).orElseThrow();
if (order.status() == OrderStatus.CONFIRMED) {
return; // уже в нужном состоянии
}
order.confirm();
orderRepo.save(order);
}
Не требует отдельной таблицы — состояние уже хранится в бизнес-объекте.
Retry с задержкой и Dead Letter Queue
Что если обработчик упал не из-за бага, а из-за временной недоступности внешнего сервиса? Нужно попробовать снова, но не немедленно.
В Spring AMQP delayed retry собирается через x-message-ttl и Dead Letter Exchange:
@Bean DirectExchange ordersRequeue() { return new DirectExchange("orders.requeue", true, false); }
@Bean Binding backToFulfillment(Queue fulfillment, DirectExchange ordersRequeue) {
return BindingBuilder.bind(fulfillment).to(ordersRequeue).with("order.created");
}
@Bean Queue retryQueue() {
return QueueBuilder.durable("orders.retry")
.withArgument("x-message-ttl", 30_000) // ждём 30 секунд
.withArgument("x-dead-letter-exchange", "orders.requeue")
.withArgument("x-dead-letter-routing-key", "order.created")
.withArgument("x-dead-letter-strategy", "at-least-once")
.withArgument("x-overflow", "reject-publish")
.quorum().build();
}
Поток: обработчик отклонил сообщение → оно попадает в retry очередь → через 30 секунд по истечении срока уходит через обменник недоставленных обратно в основную очередь → новая попытка.
Круг повтора: рабочая очередь отклоняет сообщение в orders.retry, там оно ждёт свой срок и возвращается на новую попытку через отдельный обменник orders.requeue, а смотреть стоит на orders.retry, потому что срок брокер проверяет только у головы очереди и разные задержки в одной очереди держат друг друга.
Обратите внимание на отдельный обменник orders.requeue. Возвращать сообщение в orders с ключом order.created нельзя: в топологии из раздела «Routing» этот ключ привязан и к orders.fulfillment, и к orders.audit, так что каждая повторная попытка задваивала бы запись в аудите. Отдельный обменник ведёт ровно в одну очередь — ту, где обработка и падала.
Вторая пара аргументов — про надёжность самого перекладывания. У надёжных очередей стратегия перекладывания в обменник недоставленных по умолчанию at-most-once: брокер не обещает, что сообщение туда доедет. Раз уж весь приём про «не потерять», ставим x-dead-letter-strategy: at-least-once, а он работает только вместе с x-overflow: reject-publish. Ту же пару стоит поставить и основной очереди.
Одна тонкость этого приёма: брокер проверяет срок жизни только у сообщения в голове очереди. Если впереди лежит сообщение с бо́льшим оставшимся сроком, всё, что за ним, будет ждать его — даже если само уже «созрело». Пока задержка у всех одинаковая, как здесь, проблемы нет; она появляется, когда для разных попыток задают разные задержки в одной очереди.
Лечат головную блокировку двумя способами. Первый — очередь на каждую задержку: retry.10s, retry.1m, retry.10m, и попытка выбирает очередь по своему номеру. Внутри каждой задержка одинаковая, значит порядок совпадает со сроками, и голова никого не держит. Способ работает на любом брокере и стоит трёх лишних очередей в топологии.
Второй, который в проде встречается чаще, — плагин отложенной доставки rabbitmq_delayed_message_exchange. Он включается на брокере (rabbitmq-plugins enable rabbitmq_delayed_message_exchange) и добавляет особый тип обменника, у которого задержка задаётся у каждого сообщения заголовком:
@Bean CustomExchange delayedExchange() {
return new CustomExchange("orders.delayed", "x-delayed-message", true, false,
Map.of("x-delayed-type", "direct"));
}
rabbit.convertAndSend("orders.delayed", "order.created", event, message -> {
message.getMessageProperties().setHeader("x-delay", 30_000);
return message;
});
Одна очередь, произвольные задержки, никакой головной блокировки. Цена тоже есть, и её стоит знать: плагин держит отложенные сообщения в своей встроенной базе на том узле, куда они опубликованы, — это не очередь, и она не реплицируется. Падение узла до срока доставки теряет отложенные сообщения, а большое их число заметно ест память. Отсюда правило: для повторов и коротких задержек плагин удобен, а для «напомнить через месяц» берут не брокер, а таблицу в базе и задачу по расписанию.
Число попыток лежит в заголовке x-death. После трёх кругов он выглядит так:
x-death:
- queue: orders.fulfillment reason: rejected count: 3
- queue: orders.retry reason: expired count: 3
Записей две, потому что сообщение каждый раз «умирает» дважды: его отклонили в основной очереди и у него истёк срок в retry-очереди. Считать попытки нужно по записи своей очереди с причиной rejected; в обработчике она доступна как @Header(name = "x-death", required = false) List<Map<String, Object>> deaths, и по count решают, отправлять ли сообщение на новый круг. У надёжных очередей есть и встроенный ограничитель — аргумент x-delivery-limit. С RabbitMQ 4.0 он не просто «есть», а включён по умолчанию со значением 20: после двадцатой доставки брокер сам снимет сообщение, без вашего кода. Отправит в обменник недоставленных, если тот настроен, — и молча выбросит, если нет. Это ещё один довод настроить DLX до того, как он понадобится.
Сообщения, которые не получилось обработать после всех попыток, уходят в Dead Letter Queue (DLQ) — отдельную очередь для разбора вручную или через алерты.
Гарантированная публикация — Outbox
Бывает задача: сохранить заказ в базу данных и опубликовать событие — атомарно. Если сначала сохранить, потом опубликовать, то сервис может упасть между двумя операциями. Событие потеряется.
Outbox pattern: событие сохраняется в ту же транзакцию, что и бизнес-данные. Отдельный процесс читает таблицу и публикует в AMQP.
@Transactional
public void confirm(OrderId orderId) {
var order = orderRepo.findById(orderId).orElseThrow();
order.confirm();
orderRepo.save(order);
outboxRepo.save(new OutboxEvent(
UUID.randomUUID(),
"order.confirmed",
"orders",
toJson(new OrderConfirmedEvent(orderId))
));
}
@Scheduled(fixedDelay = 500)
@Transactional
public void publishOutbox() {
// SELECT ... FOR UPDATE SKIP LOCKED: пачку, занятую соседним экземпляром, этот просто пропустит
var batch = outboxRepo.lockUnpublished(100);
for (var event : batch) {
rabbit.convertAndSend(event.exchange(), event.routingKey(), event.payload());
outboxRepo.markPublished(event.id());
}
}
Либо оба изменения зафиксированы, либо ни одного. Дубли возможны (публикация прошла, но пометить как отправленное не успело) — поэтому получатель всё равно должен быть идемпотентным, тем же способом, что разобран выше в разделе про повторную доставку. Outbox гарантирует, что событие уйдёт, а Idempotent Consumer — что второй его приход ничего не сломает; одно без другого закрывает половину задачи.
Отдельно стоит сказать про выборку пачки. Обычный SELECT здесь опасен: сервис почти всегда работает в нескольких экземплярах, и два прогона по расписанию прочитают одни и те же строки и опубликуют их дважды. Строки нужно забирать под блокировку — SELECT ... FOR UPDATE SKIP LOCKED: первый экземпляр забирает пачку, второй не ждёт её, а сразу берёт следующие свободные строки.
Остальные вопросы про отправляющий процесс решаются так.
Размер пачки и частота. Сотня строк раз в полсекунды — разумная отправная точка: пачка должна успевать отправиться до следующего запуска, иначе прогоны начнут накладываться. Признак, что пачка мала, — таблица растёт быстрее, чем пустеет; признак, что велика — одна транзакция держит блокировки слишком долго. Пачку берут с сортировкой по времени создания, чтобы события не перепрыгивали друг через друга.
Порядок. Гарантировать глобальный порядок отправкой из нескольких экземпляров нельзя. Если порядок важен внутри заказа, отправка идёт по ключу: пачка выбирается так, чтобы события одного заказа не попали одновременно к двум отправителям (например, блокировкой по идентификатору заказа), а на стороне получателя порядок держат приёмы из раздела выше.
Повторы и застрявшие записи. Отправка падает — запись остаётся неотправленной и уйдёт в следующий прогон, это нормальная работа. Ненормально, когда запись не уходит никогда: сообщение не проходит проверку получателя, ключ маршрутизации неверный, тело не разбирается. Поэтому в таблице держат счётчик попыток и время последней; после нескольких неудач запись помечают отложенной с растущей выдержкой, а после потолка — «требует внимания», и на это ставят сигнал. Иначе один битый ряд будет бесконечно занимать место в каждой пачке.
Падение между отправкой и отметкой. Это и есть причина дублей: сообщение ушло, markPublished не выполнился, следующий прогон отправит снова. Полностью избавиться нельзя — брокер и база не участвуют в одной транзакции, — поэтому получателя делают идемпотентным. Уменьшить число дублей помогает отметка в той же транзакции, что и выборка (пачка помечается отправленной сразу после успешного подтверждения от брокера), а не отдельным запросом на каждую строку.
Уборка. Отправленные записи не хранят вечно: либо удаляют пачками по расписанию, либо секционируют таблицу по дате и удаляют старые секции. Таблица исходящих, которая растёт годами, однажды становится самой большой в базе.
Когда не брать: по абзацу на шаблон
У каждого шаблона выше есть граница, за которой он делает хуже. Собираем их в одном месте.
Work Queue не берут, когда сообщения связаны между собой и порядок важен: параллельные обработчики его ломают, а лечится это способами из раздела про порядок по ключу. И не берут, когда задача выполняется быстрее, чем стоит поход в брокер: очередь ради вызова на пять миллисекунд добавляет задержку и точку отказа, не давая ничего взамен.
Publish/Subscribe не берут, когда отправителю важен результат. Событие «заказ создан» ничего не знает о том, справились ли подписчики; если вам нужен ответ или нужна гарантия, что обработали все, это не события, а последовательность вызовов с проверкой. Второе ограничение: широковещательная рассылка плохо переносит большое число подписчиков с разными скоростями — медленный подписчик копит очередь, и за ней надо следить отдельно.
Routing и Topic не берут, когда маска превращается в описание бизнес-логики. Если ключ маршрутизации выглядит как order.created.premium.moscow.retry2, а привязки читаются как правила скидок, значит логика уехала в топологию брокера: её не видно в коде, не покрыть тестом и не отладить. Правило простое: в ключе живёт тип события и, может быть, регион — но не условия.
RPC не берут, когда HTTP или gRPC просто работают, и особенно не берут для изменяющих операций: причины перечислены в разделе выше. Он оправдан только при настоящем сетевом ограничении — нет публичного адреса, нужна балансировка по пулу разнородных обработчиков.
Idempotent Consumer — единственный шаблон, который берут всегда, но и у него есть граница: если действие идемпотентно само по себе (установка значения, а не прибавление; создание по естественному ключу с уникальным индексом), отдельная таблица обработанных событий — лишняя таблица и лишняя запись в каждой транзакции.
Retry с задержкой не берут для ошибок, которые не пройдут сами: неверный формат сообщения, отсутствующая сущность, отказ по правам. Такие сообщения отправляют в очередь недоставленных сразу — повторять их бессмысленно, а место в очереди они занимают. Разделять помогает деление исключений на временные и постоянные.
Outbox не берут, когда терять событие не страшно (уведомление «просмотрено», обновление счётчика для витрины) или когда отправитель и получатель и так работают с одной базой. Он добавляет таблицу, отправляющий процесс, уборку и задержку в полсекунды; всё это оправдано для событий, за потерю которых придётся отвечать.
Шпаргалка по выбору
| Задача | Паттерн | Тип обменника |
|---|---|---|
| Распределить нагрузку между обработчиками | Work Queue | direct (default) |
| Broadcast события всем сервисам | Publish/Subscribe | fanout |
| Разные события в разные очереди | Routing | direct |
| Подписка по маске на иерархические события | Topic | topic |
| Синхронный вызов через очередь | RPC | direct + reply-to |
| Защита от повторной доставки | Idempotent Consumer | любой |
| Retry с задержкой | Delayed Retry | direct + DLX |
| Атомарная публикация вместе с записью в БД | Outbox | direct |
Глубже: отметка и действие в одной транзакции: где хранить ключ и как чиститьрасширенное
Раздел про идемпотентного потребителя говорит «запомнить ключ». Три вопроса, на которых ломается реализация: где запомнить, когда именно и как не хранить вечно.
Где. В той же базе, что и бизнес-данные, таблицей вида processed_messages(message_id, consumer, processed_at) с первичным ключом по паре «идентификатор, потребитель»: одно сообщение обрабатывают несколько потребителей, и у каждого своя отметка. Redis с SETNX быстрее, но живёт в другой системе: отметка поставлена, транзакция с заказом откатилась, и сообщение больше никогда не обработается; или наоборот, заказ записан, Redis не ответил, и повтор создаст второй заказ. Отметка в другом хранилище это всегда окно, в котором оба исхода неверны.
Когда. Вставка отметки идёт в той же транзакции, что и бизнес-действие, и первой: INSERT INTO processed_messages ... бросает нарушение уникальности, если сообщение уже было, и обработчик подтверждает сообщение брокеру без действия. Если вставка прошла, дальше идёт действие, и COMMIT фиксирует оба. Подтверждение брокеру (ack) отправляют после COMMIT; падение между ними даст повтор, который упрётся в отметку. Порядок «действие, потом отметка» в той же транзакции тоже работает, а вот отметка до транзакции или после неё нет.
Какой ключ. У RabbitMQ идентификатор кладёт отправитель в свойство message_id (в Spring MessageProperties.setMessageId), и без него дедупликация невозможна: брокер при повторной доставке не генерирует ничего нового. У Kafka идентификатора нет, и берут либо topic-partition-offset, который одинаков при повторе той же записи, либо messageId в заголовке от продьюсера, который переживает и повторную отправку самим продьюсером. Естественный ключ из бизнеса (order_id плюс версия события) лучше обоих: он защищает и от дублей, которые породил не брокер, а повторный вызов отправителя.
Как чистить. Таблица растёт на каждое сообщение, и её чистят по возрасту: DELETE FROM processed_messages WHERE processed_at < now() - interval '7 days' по расписанию, пакетами. Срок хранения обязан быть длиннее максимального окна повторной доставки: срока жизни сообщения в очереди недоставленных, срока хранения топика, длительности самого долгого инцидента, после которого брокер отдаст старое. Семь дней это обычный минимум; при секционировании таблицы по дню чистка становится удалением секции.
Глубже: версия события: X-Event-Version, совместимость и удалённое полерасширенное
Заголовок X-Event-Version уже стоит в примерах, и раздел про Kafka объясняет Schema Registry, а у RabbitMQ реестра нет, и контракт события держится на договорённостях. Их четыре.
Читатель терпим к лишнему. Потребитель игнорирует незнакомые поля (FAIL_ON_UNKNOWN_PROPERTIES выключен для событий, в отличие от команд), поэтому добавить поле можно всегда, и это единственное изменение, которое не требует ничего от потребителей. Новое поле обязано быть необязательным: старые события в очередях и в DLQ его не содержат.
Удаление и переименование через расширение и сжатие. Поле убирают в три выката: перестают на него полагаться у потребителей, потом перестают писать у отправителя, потом удаляют из схемы. Переименование это добавление нового поля с публикацией обоих, пока потребители не переедут. Смена типа или смысла поля это новое событие, а не новая версия: order.paid с суммой в копейках вместо рублей ломает всех молча, и честнее выпустить order.paid.v2 с отдельным ключом маршрутизации и держать оба, пока последний потребитель не переедет на новое.
Версия в конверте, а не в теле. X-Event-Version в заголовке позволяет потребителю выбрать разборщик до десериализации, и в Spring это @RabbitListener с условием по заголовку или один слушатель с ветвлением по @Header("X-Event-Version"). Мажорную версию поднимают только на несовместимое изменение, и в тот момент она уже означает другое событие.
Схема в общем месте с проверкой в сборке. Без реестра его роль играет репозиторий контрактов: JSON Schema или Avro-схема на каждое событие, из неё генерируются классы у отправителя и потребителей, а сборка сравнивает новую схему со старой и падает на ломающем изменении, как это делает oasdiff для REST в статье про API-first. Контрактные тесты для сообщений (Pact умеет и их) закрывают последнее: потребитель описывает, какое событие он ждёт, и сборка отправителя проверяет, что оно ещё такое. Всё это дешевле, чем Schema Registry рядом с RabbitMQ, и достаточно, пока событий десятки, а не сотни.
Коротко
- Work Queue — одна очередь, несколько обработчиков, каждое сообщение получает ровно один. Для фоновых задач. Publish/Subscribe — fanout exchange копирует сообщение во все привязанные очереди. Для broadcast-событий.
- Routing — direct exchange смотрит на ключ маршрутизации. Для точного разделения потоков. Topic — как routing, но с шаблонами
*и#. Для иерархических событий с гибкой подпиской. - RPC через очередь — запрос-ответ через брокер с
reply-toиcorrelation-id. Для вызовов без HTTP. У RPC через брокер ломается не только таймаут: ответ после таймаута, потерянный ответ при разрыве, запросы без обработчиков и отсутствие предела одновременных вызовов. - Idempotent Consumer — at-least-once означает возможные дубли. Защита: ключ идемпотентности в БД или проверка состояния объекта. Отметка идемпотентности живёт в той же базе и в той же транзакции, что действие, первой вставкой с уникальным ключом; ключ это
message_idот отправителя или естественный ключ события; чистка по возрасту длиннее окна повторной доставки. - Delayed Retry — TTL + DLX: сообщение «паркуется» на время, потом возвращается. Разные задержки в одной очереди держит голова очереди: либо очередь на каждую задержку, либо плагин отложенной доставки, у которого свои ограничения — узел, память и отсутствие репликации.
- Outbox — событие сохраняется в той же транзакции, что и данные. Атомарность без двухфазного коммита. Отправляющий процесс таблицы исходящих нуждается в размере пачки, счётчике попыток с выдержкой, сигнале на застрявшие записи и уборке отправленного; дубли на падении между отправкой и отметкой неизбежны.
- Контракт события без реестра: терпимый читатель, только необязательные добавления, удаление через расширение и сжатие, несовместимое изменение это новое событие, схема в репозитории со сравнением в сборке.
- Порядок по ключу при нескольких обработчиках даёт
x-single-active-consumer(порядок без параллелизма) или обменник по хешу ключа с очередью на ключ; несколько потоков у одного потребителя порядок не сохраняют.
Что почитать дальше
- Протокол AMQP — модель exchange/binding/queue изнутри.
- Spring AMQP — конфигурация, RabbitTemplate, аннотации.
- RabbitMQ в production — Quorum Queues, кластеризация, мониторинг.
- AMQP vs Kafka — какой брокер и когда брать.