Guide · RabbitMQ

RabbitMQ for AI Agents: Patterns, Setup and Pitfalls

AI agents spend most of their time waiting on model APIs and tools that are slow, rate-limited and sometimes down. RabbitMQ handles that kind of work well: it holds tasks durably, hands them to as many workers as the rate limit allows, retries what fails and parks what never will. This guide sets out the patterns in the order you build them, with a worker to start from and the pitfalls that appear in production. Broker behavior is from the RabbitMQ 4.3 documentation, checked on 10 October 2026.

Tyler Eastridge

By Tyler Eastridge, Head of Operations

LinkedIn · Updated

10 min read10 sections
On this page
RabbitMQ for agents in one paragraph

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-length and the reject-publish overflow, so a publisher using confirms gets a basic.nack and 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-timeout per 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 redelivered flag 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.

RabbitMQ services

Where this gets done

The work behind this page, run by the same engineers who wrote it.

More resources

Other RabbitMQ guides, comparisons and research

From the blog

Recent RabbitMQ articles

Next step

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.