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

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

Разница между ними не в коде отправителя, а в том, сколько получателей увидит одну публикацию.

код отправителя один — число доставок решает топология producer 3 задания 1 событие 3 события order.cancelled.eu Work Queueодна очередь на всехкаждое сообщение —ровно одному из пула Publish/Subscribefanoutкопия — в каждуюпривязанную очередь Routingdirectсовпал binding key —ушла копия Topictopic · маски * и #* — ровно одно слово# — ноль и более слов consumer 1взял задание 1consumer 2взял задание 2consumer 3взял задание 3 queue svc-a.cacheкопия событияqueue svc-b.cacheкопия событияqueue svc-c.cacheкопия события orders.fulfillmentorder.created → 1orders.auditcreated + cancelled → 2orders.alertspayment-failed → 1 audit.ordersorder.# → копияdashboard.eu*.*.eu → копияalerts.criticalpayment.failed.# → мимо 3 задания → 3 доставки: каждое взял ровно один 1 событие → 3 доставки: копию получили все 3 события → 4 доставки: order.created ушёл в две очереди 1 событие → 2 доставки: маска решила, alerts не подписан паттерн обменник отправлено доставлено Work Queue direct (default) 3 сообщения 3 доставки Publish/Subscribefanout1 сообщение3 доставки Routingdirect3 сообщения4 доставки Topictopic + маска1 сообщение2 доставки

Отправитель во всех четырёх тактах делает одно и то же — публикует. Сколько получателей увидит публикацию, решает топология: work queue делит 3 задания на 3 обработчика, fanout превращает одно событие в 3 копии, direct по точному ключу даёт 4 доставки из 3 событий, topic по маске — 2 из одного.

Обязательно

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

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

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

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

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

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

@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.fulfillment orders.retry orders.requeue отклонил истёк срок новая попытка orders.dlq попытки кончились

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

Глубже: отметка и действие в одной транзакции: где хранить ключ и как чиститьрасширенное

Раздел про идемпотентного потребителя говорит «запомнить ключ». Три вопроса, на которых ломается реализация: где запомнить, когда именно и как не хранить вечно.

Где. В той же базе, что и бизнес-данные, таблицей вида 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 (порядок без параллелизма) или обменник по хешу ключа с очередью на ключ; несколько потоков у одного потребителя порядок не сохраняют.

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