RabbitMQ с aio-pika: producer, consumer, ack, DLQ

В прошлом уроке узнали про message brokers вообще. Теперь конкретика - aio-pika, asynchronous RabbitMQ client для Python. Подходит для FastAPI и других async stack. В этом уроке - подключение, publisher, consumer, acknowledgements, prefetch и dead letter queue.

Установка и базовая структура

pip install aio-pika

Запустим RabbitMQ через Docker:

docker run -d --name rabbit -p 5672:5672 -p 15672:15672 rabbitmq:3-management

Управление через web UI на http://localhost:15672 (guest/guest).

Подключение

import aio_pika

connection = await aio_pika.connect_robust("amqp://guest:guest@localhost/")
channel = await connection.channel()
  • connect_robust - auto-reconnect при потере связи
  • channel - канал для операций (queue.declare, publish, consume)

Channel это logical connection - один TCP может иметь много channels. Channel cheap, создавай по одному на consumer/producer.

Простой producer

import aio_pika
import asyncio

async def main():
    connection = await aio_pika.connect_robust("amqp://guest:guest@localhost/")
    async with connection:
        channel = await connection.channel()
        await channel.declare_queue("hello", durable=True)

        message = aio_pika.Message(
            body=b"Hello, RabbitMQ!",
            delivery_mode=aio_pika.DeliveryMode.PERSISTENT,
        )
        await channel.default_exchange.publish(message, routing_key="hello")
        print("Sent")

asyncio.run(main())
  • declare_queue("hello", durable=True) - создаёт queue (idempotent), durable - переживёт restart broker
  • delivery_mode=PERSISTENT - сообщение записано на диск, не теряется при restart
  • default_exchange (имя "") - direct routing к queue с именем как routing_key

Простой consumer

async def consumer():
    connection = await aio_pika.connect_robust("amqp://guest:guest@localhost/")
    async with connection:
        channel = await connection.channel()
        await channel.set_qos(prefetch_count=10)   # max 10 unacknowledged at once

        queue = await channel.declare_queue("hello", durable=True)

        async with queue.iterator() as queue_iter:
            async for message in queue_iter:
                async with message.process():
                    print(f"Received: {message.body.decode()}")
                    # message.process() автоматически ack при exit
                    # rollback (nack/reject) если exception

asyncio.run(consumer())
  • set_qos(prefetch_count=10) - не брать в обработку >10 unacknowledged messages
  • queue.iterator() - async iterator поверх consume
  • message.process() - context manager: ack on success, nack on exception

Manual acknowledgement

Для finer control:

async for message in queue_iter:
    try:
        # process
        result = process(message.body)
        await message.ack()
    except SomeRecoverableError:
        await message.nack(requeue=True)   # retry
    except Exception:
        await message.reject(requeue=False)   # discard или в DLQ
  • ack() - подтверждение, broker удаляет message
  • nack(requeue=True) - вернуть в queue для retry
  • nack(requeue=False) - не возвращать, в DLQ если настроено

Exchanges и routing

Direct exchange - точное соответствие routing key:

# Producer
exchange = await channel.declare_exchange("logs", aio_pika.ExchangeType.DIRECT)
await exchange.publish(
    aio_pika.Message(b"error message"),
    routing_key="error",
)

# Consumer для error queue
queue = await channel.declare_queue("error_queue")
await queue.bind(exchange, routing_key="error")

Topic exchange - wildcard pattern:

exchange = await channel.declare_exchange("logs.topic", aio_pika.ExchangeType.TOPIC)

# Publisher
await exchange.publish(
    aio_pika.Message(b"order created"),
    routing_key="order.created",
)
await exchange.publish(
    aio_pika.Message(b"order updated"),
    routing_key="order.updated",
)

# Consumer всех order events
queue = await channel.declare_queue("order_events")
await queue.bind(exchange, routing_key="order.#")   # # = любые слова
# или routing_key="order.*"  # * = одно слово

Topic patterns:

  • order.* - матчит order.created, order.deleted (одно слово после order.)
  • order.# - матчит order.created, order.user.banned (несколько слов)
  • *.error - любое слово + .error

Fanout - broadcast всем queues:

exchange = await channel.declare_exchange("broadcast", aio_pika.ExchangeType.FANOUT)
# Все bound queues получают копию каждого message

Полезно для system-wide notifications, cache invalidation.

Dead Letter Exchange

Сообщения, которые reject (без requeue) или TTL expired - в Dead Letter Exchange:

# Setup
async with connection:
    channel = await connection.channel()

    # DLX и DLQ
    dlx = await channel.declare_exchange("orders.dlx", aio_pika.ExchangeType.DIRECT)
    dlq = await channel.declare_queue("orders.failed", durable=True)
    await dlq.bind(dlx, routing_key="failed")

    # Main queue с настройкой DLX
    main_queue = await channel.declare_queue(
        "orders",
        durable=True,
        arguments={
            "x-dead-letter-exchange": "orders.dlx",
            "x-dead-letter-routing-key": "failed",
            "x-message-ttl": 60000,   # 60 sec - после reject
        },
    )

