RabbitMQ in Production: Acks, Prefetch, Durability, Quorum Queues and Dead-Lettering

Key takeaways

RabbitMQ gives you at-least-once delivery only if every link is configured: a durable queue, persistent messages, publisher confirms on the producer and manual acks with a bounded prefetch on the consumer. This guide explains each link, the exact errors you hit when one is wrong, and how to handle poison messages with dead-letter exchanges and quorum queue delivery limits.

RabbitMQ is easy to get working and surprisingly easy to get subtly wrong. A “hello world” producer and consumer run fine on a laptop, then in production messages vanish after a broker restart, get processed twice, or pile up as thousands of unacked deliveries on one consumer. Almost every one of those problems comes from one of five settings: acknowledgement mode, prefetch, queue durability, message persistence, and publisher confirms.

This guide walks through the AMQP model briefly, then spends most of its time on those reliability settings, the error messages RabbitMQ gives when they are wrong, and how to deal with messages that can never succeed. Examples use Node.js with amqplib, with a short Python pika section at the end.


The Model in One Paragraph

Producers never write to a queue directly. They publish to an exchange with a routing key; the exchange looks at its bindings and copies the message into zero or more queues; consumers subscribe to queues and acknowledge each message when they are done with it. Even channel.sendToQueue('tasks', ...) goes through an exchange: it publishes to the unnamed default exchange, which has an implicit binding for every queue by name.

That “zero or more queues” part matters. If no binding matches, the exchange silently drops the message. Nothing errors on the producer unless you ask for it (see publisher confirms and mandatory).


Local Setup

docker run -d --name rabbitmq \
  -p 5672:5672 -p 15672:15672 \
  rabbitmq:4-management

The management UI is at http://localhost:15672. The default guest user only works from localhost. If your app runs in another container or host and tries guest, the connection is refused with:

ACCESS_REFUSED - Login was refused using authentication mechanism PLAIN.

Create a real user (rabbitmqctl add_user, set_permissions) or set RABBITMQ_DEFAULT_USER / RABBITMQ_DEFAULT_PASS on the container instead of loosening the guest restriction.


A Producer That Does Not Leak Connections

npm install amqplib

A common first version opens a connection, publishes one message and closes it. That works, but a connection is a TCP connection plus an AMQP handshake, and opening one per message will eventually exhaust file descriptors or trip connection-churn alarms. Open one connection per process, keep it, and reuse a channel:

import amqp from 'amqplib';

const connection = await amqp.connect(process.env.AMQP_URL ?? 'amqp://localhost');
const channel = await connection.createConfirmChannel();

await channel.assertQueue('tasks', { durable: true });

export function publishTask(task) {
  return new Promise((resolve, reject) => {
    channel.sendToQueue(
      'tasks',
      Buffer.from(JSON.stringify(task)),
      { persistent: true, contentType: 'application/json' },
      (err) => (err ? reject(err) : resolve())
    );
  });
}

connection.on('error', (err) => console.error('AMQP connection error', err));
connection.on('close', () => {
  // amqplib does not reconnect for you. Exit and let the supervisor restart,
  // or implement reconnect-and-redeclare logic here.
  process.exit(1);
});

Two details worth knowing:

  • amqplib does not auto-reconnect. When the connection drops, the channel objects are dead. You either restart the process or rebuild the connection, channels, and topology yourself.
  • sendToQueue / publish return false when the write buffer is full. The message is still buffered, but if you keep publishing in a tight loop you grow memory without bound. Wait for the channel’s 'drain' event before publishing more in bulk jobs.

Durable Queue AND Persistent Message

These two flags are independent, and both are required for a message to survive a broker restart:

SettingWhereWhat it controls
durable: trueassertQueueThe queue definition survives a restart
persistent: trueeach publishThe message is written to disk

A durable queue holding transient messages comes back after a restart, but empty. Persistent messages published to a non-durable queue are lost along with the queue. Quorum queues are always durable and always store messages on disk, which removes one way to get this wrong.

Queue arguments are also fixed at declaration time. If a queue already exists as non-durable and your code now asserts it as durable, or you add a dead-letter argument to an existing queue, the channel is closed with:

PRECONDITION_FAILED - inequivalent arg 'durable' for queue 'tasks' in vhost '/': received 'true' but current is 'false'

Deleting and recreating the queue loses its messages. For settings that can change over time (dead-letter exchange, TTL, max length, delivery limit), prefer policies (rabbitmqctl set_policy or the management UI), which can be updated without redeclaring the queue.


