Когда сервисы начинают общаться через брокер, возникает вопрос: как правильно организовать очереди, обменники и подписки? Большинство задач укладываются в несколько типовых схем. Разберём каждую на примерах aio-pika.
Раздать задачи нескольким обработчикам — Work Queue
Представьте: пользователи загружают фотографии, и каждую надо сжать и нарезать в несколько размеров. Это долго. Делать это прямо в HTTP-запросе нельзя — пользователь будет ждать минуты.
Решение — work queue (очередь задач): сохранить задание в очередь и отдать её одному из пула обработчиков.
Очередь одна, потребителей несколько: брокер раздаёт сообщения по одному, и каждое достаётся ровно кому-то одному. Больше потребителей — быстрее разбирается очередь, но общий порядок между ними уже не сохраняется.
Брокер сам распределяет задачи между обработчиками. Каждое сообщение получает ровно один из них.
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 Queue | direct (default) |
| Broadcast события всем сервисам | Publish/Subscribe | fanout |
| Разные события в разные очереди | Routing | direct |
| Подписка по маске на иерархические события | Topic | topic |
| Синхронный вызов через очередь | RPC | direct + reply-to |
| Защита от повторной доставки | Idempotent Consumer | любой |
| Retry с задержкой | Delayed Retry | direct + DLX |
| Атомарная публикация вместе с записью в БД | Outbox | direct |
Глубже: соединение, подтверждения публикации и возвратырасширенное
Соединение с брокером открывают через 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.
Что почитать дальше
- Протокол AMQP — модель exchange/binding/queue изнутри.
- Распределённые паттерны на Python — Outbox, Saga и Idempotent Consumer подробнее.
- RabbitMQ в production — Quorum Queues, кластеризация, мониторинг.
- AMQP vs Kafka — какой брокер и когда брать.