On this page
RabbitMQ for agents in one paragraph
Put each LLM or tool job on a quorum queue as a persistent message with a unique ID, published with confirms. Run competing consumers with a prefetch of 1 and manual acknowledgements, sized to the provider's rate limit rather than to demand. Set a delivery limit and a dead letter exchange so a poison prompt fails a few times and stops. Route by intent or skill through a topic exchange, use reply-to and correlation IDs for tool calls, and keep a stream for replay. Make handlers idempotent, cap queue length, pass references instead of large payloads, and size the consumer timeout to the slowest legitimate call.
Why agents need a broker
An agent loop is mostly waiting. A model call takes seconds or minutes, a tool times out, and the provider throttles you when traffic spikes. Inline calls tie every caller to the slowest dependency. A broker breaks the coupling: the caller publishes a task and moves on, and workers take it when they have capacity.
That one change buys the rest. Failed calls are retried by redelivery, not hand-written loops. Backpressure is a queue that grows, not a thread pool that falls over. Fan-out to several tools is one publish. A task survives a worker crash because it was never acknowledged. And a human approval step is a queue a reviewer's tool consumes from, so the plan waits without holding a process open.
Run LLM jobs as a work queue
The basic shape is the work queue from the RabbitMQ tutorials: one queue, many competing consumers, each message delivered to one of them. Two settings make it safe for model calls. Use manual acknowledgements, because in automatic mode RabbitMQ treats a message as delivered the moment it is sent, and a worker that crashes mid-call loses it. And set the prefetch to 1. The documentation calls a prefetch of 1 the most conservative value and warns that it costs throughput, which is the right trade when each message is a long, billed call rather than a quick database write.
The consumer count then becomes your concurrency limit. If the provider allows 50 concurrent requests, run 50 consumers with a prefetch of 1 and let the queue absorb the rest. Throughput becomes a capacity decision made on purpose rather than a side effect of a traffic spike, which makes it one of the simplest cost controls available.
Make agent tasks durable with quorum queues
A customer request, a multi-step plan or a pending approval should survive a node failure. Declare those queues as quorum queues (x-queue-type set to quorum), publish the messages as persistent, and turn on publisher confirms so the producer knows the broker has the task before it reports the work as queued. Quorum queues are always durable and replicated, and they carry the two features agent work leans on most: a delivery limit, which defaults to 20 since RabbitMQ 4.0, and an at-least-once dead-lettering mode. They do not support global QoS, so set prefetch per consumer.
Stop poison prompts with a delivery limit and a dead letter queue
Some inputs will never succeed: a prompt the model refuses, a document the parser cannot read, a tool call whose arguments fail validation. Requeued forever, they loop and burn tokens. A quorum queue counts failed deliveries in the x-delivery-count header and, once the delivery limit is exceeded, drops the message, or dead-letters it if a dead letter exchange is configured. For billed model calls the default of 20 is generous; a limit of 3 to 5 keeps the cost of one bad input bounded.
Point the queue at a dead letter exchange and bind a queue behind it that a person or a triage agent reads. RabbitMQ adds an x-death header recording the source queue, the reason and the count, so the dead letter queue explains why each task landed there. When a worker already knows a retry is pointless, it should reject with requeue set to false and send the message straight there.
A worker to start from
The worker below uses aio-pika, the asyncio client for Python, so the connection keeps being serviced while a model call is awaited; pika's own documentation warns that a long-running callback can drop the connection on heartbeat timeout. It consumes a quorum queue with a delivery limit of 5, takes one message at a time, skips duplicates by message ID, backs off before returning a rate-limited job, and rejects anything that cannot succeed straight to the dead letter exchange. run_agent_step, already_done and save_result stand in for your own code.
import asyncio
import json
import aio_pika
class Transient(Exception):
"""Rate limit, timeout or 5xx from the model or a tool: worth retrying."""
async def handle(msg):
key = msg.message_id # stable per task, set by the publisher
if await already_done(key): # duplicate delivery: do not repeat the work
await msg.ack()
return
try:
job = json.loads(msg.body)
result = await run_agent_step(job) # the slow model or tool call
await save_result(key, result) # record the outcome before acknowledging
await msg.ack()
except Transient:
attempt = int((msg.headers or {}).get("x-delivery-count", 0))
await asyncio.sleep(min(2 ** attempt, 60)) # back off before handing it back
await msg.nack(requeue=True) # counts toward the delivery limit
except Exception:
await msg.reject(requeue=False) # cannot succeed: straight to the DLX
async def main():
conn = await aio_pika.connect_robust("amqps://agent-worker:PASSWORD@broker.example.com/agents")
ch = await conn.channel()
await ch.set_qos(prefetch_count=1) # one job in flight per worker
dlx = await ch.declare_exchange("agents.dlx", aio_pika.ExchangeType.FANOUT, durable=True)
dead = await ch.declare_queue("llm-jobs.dead", durable=True,
arguments={"x-queue-type": "quorum"})
await dead.bind(dlx)
jobs = await ch.declare_queue("llm-jobs", durable=True, arguments={
"x-queue-type": "quorum",
"x-delivery-limit": 5,
"x-dead-letter-exchange": "agents.dlx",
})
async with jobs.iterator() as messages:
async for msg in messages:
await handle(msg)
asyncio.run(main())The publisher sets the message ID the idempotency check depends on, marks the message persistent and gives it an expiry (aio-pika converts a timedelta to the milliseconds the protocol uses, and its channels have publisher confirms on by default).
from datetime import timedelta
await ch.default_exchange.publish(
aio_pika.Message(
json.dumps(job).encode(),
message_id=job["task_id"], # idempotency key
delivery_mode=aio_pika.DeliveryMode.PERSISTENT,
expiration=timedelta(minutes=2), # stale after two minutes
),
routing_key="llm-jobs",
)The queue arguments are inline so the example runs on its own. In production, set delivery-limit and dead-letter-exchange by policy: the documentation recommends policies because they can be changed without redeploying applications.
Tool calls and routing: reply-to and topic exchanges
When an agent calls a tool through the broker and needs the answer, use the request-reply pattern from the RabbitMQ RPC tutorial: set reply_to to a callback queue and correlation_id to a unique value per request, and have the tool publish its result to that queue with the same ID. Discard replies with a correlation ID you do not recognize. Direct reply-to (amq.rabbitmq.reply-to) avoids declaring a queue for short interactive calls, but its replies are not stored and are dropped if the requester disconnects, so anything you cannot afford to lose needs a durable reply queue.
To reach the right agent or tool, publish to a topic exchange with a routing key built from intent and skill, such as agent.billing.refund or tool.search.web. In a binding, * matches exactly one word and # zero or more, so specialists bind narrow keys while an audit consumer binds agent.#. A message goes to every queue whose binding matches, which is how one publish fans out to several tools. Routing keys are limited to 255 bytes.
Priorities, stale requests and long model calls
As of RabbitMQ 4.3, quorum queues support strict priorities with 32 levels (0 to 31) and no opt-in; classic queues use x-max-priority, ideally with a low single-digit number of levels. Publish user-facing requests above overnight summarization and workers serve them first.
A request nobody is waiting for is wasted spend. Give it a per-message TTL with the expiration property so a chat turn not started within, say, two minutes expires instead of running. Expired messages are only removed, or dead-lettered, once they reach the head of the queue, so a deep backlog still holds them until consumers get there.
Then the consumer timeout. RabbitMQ closes a channel with a PRECONDITION_FAILED error when a delivery is not acknowledged within 30 minutes by default, and requeues every delivery on that channel. A long agent run can hit that limit. Since RabbitMQ 3.12 the timeout can be set per queue, with the consumer-timeout policy key or the x-consumer-timeout argument in milliseconds, so raise it only on the queues that carry long jobs, with a margin above the slowest legitimate run.
Streams for replay and agent memory
A queue forgets a message once it is acknowledged. Agents often need the opposite: a record of every event a workflow produced, which a new agent can read from the start or a debugger can replay from a point in time. That is a RabbitMQ stream, a replicated, append-only log read without removing messages. Bind a stream to the exchange your agents already publish to, set retention with max-age or max-length-bytes, and let each reader choose its starting point with x-stream-offset: first, last, next, an offset, a timestamp or an interval. Streams do not support dead-lettering or TTL, so keep work on quorum queues and the history in the stream.
Frameworks and MCP
For Python teams, Celery is a direct route: RabbitMQ is its default broker, and setting task_default_queue_type to quorum moves the default queue onto quorum queues. Celery acknowledges a task just before running it by default; acks_late acknowledges afterwards, so a crash mid-task means it runs again, which is why Celery's documentation asks for idempotent tasks.
LangGraph, CrewAI and AutoGen orchestrate the agent itself, and the broker sits around them: a worker takes a task off a queue and invokes the graph or crew. LangGraph keeps graph state in a checkpointer keyed by a thread_id (its documentation lists in-memory, SQLite and Postgres savers), so carry the thread ID in the message and a redelivered task can resume from the last checkpoint. CrewAI crews start with kickoff() or its async variants. AutoGen's distributed runtime uses gRPC and is marked experimental. In each case RabbitMQ carries work between services; it does not replace the framework's runtime.
MCP, the Model Context Protocol, is an open standard for connecting AI applications to external tools and data, with stdio and Streamable HTTP as its standard transports. Community-built MCP servers for RabbitMQ exist, typically wrapping the management API and AMQP so an assistant can inspect queues or publish. Give any of them its own RabbitMQ user with the narrowest permissions that work, because it acts with whatever access it holds.
Pitfalls that show up in production
- Unbounded queues when the model API throttles. Once the provider returns rate-limit errors, consumers slow down and the queue grows until a memory or disk alarm blocks every publishing connection on the node. Cap queues that take user traffic with
max-lengthand thereject-publishoverflow, so a publisher using confirms gets abasic.nackand can ask the user to retry. The default overflow,drop-head, discards the oldest message, or dead-letters it if a dead letter exchange is set. - Long runs hitting the consumer timeout. Size
consumer-timeoutper queue, keep prefetch at 1 for slow work, and acknowledge on every code path, including errors. - Large payloads. The default maximum message size is 16 MiB (
max_message_size, configurable up to 512 MiB); larger messages are rejected. Keep documents, images and embedding batches in object storage or a vector database and put a reference in the message. - Ordering assumptions. A queue keeps publication order, but with several consumers and requeues each consumer can see messages out of order. For strict per-conversation order, use a queue with a single active consumer, or carry a sequence number.
- Retries that repeat side effects. Delivery is at least once: a worker can finish the work and crash before acknowledging, and the message returns with the
redeliveredflag set. Key every side effect (an email sent, a ticket opened, a refund issued) on the task's message ID and check it first, as the worker above does.
Frequently asked questions
Can RabbitMQ be used for AI agents?
Yes. Workers consume LLM and tool jobs with manual acknowledgements, quorum queues keep tasks durable, delivery limits and dead letter exchanges stop poison prompts, topic exchanges route by skill, and streams keep a replayable history.
Which message queue should I use for AI agents?
It depends on the workload. Task hand-off with retries, priorities and per-task acknowledgement suits a queue broker such as RabbitMQ; high-volume event streams that many consumers replay suit a log, such as Kafka or RabbitMQ streams.
How do I stop long LLM calls hitting the RabbitMQ consumer timeout?
The default timeout is 30 minutes. On RabbitMQ 3.12 and later, set consumer-timeout by policy, or x-consumer-timeout as a queue argument, on the queues that carry long jobs, keep prefetch at 1, and acknowledge on every code path.
Does LangChain or LangGraph work with RabbitMQ?
They are used together, with RabbitMQ as the task layer around the graph. A worker consumes a task and invokes the graph, while LangGraph's checkpointer, which its documentation lists as in-memory, SQLite or Postgres, keeps the graph state under a thread_id carried in the message.
How should an AI agent handle duplicate RabbitMQ messages?
Assume at-least-once delivery. Give every task a stable message ID, record completed IDs alongside the result, and skip any delivery whose ID is already recorded before repeating a side effect.
Related
Where this gets done
The work behind this page, run by the same engineers who wrote it.
- 24/7 RabbitMQ support15-minute emergency SLA, versions back to 3.8.x
- Managed RabbitMQ servicesWe run the brokers, on your infrastructure or hosted
- RabbitMQ consultingArchitecture, migration and remediation from senior engineers
- RabbitMQ health checkEngineer-led assessment with a prioritised fix list
- Extended LTS support for RabbitMQ 3.xCVE backports for versions the community no longer patches
- RabbitMQ commercial licensingTanzu RabbitMQ licences from an authorized Broadcom partner
- RabbitMQ troubleshootingLive incidents and recurring faults
- RabbitMQ upgrades3.x to 4.x, planned and executed in your window
- RabbitMQ migrationsFrom IBM MQ, Kafka, cloud brokers or older RabbitMQ
- RabbitMQ implementation and architectureCluster design, DR and go-live
- RabbitMQ corporate trainingAdmin and developer courses taught by working engineers
Other RabbitMQ guides, comparisons and research
- GuideThe RabbitMQ Performance Tuning GuideRead the guide
- GuideThe RabbitMQ Reliability Guide: Ten Failure Patterns and Their FixesRead the guide
- GuideThe RabbitMQ Disaster Recovery GuideRead the guide
- GuideThe RabbitMQ Clustering and Sizing GuideRead the guide
- GuideThe RabbitMQ on Kubernetes GuideRead the guide
- GuideThe RabbitMQ Migration GuideRead the guide
- GuideThe RabbitMQ Security and Hardening GuideRead the guide
- GuideThe RabbitMQ Monitoring and Alerting GuideRead the guide
- ComparisonRabbitMQ vs Amazon SQS ComparedSee the comparison
- ComparisonRabbitMQ vs Redis ComparedSee the comparison
- ComparisonManaged RabbitMQ Options ComparedSee the comparison
- ComparisonMessage Broker Support Options ComparedSee the comparison
- ResearchWhat Breaks in Production RabbitMQ: 145 Support Tickets, 2023 to 2026Read the research
- ResearchRabbitMQ in Production 2026: What 22 Assessed Estates Actually RunRead the research
- ResearchThe RabbitMQ CVE Register, 2026 EditionRead the research
Recent RabbitMQ articles
- RabbitMQ Cluster Operator vs Helm Chart on KubernetesOct 2026
- RabbitMQ Alternatives: 9 Options and When to StayOct 2026
- RabbitMQ End of Life and End of Support Dates (3.6 to 4.3)Oct 2026
- VMware Licensing Cost in 2026Sep 2026
- The Tanzu Software in Your VCF You Are Not UsingSep 2026
- Production RabbitMQ Architecture for Enterprise TeamsSep 2026
Need this done on your cluster?
AceMQ's senior RabbitMQ engineers support 130+ enterprise clients in 26+ countries under a 15-minute emergency SLA, with direct escalation to the RabbitMQ core team.