Consumer Acks and Prefetch

How acknowledgement actually works

In amqplib, channel.consume uses manual acknowledgement by default (noAck: false). Each delivered message is unacked until you call ack, nack, or reject. If the channel closes while a message is unacked (process crash, deploy, network blip), RabbitMQ puts it back on the queue and redelivers it with msg.fields.redelivered === true.

That is the mechanism that gives you at-least-once delivery, and it has a consequence: your consumer will see duplicates. A message processed successfully but not yet acked when the process dies will be processed again. Make handlers idempotent, for example by storing a message or business ID with a unique constraint and skipping already-processed ones.

const channel = await connection.createChannel();
await channel.assertQueue('tasks', { durable: true });
await channel.prefetch(10);

await channel.consume('tasks', async (msg) => {
  if (msg === null) return; // consumer was cancelled by the broker (queue deleted, etc.)

  let task;
  try {
    task = JSON.parse(msg.content.toString());
  } catch {
    channel.nack(msg, false, false); // malformed: never retry, dead-letter it
    return;
  }

  try {
    await processTask(task);
    channel.ack(msg);
  } catch (err) {
    console.error('task failed', { redelivered: msg.fields.redelivered, err });
    channel.nack(msg, false, false); // to the DLX, see below
  }
});

Three failure modes show up repeatedly:

  1. A code path that never acks. An early return or an unhandled promise rejection leaves the message unacked forever. Since RabbitMQ 3.8.15 there is a consumer delivery acknowledgement timeout (30 minutes by default); when it fires, the broker closes the whole channel with PRECONDITION_FAILED - delivery acknowledgement on channel 1 timed out, and every unacked message on that channel is requeued.
  2. Acking twice or on the wrong channel. Delivery tags are scoped to the channel. Acking a tag twice, or acking on a different channel than the one that received the message, closes the channel with PRECONDITION_FAILED - unknown delivery tag 1.
  3. nack(msg, false, true) on a message that will always fail. Requeue puts it back at the head of the queue, the consumer receives it immediately, fails again, and you get a hot loop burning CPU and filling logs. That is the poison-message problem covered below.

Why prefetch matters

Prefetch (basic.qos) caps how many unacked messages a consumer may hold. With no prefetch set, the limit is unbounded: RabbitMQ pushes as much of the queue as it can into the first consumer to connect. You see this as one worker with thousands of unacked messages in the management UI while newly started workers get nothing, and when that worker crashes all of those messages are redelivered at once.

  • prefetch(1) gives the fairest distribution between workers of uneven speed but costs a network round trip per message.
  • Something in the range of 10 to a few hundred is typical for fast handlers; tune it by watching consumer utilisation in the management UI rather than guessing.
  • Prefetch should be at least as large as the concurrency you actually use. An async handler in amqplib does not block the next delivery, so with prefetch(10) you may have 10 handlers running at once. If the downstream work cannot handle that parallelism, lower prefetch.

I have seen the unbounded-prefetch trap more than any other RabbitMQ mistake: everything looks healthy in development because there is only one consumer and a few messages, and the problem only shows up when a backlog forms and autoscaling adds workers that sit idle while one pod holds the whole queue in memory. Setting prefetch explicitly on every consumer is the first thing I check when reviewing RabbitMQ code.


Publisher Confirms and Unroutable Messages

Acks protect the consumer side. On the producer side, channel.publish returns before the broker has done anything with the message. If the broker crashes or the connection drops right after, the message is gone and the producer never knows.

A confirm channel makes the broker acknowledge each publish. For a persistent message on a durable queue, the confirm is sent after the message has been written (for quorum queues, after a majority of replicas have it). The per-message callback shown earlier is one option; for batches, publish several messages and wait once:

const ch = await connection.createConfirmChannel();

for (const task of batch) {
  ch.sendToQueue('tasks', Buffer.from(JSON.stringify(task)), { persistent: true });
}
await ch.waitForConfirms(); // rejects if any message was nacked by the broker

A confirm does not mean the message reached a queue. If the routing key matches no binding, the broker drops the message and still confirms it. To detect that, publish with mandatory: true and listen for returned messages:

ch.on('return', (msg) => {
  console.error('unroutable message', msg.fields.exchange, msg.fields.routingKey);
});

ch.publish('orders', 'created', Buffer.from(body), { persistent: true, mandatory: true });

The typical bug this catches is a typo in a routing key or a consumer service that has not yet declared its queue and binding when the producer starts publishing.


