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

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

Обязательно

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

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

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

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

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

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

channel = await connection.channel()
await channel.set_qos(prefetch_count=20)

images_queue = await channel.declare_queue(
    "images.to-process",
    durable=True,
    arguments={"x-queue-type": "quorum"},
)

async def process(message: AbstractIncomingMessage) -> None:
    async with message.process():
        job = deserialize(message.body)
        # обработка изображения

await images_queue.consume(process)

prefetch_count=20 означает: брокер выдаёт одному потребителю до 20 неподтверждённых сообщений одновременно. Если запустить 5 копий сервиса — получится до 100 параллельных обработчиков.

message.process() подтверждает сообщение, если блок завершился без исключения. При исключении оно отклоняется без возврата в очередь (requeue=False по умолчанию): без очереди мёртвых писем такое сообщение пропадёт, а process(requeue=True) вернёт его в голову очереди, и обработчик с багом будет крутить его бесконечно. Поэтому повторы делают через DLX, как показано ниже.

Когда брать: фоновые задачи — обработка файлов, отправка писем, генерация отчётов, любые «положили в очередь — кто-то возьмёт».

Отправить событие всем сервисам сразу — Publish/Subscribe

Другая задача: произошло событие «конфигурация обновлена» — каждый сервис должен обновить свой кеш. Нельзя знать заранее, кто именно подписан и сколько сервисов запущено.

Это publish/subscribe (рассылка всем): издатель отправляет одно сообщение, а копию получают все подписчики одновременно.

Здесь нужен fanout exchange — он копирует каждое сообщение во все привязанные очереди. Каждый сервис объявляет свою очередь и привязывает её к общему обменнику.

cache_invalidation = await channel.declare_exchange(
    "cache.invalidation", ExchangeType.FANOUT, durable=True
)

service_a_cache = await channel.declare_queue(exclusive=True, auto_delete=True)
await service_a_cache.bind(cache_invalidation)

async def invalidate(message: AbstractIncomingMessage) -> None:
    async with message.process():
        event = deserialize(message.body)
        cache.evict(event.key)

await service_a_cache.consume(invalidate)

exclusive=True, auto_delete=True — очередь принадлежит одному соединению и удаляется при отключении. При перезапуске сервиса не накапливается мусор в брокере.

Когда брать: инвалидация кешей, broadcast-уведомления всему кластеру, обновления конфигурации.

Направить событие в нужный обработчик — Routing

Иногда нужно не «всем», а «именно тому, кому надо». Например: событие order.created должно идти в сервис выполнения и в аудит, а order.payment-failed — только в алерты.

Это routing (точечная маршрутизация): direct exchange смотрит на ключ маршрутизации сообщения и отправляет его только в очереди с совпадающим binding key.

orders = await channel.declare_exchange("orders", ExchangeType.DIRECT, durable=True)

quorum = {"x-queue-type": "quorum"}
fulfillment = await channel.declare_queue("orders.fulfillment", durable=True, arguments=quorum)
audit = await channel.declare_queue("orders.audit", durable=True, arguments=quorum)
alerts = await channel.declare_queue("orders.alerts", durable=True, arguments=quorum)

await fulfillment.bind(orders, routing_key="order.created")
await audit.bind(orders, routing_key="order.created")
await audit.bind(orders, routing_key="order.cancelled")
await alerts.bind(orders, routing_key="order.payment-failed")
  • order.created → fulfillment + audit.
  • order.cancelled → только audit.
  • order.payment-failed → только alerts.

Когда брать: явное разделение потоков — алерты отдельно от аудита, основной обработчик отдельно от мониторинга.

Подписаться по маске — Topic

Routing хорош для жёстких правил. Но что если сервис хочет подписаться на «все события по заказам»? Или «всё из EU-региона»?

Topic exchange позволяет задавать подписки через шаблоны. Ключи сообщений строятся через точку (order.created.eu), а в подписке можно использовать:

  • * — ровно одно слово,
  • # — ноль и более слов.