При reject (или TTL expired) сообщение автоматически попадает в DLQ для inspection. Как DLQ устроен со стороны самого RabbitMQ - в отдельном уроке трека message-brokers.

Producer с confirms

Гарантия что message достиг broker:

async with connection:
    channel = await connection.channel(publisher_confirms=True)

    await channel.default_exchange.publish(
        aio_pika.Message(b"important"),
        routing_key="critical_queue",
        mandatory=True,   # if no route - returns to publisher
    )
    # Если broker не подтвердил - exception

publisher_confirms - acknowledge от broker. mandatory - error если нет matching queue.

Headers exchange

Routing по headers, не routing key:

exchange = await channel.declare_exchange("attrs", aio_pika.ExchangeType.HEADERS)

queue = await channel.declare_queue("priority_orders")
await queue.bind(exchange, arguments={
    "x-match": "all",   # требуются все headers
    "format": "json",
    "priority": "high",
})

await exchange.publish(
    aio_pika.Message(b"data", headers={"format": "json", "priority": "high"}),
    routing_key="",   # routing key игнорируется в headers exchange
)

Reдко используется - сложнее managment чем topic exchanges.

Connection pooling и channel pooling

Для high-throughput producers - reuse:

connection_pool = []
channel_pool = []

async def publish(message_body):
    async with channel_pool.acquire() as channel:
        await channel.default_exchange.publish(
            aio_pika.Message(message_body),
            routing_key="tasks",
        )

aio-pika имеет Pool для этого:

from aio_pika.pool import Pool

async def get_connection():
    return await aio_pika.connect_robust("amqp://...")

connection_pool: Pool = Pool(get_connection, max_size=10)

async def get_channel() -> aio_pika.Channel:
    async with connection_pool.acquire() as conn:
        return await conn.channel()

channel_pool: Pool = Pool(get_channel, max_size=20)

Использование в FastAPI

Подключение через lifespan и Depends:

from fastapi import FastAPI, Depends
from contextlib import asynccontextmanager

@asynccontextmanager
async def lifespan(app: FastAPI):
    app.state.broker = await aio_pika.connect_robust("amqp://...")
    yield
    await app.state.broker.close()

app = FastAPI(lifespan=lifespan)

async def get_broker(request):
    return request.app.state.broker

@app.post("/orders")
async def create_order(data: OrderCreate, broker = Depends(get_broker)):
    order = await save_to_db(data)
    channel = await broker.channel()
    try:
        await channel.default_exchange.publish(
            aio_pika.Message(json.dumps({"order_id": order.id}).encode()),
            routing_key="orders.created",
        )
    finally:
        await channel.close()
    return order

Standalone consumer (worker)

Consumer обычно отдельный процесс/контейнер:

# worker.py
import asyncio
import aio_pika
import json
import logging

logger = logging.getLogger(__name__)

async def process_order(message_body: bytes):
    data = json.loads(message_body.decode())
    order_id = data["order_id"]

    # Real processing
    await send_confirmation_email(order_id)
    await update_inventory(order_id)
    await notify_warehouse(order_id)

async def main():
    connection = await aio_pika.connect_robust("amqp://...")
    async with connection:
        channel = await connection.channel()
        await channel.set_qos(prefetch_count=20)

        queue = await channel.declare_queue("orders.created", durable=True)

        async with queue.iterator() as queue_iter:
            async for message in queue_iter:
                async with message.process(requeue=True):
                    try:
                        await process_order(message.body)
                    except Exception:
                        logger.exception("Failed to process order")
                        raise   # триггер nack(requeue) - retry

if __name__ == "__main__":
    asyncio.run(main())

Запускается отдельно от web server. В Docker - отдельный контейнер. В Kubernetes - отдельный Deployment.

Graceful shutdown

import signal

async def main():
    connection = await aio_pika.connect_robust("amqp://...")

    stop_event = asyncio.Event()

    def signal_handler():
        stop_event.set()

    loop = asyncio.get_running_loop()
    for sig in [signal.SIGTERM, signal.SIGINT]:
        loop.add_signal_handler(sig, signal_handler)

    async with connection:
        channel = await connection.channel()
        queue = await channel.declare_queue("orders.created", durable=True)

        async with queue.iterator() as queue_iter:
            async for message in queue_iter:
                if stop_event.is_set():
                    break
                async with message.process():
                    await process_order(message.body)

При SIGTERM (Docker stop, k8s rolling update) consumer заканчивает текущую обработку и shutdown gracefully. Сообщения не теряются.

Идемпотентность

Поскольку at-least-once, consumer может получить одно message дважды (retry, restart). Делай idempotent:

async def process_order(message_body: bytes):
    data = json.loads(message_body.decode())
    order_id = data["order_id"]
    message_id = data.get("message_id", str(uuid.uuid4()))

    # Check уже processed
    if await db.exists("processed_messages", message_id):
        logger.info(f"Skip already processed: {message_id}")
        return

    # Atomic: process + mark
    async with db.transaction():
        await process(data)
        await db.insert("processed_messages", message_id)

