Tuning RabbitMQ Prefetch for Async Consumers¶
In an aio-pika consumer, the channel's prefetch count is two settings at once. It is how many unacknowledged messages RabbitMQ will push to the consumer, and — because aio-pika runs a handler for each delivered message concurrently — it is also the consumer's concurrency. Measured with RabbitMQ 4.3 and handlers that awaited 10 ms of simulated I/O, 2,000 messages were consumed at 90 messages per second with prefetch_count=1, 910 at 10, 4,158 at 50 and 12,579 at 200; the peak number of handlers running at once matched the prefetch exactly. But prefetch also decides how work is shared: with a fast and a slow consumer on one queue, each processing at most 5 messages at a time, prefetch 200 let the slow consumer hoard 210 of 600 messages and the queue took 4.24 s to drain; prefetch 5 gave the slow one 35, and the queue drained in 0.71 s. This guide sets prefetch from both effects.
Prerequisites¶
- Python 3.11+,
pip install aio-pika; measured with RabbitMQ 4.3. - A consumer, from processing RabbitMQ messages with aio-pika.
- Concurrency limits, from limiting concurrent requests with asyncio.Semaphore.
1. Know that prefetch is your concurrency¶
Set prefetch on the channel before consuming. Every delivered, unacknowledged message has a handler running for it:
channel = await connection.channel()
await channel.set_qos(prefetch_count=50) # at most 50 unacked messages, 50 handlers
queue = await channel.declare_queue("work", durable=True)
async def handle(message: aio_pika.IncomingMessage) -> None:
async with message.process(): # acks on success, rejects on exception
await call_downstream(message.body) # ~10 ms of I/O
await queue.consume(handle)
Measured: with 10 ms handlers, throughput was close to prefetch / 0.01 s until other costs dominated — 90 msg/s at 1, 910 at 10, 4,158 at 50, 12,579 at 200. With prefetch 1, an async consumer behaves like a synchronous one: it handles one message, acknowledges it, and only then receives the next. Without set_qos at all, RabbitMQ pushes the entire queue to the consumer, and aio-pika starts a handler for every message — unbounded concurrency and memory.
Verify: export the number of running handlers; its peak equals the prefetch count under load.
2. Size it from what the handler waits on¶
Use Little's law: the concurrency you need is the target throughput times the time each message spends in the handler. Then cap it by what the handler's dependencies can take:
target_rate = 500 # messages per second this consumer should handle
handler_time = 0.040 # seconds per message, measured p50, mostly I/O
needed = target_rate * handler_time # 20 in flight
db_pool_size = 10 # each handler holds one connection while it works
prefetch = min(int(needed * 1.5), db_pool_size * 2) # headroom, but not far beyond the pool
await channel.set_qos(prefetch_count=prefetch)
If every handler needs a database connection from a pool of 10, a prefetch of 200 does not make 200 handlers useful: 190 of them wait on the pool while holding messages that another consumer could be processing. Size prefetch near the real concurrency limit downstream — the pool, an HTTP client's connection limit, an API's rate limit — as in sizing async connection pools for throughput.
Verify: with the chosen prefetch, the downstream pool's wait time stays low and the queue drains at the target rate.
3. Keep slow consumers from hoarding¶
Prefetched messages belong to that consumer until it acknowledges or disconnects; other consumers cannot take them. A high prefetch on a slow consumer parks work there:
# Two consumers on one queue, each limited internally to 5 concurrent handlers
fast_channel = await conn.channel(); await fast_channel.set_qos(prefetch_count=200)
slow_channel = await conn.channel(); await slow_channel.set_qos(prefetch_count=200)
# measured: slow got 210 of 600 messages; the queue drained in 4.24 s
# Same consumers, prefetch 5: slow got 35 of 600; drained in 0.71 s
The slow consumer in the test took 100 ms per message against 5 ms for the fast one. With prefetch 200, it grabbed a large share up front and worked through it slowly while the fast consumer sat idle at the end. With prefetch 5 — matching what each could actually run at once — the fast consumer took most of the work and the queue drained six times sooner. Mixed fleets are common: different instance sizes, a node with a noisy neighbour, a consumer stuck on a slow dependency.
Verify: during a drain test with mixed consumers, no consumer holds many more unacked messages than it is actively processing.
4. Separate prefetch from concurrency when they must differ¶
Sometimes you want a small amount of buffering (to hide network latency between deliveries) but a smaller, fixed number of handlers — for example, CPU-heavy steps, or a strict downstream limit. Put a semaphore inside the handler:
work_slots = asyncio.Semaphore(8) # at most 8 running
async def handle(message: aio_pika.IncomingMessage) -> None:
async with work_slots:
async with message.process():
await process(message.body)
await channel.set_qos(prefetch_count=12) # 8 running + 4 buffered
Keep the buffer small: every buffered message is invisible to other consumers and is redelivered — out of its original order — if this consumer dies. A prefetch of "concurrency plus a few" keeps the pipeline fed without hoarding. For CPU-bound processing, offload to a process pool and size both numbers from the pool, as in Concurrent Execution & Worker Patterns.
Verify: the number of running handlers never exceeds the semaphore, and unacked messages never exceed the prefetch.
5. Account for poison messages, ordering and shutdown¶
Prefetch interacts with three other behaviours:
async def shutdown(channel, queue, consumer_tag) -> None:
await queue.cancel(consumer_tag) # stop new deliveries
await drain_running_handlers() # let in-flight handlers ack
await channel.close() # unacked remainder returns to the queue
A poison message that is requeued returns to the queue and can be delivered again immediately; with prefetch 1, a requeue loop blocks the consumer completely, and with higher prefetch it consumes one slot, as measured in handling poison messages in async consumers. Concurrency means messages finish out of order; if order matters per key, prefetch must be 1 per ordered stream, or ordering must be restored by key. And on shutdown, cancel the consumer first so no new messages arrive, then let running handlers finish — anything still unacknowledged when the channel closes is redelivered to another consumer.
Verify: a rolling restart under load produces no lost messages and only the redeliveries of messages that were unacknowledged at shutdown.
Verification¶
Prefetch is tuned when:
- It is always set, never left unbounded.
- It matches the real concurrency the handler's dependencies support.
- Slow consumers do not hoard messages, shown by a drain test with mixed consumers.
- Shutdown cancels consumption first and lets in-flight handlers finish.
Diagnostic Hook: compare each consumer's unacked message count (RabbitMQ management API) with its number of running handlers. Unacked far above running means prefetch is hoarding; running pinned at prefetch while the queue grows means the consumer needs more prefetch or more instances; running pinned while a downstream pool shows waits means prefetch is already above what the dependency allows.
Pitfalls & edge cases¶
- No
set_qos. The whole queue is pushed and a handler starts for every message. - Prefetch 1 with async handlers. Measured at 90 msg/s against 12,579 at 200.
- Very high prefetch on mixed consumers. Measured: 4.24 s to drain against 0.71 s.
- Prefetch above the downstream limit. Handlers queue on the pool while holding messages.
Frequently Asked Questions¶
What prefetch count should I use with aio-pika?
Start from target throughput times handler duration, then cap it near the concurrency your dependencies allow. In testing with 10 ms handlers, prefetch 50 gave 4,158 msg/s and 200 gave 12,579.
Does aio-pika process messages concurrently?
Yes. Each delivered message gets its own handler, so the prefetch count is the maximum number of handlers running at once; the measured peak matched the prefetch exactly.
Why is one RabbitMQ consumer getting most of the messages?
A high prefetch lets a consumer take many messages up front, even if it processes them slowly. With prefetch 200, a slow consumer held 210 of 600 messages; with prefetch 5 it held 35.
What happens if I don't set prefetch in aio-pika?
RabbitMQ delivers without limit, so a large queue is pushed into the consumer's memory and a handler is started for every message.
Related¶
- Message Brokers & Event Streams — up to the topic overview.
- Using NATS from asyncio with nats-py — pull consumers as an alternative form of backpressure.
- Network I/O & Protocol Handling — the section overview.