events = await channel.declare_exchange("events", ExchangeType.TOPIC, durable=True)

quorum = {"x-queue-type": "quorum"}
audit_all_orders = await channel.declare_queue("audit.orders", durable=True, arguments=quorum)
eu_dashboard = await channel.declare_queue("dashboard.eu", durable=True, arguments=quorum)
alerts = await channel.declare_queue("alerts.critical", durable=True, arguments=quorum)

await audit_all_orders.bind(events, routing_key="order.#")
await eu_dashboard.bind(events, routing_key="*.*.eu")
await alerts.bind(events, routing_key="payment.failed.#")

Сообщение с ключом order.cancelled.eu попадёт в audit_all_orders (по order.#) и в eu_dashboard (по *.*.eu).

Когда брать: события с иерархической структурой, когда нужно гибко подписываться без переделки топологии при добавлении новых типов событий.

Порядок по ключу при нескольких обработчиках

Work Queue выше раздаёт сообщения нескольким обработчикам, и это ровно то, что нужно, пока сообщения независимы. Как только они связаны (два обновления одного заказа, создание и отмена одной брони), параллельная обработка ломает порядок: второе сообщение попадает к свободному обработчику раньше, чем первый закончил.

В Kafka порядок держит ключ партиции. У RabbitMQ порядок гарантирован только внутри одной очереди с одним потребителем, и из этого растут два приёма.

Single active consumer. Аргумент очереди x-single-active-consumer: true означает: подписчиков может быть много, но сообщения получает один, остальные ждут и мгновенно подхватывают работу, если активный отвалился. Порядок сохранён, отказоустойчивость есть, параллелизма нет. Это правильный выбор, когда очередь и так не перегружена, а порядок важен целиком.

bookings = await channel.declare_queue(
    "bookings",
    durable=True,
    arguments={"x-queue-type": "quorum", "x-single-active-consumer": True},
)

Consistent hash exchange. Плагин rabbitmq_consistent_hash_exchange добавляет обменник, который раскладывает сообщения по нескольким очередям по хешу ключа маршрутизации или заголовка. Сообщения одного заказа всегда попадают в одну очередь, у каждой очереди свой единственный потребитель, и получается то же, что партиции в Kafka: порядок внутри ключа плюс параллелизм между ключами. Цена ручное управление: очереди создаёте вы, число очередей меняется только с остановкой, а перекос по ключам («один крупный клиент») складывается в одну очередь.

Что не работает: «сделаем один потребитель, но в десять задач». Они разберут сообщения из общей выдачи и обработают их вперемешку, порядок теряется так же, как при нескольких потребителях. Если нужна и скорость, и порядок, ключ обязан определять, кто обрабатывает; других вариантов нет.

Запрос-ответ через очередь — RPC

Иногда нужен синхронный ответ, но HTTP не подходит: сервис за NAT, нет публичного адреса, или хочется балансировки по пулу обработчиков.

RPC через очередь: клиент отправляет запрос и ждёт ответа. Брокер доставляет запрос одному из обработчиков, тот отвечает в отдельную очередь-ответ. Для сопоставления запроса и ответа используется correlation-id.

В aio-pika это скрыто за паттерном RPC:

from aio_pika.patterns import RPC

# Клиент
rpc = await RPC.create(channel)
quote = await rpc.call("pricing.quote", kwargs={"request": request})

# Сервер
async def handle(*, request: QuoteRequest) -> PriceQuote:
    return PriceQuote.compute(request)  # возврат автоматически уходит в reply-to

rpc = await RPC.create(channel)
await rpc.register("pricing.quote", handle, auto_delete=True)

aio-pika сам создаёт временную очередь-ответ, проставляет reply-to и correlation-id, ждёт ответа. Возврат из зарегистрированного обработчика автоматически публикуется обратно.

Аргументы и ответ RPC по умолчанию сериализует через pickle. Внутри одной кодовой базы это удобно, но принимать pickle от чужого сервиса нельзя: распаковка выполняет произвольный код. Между сервисами разных команд берут JsonRPC из того же модуля или свой сериализатор.

Когда брать: нужен синхронный вызов, но HTTP не работает (NAT, firewall, нет публичного адреса); нужна балансировка запросов по пулу обработчиков.

Когда не брать: если HTTP/gRPC просто работает — RPC через брокер сложнее в отладке и дороже.

Что делать с повторной доставкой — Idempotent Consumer

AMQP гарантирует at-least-once доставку: одно и то же сообщение может прийти дважды. Это случается, когда брокер не получил подтверждение обработки (например, из-за сетевой проблемы) и переотправляет сообщение.

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

Ключ идемпотентности в базе данных

Самый надёжный способ — запоминать уже обработанные сообщения:

async def process(message: AbstractIncomingMessage) -> None:
    async with message.process():
        event = deserialize(message.body)
        async with session.begin():
            if await processed_events_repo.exists(event.idempotency_key):
                return  # уже обработано — просто подтверждаем получение
            await processed_events_repo.save(ProcessedEvent(event.idempotency_key))
            await account_repo.debit(event.account_id, event.amount)

Таблица processed_events с уникальным индексом на idempotency_key. Если два одинаковых сообщения придут одновременно — база данных поймает дубль через нарушение уникального ограничения.

Проверка состояния объекта

Если событие переводит объект в новое состояние, достаточно проверить текущее:

async def on_order_confirmed(event: OrderConfirmedEvent) -> None:
    async with session.begin():
        order = await order_repo.get(event.order_id)
        if order.status == OrderStatus.CONFIRMED:
            return  # уже в нужном состоянии
        order.confirm()
        await order_repo.save(order)

Не требует отдельной таблицы — состояние уже хранится в бизнес-объекте.

Retry с задержкой и Dead Letter Queue

Что если обработчик упал не из-за бага, а из-за временной недоступности внешнего сервиса? Нужно попробовать снова, но не немедленно.

Delayed retry собирается через x-message-ttl и Dead Letter Exchange:

retry_queue = await channel.declare_queue(
    "orders.retry",
    durable=True,
    arguments={
        "x-message-ttl": 30_000,  # ждём 30 секунд
        "x-dead-letter-exchange": "orders",
        "x-dead-letter-routing-key": "order.created",
        "x-queue-type": "quorum",
    },
)

Поток: обработчик отклонил сообщение → оно попадает в retry очередь → через 30 секунд по истечении TTL уходит через DLX обратно в основную очередь → новая попытка.

Чтобы отклонённое сообщение вообще попало в orders.retry, у основной очереди должен быть свой x-dead-letter-exchange — обменник, к которому привязана очередь повторов. Без него reject просто удалит сообщение.

Количество попыток считается через заголовок x-death.count — его нужно проверять вручную, встроенного ограничителя нет.

Ещё одна ловушка: очередь с TTL отдаёт сообщения строго по порядку, и сообщение с задержкой 5 секунд, вставшее за десятиминутным, дождётся его. Поэтому под каждую задержку заводят свою очередь повторов (orders.retry.5s, orders.retry.30s, orders.retry.10m) либо ставят плагин rabbitmq_delayed_message_exchange. Он включается на брокере (rabbitmq-plugins enable rabbitmq_delayed_message_exchange) и добавляет особый тип обменника, у которого задержка задаётся у каждого сообщения заголовком:

delayed = await channel.declare_exchange(
    "orders.delayed", "x-delayed-message", durable=True,
    arguments={"x-delayed-type": "direct"},
)

await delayed.publish(
    Message(body, headers={"x-delay": 30_000}, delivery_mode=DeliveryMode.PERSISTENT),  # задержка этого сообщения, мс
    routing_key="order.created",
)

Одна очередь, произвольные задержки, никакой головной блокировки. Цена тоже есть: плагин держит отложенные сообщения в своей встроенной базе на том узле, куда они опубликованы, это не очередь, и она не реплицируется. Падение узла до срока доставки теряет отложенные сообщения, а большое их число заметно ест память. Для повторов и коротких задержек плагин удобен, а для «напомнить через месяц» берут не брокер, а таблицу в базе и задачу по расписанию.

Сообщения, которые не получилось обработать после всех попыток, уходят в Dead Letter Queue (DLQ) — отдельную очередь для разбора вручную или через алерты.

Гарантированная публикация — Outbox

Бывает задача: сохранить заказ в базу данных и опубликовать событие — атомарно. Если сначала сохранить, потом опубликовать, то сервис может упасть между двумя операциями. Событие потеряется.

Outbox pattern: событие сохраняется в ту же транзакцию, что и бизнес-данные. Отдельный процесс читает таблицу и публикует в AMQP.

async def confirm(order_id: OrderId) -> None:
    async with session.begin():
        order = await order_repo.get(order_id)
        order.confirm()
        await order_repo.save(order)
        await outbox_repo.save(OutboxEvent(
            id=uuid.uuid4(),
            routing_key="order.confirmed",
            exchange="orders",
            payload=to_json(OrderConfirmedEvent(order_id)),
        ))

async def publish_outbox() -> None:  # фоновая задача
    while True:
        async with session.begin():
            batch = await outbox_repo.fetch_unpublished(100)
            for event in batch:
                exchange = await channel.get_exchange(event.exchange)
                await exchange.publish(
                    Message(event.payload.encode()), routing_key=event.routing_key
                )
                await outbox_repo.mark_published(event.id)
        await asyncio.sleep(0.5)

Либо оба изменения зафиксированы, либо ни одного. Дубли возможны (публикация прошла, но пометить как отправленное не успело) — поэтому получатель всё равно должен быть идемпотентным.

Шпаргалка по выбору

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

Глубже: соединение, подтверждения публикации и возвратырасширенное

Соединение с брокером открывают через aio_pika.connect_robust(...), а не connect: при обрыве сети оно само переподключается, заново объявляет очереди и обменники и восстанавливает подписки. Обычный connect после обрыва молча оставляет приложение без потребителей.

Канал в aio-pika создаётся с publisher_confirms=True: await exchange.publish(...) возвращается только после подтверждения брокером, что сообщение принято, а для надёжной очереди — записано. Поэтому publish в outbox-цикле выше можно считать «доставлено брокеру» и только после него помечать событие отправленным.

Сообщение, которое не подошло ни одной очереди, брокер молча отбрасывает, и подтверждение при этом всё равно приходит. Чтобы узнать о таком, публикуют с mandatory=True, а канал открывают с on_return_raises=True: тогда publish поднимет исключение вместо тихой потери.

Коротко

  • Work Queue — одна очередь, несколько обработчиков, каждое сообщение получает ровно один. Для фоновых задач.
  • Publish/Subscribe — fanout exchange копирует сообщение во все привязанные очереди. Для broadcast-событий.
  • Routing — direct exchange смотрит на ключ маршрутизации. Для точного разделения потоков.
  • Topic — как routing, но с шаблонами * и #. Для иерархических событий с гибкой подпиской.
  • RPC через очередь — запрос-ответ через брокер с reply-to и correlation-id. Для вызовов без HTTP.
  • Idempotent Consumer — at-least-once означает возможные дубли. Защита: ключ идемпотентности в БД или проверка состояния объекта.
  • Delayed Retry — TTL + DLX: сообщение «паркуется» на время, потом возвращается.
  • Outbox — событие сохраняется в той же транзакции, что и данные. Атомарность без двухфазного коммита.
  • Соединение через connect_robust; публикация ждёт подтверждения брокера; непривязанное сообщение ловят через mandatory=True.

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