Подробно про идемпотентность - в следующем уроке.

Распространённые ошибки

1. Forget durable=True

queue = await channel.declare_queue("orders")   # not durable

После restart broker queue исчезает с сообщениями. Always durable=True для важных queues.

2. Сообщения не persistent

await channel.default_exchange.publish(
    aio_pika.Message(b"data"),   # not persistent
    routing_key="orders",
)

Без delivery_mode=PERSISTENT сообщения только в memory. Restart broker = потеря. Always PERSISTENT для важных.

3. Высокий prefetch без причины

await channel.set_qos(prefetch_count=10000)

Если consumer медленный, накопит много unacked. При crash все потеряны (will be requeued, но возможны дубли). Низкий prefetch (10-100) для балансировки.

4. Acknowledge до фактической обработки

async with queue.iterator() as queue_iter:
    async for message in queue_iter:
        await message.ack()   # ACK перед обработкой!
        await process(message.body)   # если упадёт - message потерян

ACK ПОСЛЕ успешной обработки. Иначе message потеряется если crash во время processing.

5. Не handle connection drops

connection = await aio_pika.connect("amqp://...")   # не robust

Используй connect_robust для auto-reconnect. Иначе network blip и connection dead.

Production tips

  • Используй connect_robust для resilience
  • durable queues + persistent messages для важных данных
  • Reasonable prefetch_count (10-100)
  • DLQ для failed messages
  • Graceful shutdown handling
  • Idempotent consumers
  • Monitor queue depth (Prometheus + RabbitMQ exporter)
  • Backup планов: если broker недоступен, не блокировать requests (publish failure → fallback в DB outbox)

Сравнение с Celery

Celery высокоуровневый task queue framework над брокерами (RabbitMQ, Redis):

from celery import Celery

app = Celery("tasks", broker="amqp://...")

@app.task
def send_email(user_id):
    # ...

# Использование
send_email.delay(123)   # async

Преимущества Celery:

  • Retries с backoff
  • Scheduling (periodic tasks через Celery Beat)
  • Progress tracking
  • Result backend

Преимущества raw aio-pika:

  • Tighter control
  • Async-first
  • Меньше зависимостей
  • Лучше для event-driven (не tasks)

Для backend в FastAPI обычно aio-pika если event-driven, Celery если task-heavy workflow.

Сравнение с Go и PHP

В Go - streadway/amqp (legacy) или rabbitmq/amqp091-go (fork). Async через горутины:

ch.QueueDeclare("hello", true, false, false, false, nil)
ch.Publish("", "hello", false, false, amqp.Publishing{
    Body: []byte("hello"),
})

В PHP - php-amqplib популярный. Также enqueue PSR абстракция.

aio-pika в Python специально для async. Принципы AMQP одинаковы во всех языках - knowledge transferable.

Мини-задание

  1. Producer и consumer в одном файле:
import asyncio
import aio_pika

async def producer():
    connection = await aio_pika.connect_robust("amqp://guest:guest@localhost/")
    async with connection:
        channel = await connection.channel()
        await channel.declare_queue("test_queue", durable=True)

        for i in range(10):
            await channel.default_exchange.publish(
                aio_pika.Message(
                    body=f"Message {i}".encode(),
                    delivery_mode=aio_pika.DeliveryMode.PERSISTENT,
                ),
                routing_key="test_queue",
            )
        print("Published 10 messages")

async def consumer():
    connection = await aio_pika.connect_robust("amqp://guest:guest@localhost/")
    async with connection:
        channel = await connection.channel()
        await channel.set_qos(prefetch_count=5)

        queue = await channel.declare_queue("test_queue", durable=True)

        async with queue.iterator() as queue_iter:
            count = 0
            async for message in queue_iter:
                async with message.process():
                    print(f"Got: {message.body.decode()}")
                    count += 1
                    if count >= 10:
                        break

# Запуск (по очереди):
asyncio.run(producer())
asyncio.run(consumer())
  1. Topic routing:
async def setup_topics():
    connection = await aio_pika.connect_robust("amqp://...")
    async with connection:
        channel = await connection.channel()

        exchange = await channel.declare_exchange("events", aio_pika.ExchangeType.TOPIC)

        # Все order events
        all_orders = await channel.declare_queue("all_orders")
        await all_orders.bind(exchange, "order.#")

        # Только cancellations
        cancellations = await channel.declare_queue("cancellations")
        await cancellations.bind(exchange, "*.cancelled")

        # Publish
        await exchange.publish(aio_pika.Message(b"order created"), routing_key="order.created")
        await exchange.publish(aio_pika.Message(b"order cancelled"), routing_key="order.cancelled")

Что дальше

Освоили RabbitMQ. В следующем уроке - Kafka с aiokafka: высокопроизводительный event streaming для big data scenarios. А retry, идемпотентность и graceful shutdown, общие для обоих брокеров, разбираем в уроке про consumer patterns.

Зарегистрируйтесь бесплатно, чтобы пройти квиз, решить задание с автопроверкой и вести прогресс.