Using a Message Queue for Agent Coordination
Most agent tutorials wire agents together with direct calls, and it's wrong at scale. An agent message queue decouples work from workers and survives restarts. Here's how and when.
class AgentQueue:
def __init__(self):
self._queue = deque()
self._in_flight = {} # msg_id -> message, until acked
def push(self, task):
msg_id = new_id()
self._queue.append((msg_id, task))
return msg_id
def pull(self):
msg_id, task = self._queue.popleft()
self._in_flight[msg_id] = task # held until ack
return msg_id, task
def ack(self, msg_id):
self._in_flight.pop(msg_id, None)
def nack(self, msg_id):
task = self._in_flight.pop(msg_id)
self._queue.append((msg_id, task)) # requeue on failureMost agent tutorials wire agents together with direct function calls, and for a demo that's fine. At scale it's the wrong architecture, and I'll defend that claim. Direct calls mean the caller waits for the callee, knows the callee exists, and loses all in-flight work if either crashes. An agent message queue fixes all three problems at once - it decouples the work from the worker, so agents don't wait on each other, don't need to know each other, and don't lose work when something restarts.
This isn't about adding infrastructure for its own sake. It's about a specific set of problems that direct coordination can't solve and a queue solves almost for free. Let me show you what an agent message queue buys you, and - just as important - when you don't need one.
What Is an Agent Message Queue?
An agent message queue is a durable buffer between agents that produce work and agents that consume it. A producer agent pushes a message (a task) onto the queue and moves on. A consumer agent pulls the next message when it's ready, processes it, and acknowledges completion. The queue holds messages until they're acknowledged, so nothing is lost if a consumer crashes mid-task.
The core shift is from push to pull, and from synchronous to asynchronous. With direct calls, a coordinator pushes work to a specific worker and waits. With a queue, the coordinator drops work into the queue and any available worker pulls it. The producer doesn't know or care which consumer handles the message.
[object Object], ,[object Object],:
,[object Object], ,[object Object],(,[object Object],):
,[object Object],._queue = deque()
,[object Object],._in_flight = {} ,[object Object],
,[object Object], ,[object Object],(,[object Object],):
msg_id = new_id()
,[object Object],._queue.append((msg_id, task))
,[object Object], msg_id
,[object Object], ,[object Object],(,[object Object],):
msg_id, task = ,[object Object],._queue.popleft()
,[object Object],._in_flight[msg_id] = task ,[object Object],
,[object Object], msg_id, task
,[object Object], ,[object Object],(,[object Object],):
,[object Object],._in_flight.pop(msg_id, ,[object Object],)
,[object Object], ,[object Object],(,[object Object],):
task = ,[object Object],._in_flight.pop(msg_id)
,[object Object],._queue.append((msg_id, task)) ,[object Object],What this does: It implements the essential queue contract - push work, pull it, acknowledge success, and requeue on failure. The in_flight tracking is what makes it durable: a message pulled but never acked (because the consumer crashed) can be requeued, so no task silently vanishes.
Why It Matters
Direct coordination fails in three ways that a queue quietly handles.
First, load imbalance. With direct calls to named workers, one worker can be swamped while others idle. A queue naturally balances - free workers pull the next message, so work distributes to whoever's available. No routing logic required.
Second, crash resilience. If a directly-called worker crashes mid-task, the task is gone. With a queue, an unacknowledged message returns to the queue and another worker picks it up. The task survives the worker.
Third, backpressure. When work arrives faster than agents can process it, direct calls either block the producer or drop work. A queue absorbs the burst, holding messages until consumers catch up. This is the difference between a system that degrades gracefully under load and one that falls over.
Three scenarios where this matters:
A logistics company ran an agent team processing shipment exceptions. During peak hours, exceptions arrived in bursts. Direct coordination meant bursts overwhelmed the agents and tasks were dropped. An agent message queue absorbed the bursts, and agents worked through the backlog smoothly.
A media firm had a content-moderation agent pipeline where one classification agent occasionally hung. With direct calls, a hung agent blocked the whole pipeline. With a queue, its unacked message requeued and another instance handled it while the hung one was recycled.
A fintech team needed to scale their analysis agents up and down with load. Because a queue decouples producers from consumers, they added and removed consumer instances freely - the queue didn't care how many workers pulled from it.
⚡ Pro tip: The requeue-on-crash behavior is worth more than the load balancing for most teams. Agents crash, hang, and time out constantly - far more than traditional services - because they depend on flaky model APIs. A queue turns "a worker died and took the task with it" into "the task got picked up by someone else." That resilience alone justifies the queue.
How Do You Use an Agent Message Queue Well?
Three practices separate a queue that helps from one that causes new problems.
Make consumers idempotent, because at-least-once delivery is the norm. A queue that guarantees no message is lost will sometimes deliver a message twice - if a consumer processes a task but crashes before acking, the task requeues and runs again. So processing a task twice must be safe. This is the same idempotency discipline that retries require, and it's non-negotiable with queues.
[object Object], ,[object Object],(,[object Object],):
msg_id, task = queue.pull()
,[object Object], task.dedup_key ,[object Object], done_set: ,[object Object],
queue.ack(msg_id) ,[object Object],
,[object Object],
result = agent.run(task)
done_set.add(task.dedup_key)
queue.ack(msg_id)
,[object Object], resultWhat this does: It guards against duplicate delivery by checking a dedup key before processing. If the task already ran, the consumer acks the duplicate without re-running it - so at-least-once delivery doesn't turn into duplicate side effects.
Set a visibility timeout so a pulled-but-unacked message returns to the queue after a bounded wait, not never and not immediately. Too short and slow tasks get processed twice; too long and crashed tasks sit stuck. Tune it to a bit above your p95 processing time.
Use a dead-letter queue for messages that fail repeatedly. A task that fails five times is probably poisoned - malformed input, an impossible request - and requeuing it forever just burns resources. Route it to a dead-letter queue for human inspection instead of letting it cycle endlessly.
⚡ Pro tip: Monitor queue depth as your primary health metric. A steadily growing queue means consumers can't keep up - you need more agents or faster ones. A queue that's always empty means you're possibly over-provisioned. Queue depth is the single clearest signal of whether your agent fleet is sized right, and it's free to watch.
Do Agent Queues Need Priority and Ordering?
Two questions come up the moment a queue is working: can urgent tasks jump ahead, and do tasks that must run in order stay in order? Both have answers that differ from traditional queue systems in ways specific to agents.
Priority is usually worth adding, because agent workloads are genuinely uneven - a critical customer escalation shouldn't wait behind a batch of routine classifications. A priority queue lets urgent work preempt routine work. The catch is starvation: if high-priority work never stops arriving, low-priority tasks wait forever. Add aging - a task's effective priority rises the longer it waits - so nothing starves indefinitely.
[object Object], ,[object Object],(,[object Object],):
wait = now - task.enqueued_at
,[object Object],
,[object Object], task.base_priority + (wait / ,[object Object],) ,[object Object],What this does: It boosts a task's priority the longer it sits in the queue, so a low-priority task that's been waiting ten minutes eventually competes with fresh high-priority work. This prevents the classic priority-queue failure where a flood of urgent tasks permanently buries everything routine.
Ordering is subtler. A plain queue with multiple consumers gives you no ordering guarantee - two related tasks can be processed simultaneously by different agents in any order. If ordering matters (process a customer's events in sequence), you need per-key ordering: route all messages sharing a key to the same consumer, so they're processed in order relative to each other while unrelated keys still parallelize. This is how you keep the throughput of many workers without losing the ordering some tasks require.
⚡ Pro tip: Only enforce ordering where it's genuinely required, keyed as narrowly as possible. Global ordering serializes your whole queue and destroys the parallelism you added agents to get. Per-customer or per-document ordering keeps most work parallel while protecting the few sequences that truly must stay in order. Over-constraining ordering is a common way teams accidentally turn a fast queue back into a slow pipeline.
Common Mistakes
⚠️ Common mistake: Reaching for a message queue when you have three agents running once per request. The queue's benefits - decoupling, resilience, backpressure - only matter under scale, concurrency, or unreliable workers. For a small synchronous pipeline, a queue adds operational weight (a broker to run, delivery semantics to reason about, idempotency to enforce) for benefits you won't use. Add it when you feel the pain of load imbalance or crash-lost work, not before.
The second mistake is forgetting idempotency and getting bitten by duplicate delivery in production - double-processed payments, duplicate notifications. If you adopt a queue, adopt idempotent consumers in the same commit.
The third is ignoring the dead-letter queue until poisoned messages are cycling thousands of times. Set up dead-lettering when you set up the queue, not after a poisoned message has run up a bill.
Conclusion
An agent message queue decouples work from workers, giving you load balancing, crash resilience, and backpressure that direct coordination can't. The price is real - at-least-once delivery forces idempotent consumers, and you take on a broker to operate - so reach for a queue when scale or unreliable agents make those benefits worth it, not by default.
Keep the consumer prompts and idempotency-key derivation logic versioned so you reuse a queue setup you trust. I store the consumer patterns and dedup-key conventions in PromptABCD, because getting at-least-once processing to behave correctly is fiddly, and once a queue-backed agent behaves well under duplicate delivery, you want to reproduce that exact behavior in the next system rather than rediscovering the edge cases.
⚡ Pro tip: When you introduce a queue, add a synthetic "canary" task that flows through it every minute and alerts if it doesn't complete in time. Queues fail quietly - a stuck consumer or a broker hiccup doesn't error, it just stops draining - so an end-to-end canary is often the first thing to tell you the pipeline stalled, well before queue depth alarms fire.
Continue Reading
Save the prompts from this post
PromptABCD is a free prompt manager. Paste, organize, and reuse your best AI prompts — no more hunting through chat history.
