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 brokerdelivery_mode=PERSISTENT- сообщение записано на диск, не теряется при restartdefault_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 messagesqueue.iterator()- async iterator поверх consumemessage.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 удаляет messagenack(requeue=True)- вернуть в queue для retrynack(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 durablequeues +persistentmessages для важных данных- 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.
Мини-задание
- 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())
- 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.