Exchange Types, With Working Examples

Direct: exact routing key match

await channel.assertExchange('orders', 'direct', { durable: true });
await channel.assertQueue('order.created', { durable: true });
await channel.bindQueue('order.created', 'orders', 'created');

channel.publish('orders', 'created', Buffer.from(JSON.stringify(order)), { persistent: true });

Topic: pattern routing

Routing keys are dot-separated words. In bindings, * matches exactly one word and # matches zero or more.

await channel.assertExchange('logs', 'topic', { durable: true });
await channel.assertQueue('error-logs', { durable: true });
await channel.assertQueue('all-logs', { durable: true });
await channel.bindQueue('error-logs', 'logs', '*.error'); // auth.error, payment.error
await channel.bindQueue('all-logs', 'logs', '#');

channel.publish('logs', 'auth.error', Buffer.from('Login failed'));
channel.publish('logs', 'payment.warn', Buffer.from('Retry 2')); // only all-logs

Note that *.error does not match payment.gateway.error (three words); use #.error if keys can have variable depth.

Fanout: broadcast to every bound queue

Each subscriber gets its own queue. For “every running instance should receive this” (cache invalidation, config reload), an exclusive, server-named queue per instance is the usual shape:

await channel.assertExchange('notifications', 'fanout', { durable: true });

const q = await channel.assertQueue('', { exclusive: true }); // server-named, deleted on disconnect
await channel.bindQueue(q.queue, 'notifications', '');

await channel.consume(q.queue, (msg) => {
  if (msg) handleNotification(msg);
}, { noAck: true });

noAck: true is acceptable here because the queue is transient anyway; a broadcast that is missed while the instance is down is not replayed. If subscribers must not miss messages while offline, give them named durable queues instead.

Headers: match on message headers

await channel.assertExchange('reports', 'headers', { durable: true });
await channel.assertQueue('pdf-queue', { durable: true });
await channel.bindQueue('pdf-queue', 'reports', '', {
  'x-match': 'all', // 'any' matches if at least one header matches
  format: 'pdf',
  priority: 'high',
});

channel.publish('reports', '', Buffer.from(data), {
  headers: { format: 'pdf', priority: 'high' },
});

Headers exchanges are rarely needed; a topic exchange with a structured routing key usually covers the same cases and is easier to reason about.


Quorum Queues vs Classic Queues

Classic queues live on one node. In a cluster, if that node goes down, the queue is unavailable. Classic mirrored queues used to address that, but they were deprecated for years and removed in RabbitMQ 4.0.

Quorum queues replicate each queue across several nodes with the Raft consensus algorithm. A publish is confirmed once a majority of replicas have it, and the queue stays available as long as a majority of its members are up. That is why a quorum queue cluster needs an odd number of nodes, usually three or five: with three nodes you tolerate one failure.

await channel.assertQueue('payments', {
  durable: true,
  arguments: { 'x-queue-type': 'quorum' },
});

Trade-offs to know before switching everything:

  • Quorum queues must be durable and cannot be exclusive, so they are not a fit for per-connection temporary queues.
  • Every message is replicated and written to disk, so throughput per queue is lower and disk usage higher than a classic queue with transient messages.
  • They track redelivery counts, which enables delivery limits: once a message has been returned to the queue more times than the limit, it is dead-lettered or dropped instead of being redelivered forever. RabbitMQ 4.0 made a default delivery limit of 20 apply to quorum queues; set your own with the x-delivery-limit argument or a policy.

The x-queue-type argument cannot be changed on an existing queue (you get the inequivalent arg error above). Migrating means creating a new quorum queue, pointing consumers and bindings at it, and draining the old one.


Dead-Letter Exchanges and Poison Messages

A message is dead-lettered when it is rejected or nacked with requeue=false, when its TTL expires, when the queue length limit pushes it out, or when a quorum queue delivery limit is exceeded. Instead of being discarded, it is republished to the queue’s dead-letter exchange (DLX), optionally with a different routing key. RabbitMQ adds an x-death header recording the reason and the queue it came from.

await channel.assertExchange('dlx', 'direct', { durable: true });
await channel.assertQueue('tasks.failed', { durable: true });
await channel.bindQueue('tasks.failed', 'dlx', 'tasks');

await channel.assertQueue('tasks', {
  durable: true,
  arguments: {
    'x-queue-type': 'quorum',
    'x-dead-letter-exchange': 'dlx',
    'x-dead-letter-routing-key': 'tasks',
    'x-delivery-limit': 5,
  },
});

With this setup there are two ways a message ends up in tasks.failed: the consumer nacks it without requeue (a known-permanent failure such as invalid JSON), or it is requeued five times and the delivery limit dead-letters it (a failure that kept recurring, such as a bug triggered by that specific payload).

Retrying with a delay

Immediate requeue retries far too fast for transient failures like a downstream API returning 503. A common pattern without plugins is a “wait” queue that has a TTL and dead-letters back to the main queue:

await channel.assertExchange('retry', 'direct', { durable: true });
await channel.assertQueue('tasks.retry-30s', {
  durable: true,
  arguments: {
    'x-message-ttl': 30_000,
    'x-dead-letter-exchange': '',          // default exchange
    'x-dead-letter-routing-key': 'tasks',  // back to the main queue by name
  },
});
await channel.bindQueue('tasks.retry-30s', 'retry', 'tasks');

// In the consumer, on a transient error:
const retryCount = msg.properties.headers?.['x-retry-count'] ?? 0;
channel.publish('retry', 'tasks', msg.content, {
  persistent: true,
  headers: { ...msg.properties.headers, 'x-retry-count': retryCount + 1 },
});
channel.ack(msg);

Nobody consumes the retry queue; messages sit there until the TTL expires and are then routed back. Carry your own retry counter in a header and send the message to the failure queue once it exceeds your limit. One caveat: per-queue TTL works because all messages in that queue have the same TTL. Per-message TTLs with different values in one queue behave badly, because expired messages are only removed from the head of the queue.

A dead-letter queue is only useful if someone looks at it. Alert on its depth, keep enough context in the message (original routing key, x-death, error text in a header) to understand the failure, and have a way to re-publish messages back to the main queue after a fix.

The poison-message loop is the other failure I would warn anyone about. A consumer that does nack(msg, false, true) on every exception works fine until one payload triggers a deterministic bug; then that single message is redelivered as fast as the consumer can fail on it, the logs fill with the same stack trace, and the messages behind it wait. Distinguishing permanent failures (dead-letter immediately) from transient ones (delayed retry) and relying on a delivery limit as the backstop is what stops that pattern.


Python with pika

pip install pika
import json
import pika

connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
channel = connection.channel()
channel.queue_declare(queue='tasks', durable=True)

channel.confirm_delivery()  # basic_publish now raises if the broker nacks
channel.basic_publish(
    exchange='',
    routing_key='tasks',
    body=json.dumps({'type': 'send-email', 'to': '[email protected]'}),
    properties=pika.BasicProperties(delivery_mode=2),  # persistent
)

def callback(ch, method, properties, body):
    task = json.loads(body)
    handle(task)
    ch.basic_ack(delivery_tag=method.delivery_tag)

channel.basic_qos(prefetch_count=10)
channel.basic_consume(queue='tasks', on_message_callback=callback)
channel.start_consuming()

The pika-specific pitfall is heartbeats. BlockingConnection processes heartbeats on the same thread that runs your callback, so a callback that blocks for longer than the negotiated heartbeat timeout lets the broker close the connection; pika then raises StreamLostError and the unacked message is redelivered to another worker, which may also time out. For long jobs, run the work in a separate thread and ack via connection.add_callback_threadsafe, rather than disabling heartbeats. Also note that a pika connection and its channels must not be shared between threads.


Reliability Settings at a Glance

Link in the chainSettingFailure if missing
Queuedurable: true or a quorum queueQueue disappears on restart
Messagepersistent: true / delivery_mode=2Queue returns empty after restart
ProducerConfirm channel, mandatory + return handlerSilent loss on broker failure or bad routing key
ConsumerManual ack on every pathRedelivery loops, ack timeout closing channels
ConsumerExplicit prefetchOne consumer hoards the queue
QueueDLX + delivery limitPoison messages loop forever
HandlerIdempotencyDuplicate side effects on redelivery

When RabbitMQ Is Not the Right Tool

RabbitMQ deletes a message once it is acknowledged, so there is no replay: if a new service needs last week’s events, a queue cannot give them to it. If you need retained, replayable streams, look at Kafka (or RabbitMQ Streams, a separate log-based queue type). If you only need background jobs inside a Node.js application and already run Redis, a Redis-backed job queue can be simpler to operate. RabbitMQ fits best when you want flexible routing between several services, per-message acknowledgement, and a broker that pushes work to consumers.