Single Active Consumer in RabbitMQ
RabbitMQ's single active consumer (SAC) feature is the broker-native way to build an ordered lane with hot standby, and this guide shows how to use it within the ordering model laid out in Message Ordering Guarantees, part of Queue Fundamentals & Architecture. It covers enabling SAC, sharding so ordered work still scales, and the prefetch and requeue settings that silently undo ordering if left at their defaults.
Problem Statement
A document service publishes edit operations to a RabbitMQ queue, and a worker applies them to a search index. With one consumer the index is correct but a deploy or crash stops indexing until the pod restarts. With three consumers failover is instant but edits for the same document interleave, and the index occasionally stores an older revision. You want exactly one consumer processing each ordered stream at any time, a standby that takes over within seconds, and enough parallelism to keep up with 2,000 edits per second.
Prerequisites
- RabbitMQ 3.8 or later (SAC is available on classic and quorum queues; 3.13+ recommended for quorum queue SAC improvements).
- The
rabbitmq_consistent_hash_exchangeplugin enabled if you shard (it ships with RabbitMQ). - A client that reports consumer cancellation cleanly —
pika,aio-pika,amqplib, or the Java client. - A per-document key on every message, carried as the routing key or a header.
Step 1 — Declare a Queue with Single Active Consumer
SAC is a queue argument set at declaration time. Every consumer subscribes normally; the broker picks one as active and holds the rest in a waiting list.
# declare.py — pika
import pika
conn = pika.BlockingConnection(pika.ConnectionParameters("rabbitmq.internal"))
ch = conn.channel()
ch.queue_declare(
queue="doc-edits",
durable=True,
arguments={
"x-single-active-consumer": True, # only one consumer receives at a time
"x-queue-type": "quorum", # replicated across 3 nodes for durability
"x-delivery-limit": 10, # quorum queues: dead-letter after 10 redeliveries
"x-dead-letter-exchange": "doc-edits.dlx",
},
)
Queue arguments cannot be changed on an existing queue. If doc-edits already exists without SAC, declare a new queue name, rebind, and drain the old one. A mismatch between your declaration and the existing queue raises PRECONDITION_FAILED and closes the channel, which is a common surprise in rolling deploys where old and new code declare differently.
Step 2 — Shard So Ordered Work Scales
One SAC queue is one lane, and one lane is one consumer's throughput. To handle 2,000 edits per second with a 5 ms handler you need at least ten lanes. A consistent-hash exchange routes each document to one of N shard queues, each with SAC.
# shards.py — 16 ordered lanes, each with its own active consumer
SHARDS = 16
ch.exchange_declare("doc-edits", exchange_type="x-consistent-hash", durable=True,
arguments={"hash-header": "doc-id"}) # hash a header, not the routing key
for i in range(SHARDS):
q = f"doc-edits.{i:02d}"
ch.queue_declare(q, durable=True, arguments={
"x-single-active-consumer": True,
"x-queue-type": "quorum",
"x-delivery-limit": 10,
"x-dead-letter-exchange": "doc-edits.dlx",
})
ch.queue_bind(q, "doc-edits", routing_key="10") # binding weight: equal share per shard
# publishing
ch.basic_publish(
exchange="doc-edits",
routing_key="",
body=payload,
properties=pika.BasicProperties(headers={"doc-id": doc_id}, delivery_mode=2),
)
Every consumer process subscribes to all sixteen shard queues. The broker makes each process active on some shards and waiting on others, so load spreads across the fleet and any shard fails over when its active process dies. With four processes, each is active on roughly four shards.
Step 3 — Keep Prefetch and Handling Sequential
SAC guarantees one consumer receives a queue's messages; it does not guarantee that consumer processes them one at a time. With prefetch_count=50 and an async handler that spawns a task per delivery, the active consumer reorders messages itself.
# consumer.py — aio-pika, sequential per shard, prefetch for throughput
import aio_pika
async def main():
conn = await aio_pika.connect_robust("amqp://rabbitmq.internal/")
ch = await conn.channel()
await ch.set_qos(prefetch_count=20) # buffered deliveries, handled one by one
for i in range(SHARDS):
q = await ch.get_queue(f"doc-edits.{i:02d}")
await q.consume(make_sequential_handler(i))
def make_sequential_handler(shard: int):
lock = asyncio.Lock() # one in-flight message per shard
async def on_message(msg: aio_pika.IncomingMessage):
async with lock:
try:
await apply_edit(json.loads(msg.body))
await msg.ack()
except TransientError:
await asyncio.sleep(1)
await msg.nack(requeue=True) # see Step 4 about requeue position
except Exception:
await msg.reject(requeue=False) # dead-letter poison edits
return on_message
A prefetch of 10–50 keeps the pipe full so the handler never waits on the network, while the per-shard lock keeps processing sequential. The same idea — prefetch as a buffer, not as concurrency — is discussed in tuning prefetch and consumer concurrency.
Step 4 — Understand Where Requeued Messages Go
When the active consumer nacks with requeue=True, RabbitMQ puts the message back at (or near) the head of the queue, so it is the next delivery — ordering is preserved, but a poison message loops. Quorum queues count these redeliveries and dead-letter after x-delivery-limit. Classic queues have no limit, so a poison message spins forever and blocks its shard.
# Inspect delivery counts and the active consumer per shard
rabbitmqctl list_queues name messages_ready messages_unacknowledged \
single_active_consumer_tag consumers --formatter pretty_table
# Quorum queue redelivery count is visible on the message as x-delivery-count
rabbitmqadmin get queue=doc-edits.03 count=1 ackmode=ack_requeue_true
Prefer quorum queues for SAC lanes: the delivery limit converts an infinite loop into a dead-lettered message and an alert, and the lane moves on. If moving on is unacceptable for your data (a later edit depends on the failed one), detect the gap downstream as described in handling out-of-order events with sequence numbers.
Step 5 — Make Failover Fast
Failover happens when the broker notices the active consumer's channel is gone. A crashed process closes its TCP connection and the switch is immediate; a hung process or a network partition relies on heartbeats.
params = pika.ConnectionParameters(
host="rabbitmq.internal",
heartbeat=15, # broker declares the connection dead after ~2 missed beats
blocked_connection_timeout=60,
)
On shutdown, cancel consumers and let in-flight work finish before closing — the pattern in graceful shutdown and deployments. A consumer that is killed mid-message has its unacked delivery returned to the queue head, and the new active consumer processes it first, so order holds even in a crash.
Step 6 — Alert on Stalled or Orphaned Lanes
A SAC shard has two failure shapes that ordinary queue alerts miss. The stalled lane has an active consumer that is alive but stuck — a handler blocked on a lock or a slow dependency — so messages_ready grows on one shard while the others are empty. The orphaned lane has no consumers at all, typically after a deploy that changed the shard count, so messages accumulate with nobody subscribed. Both look like "some backlog" in an aggregate graph.
# Stalled: one shard's ready count keeps growing while its consumer is connected
rabbitmq_queue_messages_ready{queue=~"doc-edits\\.\\d+"} > 1000
and on (queue) rabbitmq_queue_consumers{queue=~"doc-edits\\.\\d+"} > 0
and on (queue) deriv(rabbitmq_queue_messages_ready{queue=~"doc-edits\\.\\d+"}[10m]) > 0
# Orphaned: a shard with no consumers at all
rabbitmq_queue_consumers{queue=~"doc-edits\\.\\d+"} == 0
# Unacked stuck at the prefetch ceiling with no acks: the handler is wedged
rabbitmq_queue_messages_unacked{queue=~"doc-edits\\.\\d+"} >= 20
and on (queue) rate(rabbitmq_global_messages_acknowledged_total[5m]) == 0
The per-queue metrics come from the rabbitmq_prometheus plugin with per-object metrics enabled (prometheus.return_per_object_metrics = true), which is off by default on large brokers because of cardinality; with sixteen shards it is cheap. Route the stalled-lane alert to the owning team with the shard name in the summary — the fix is almost always in the handler, not the broker. The alerting approach mirrors alerting on queue backlog with Prometheus, applied per lane.
Verification
Check that each shard has exactly one active consumer and that failover works:
# Expect one non-empty single_active_consumer_tag per shard
rabbitmqctl list_queues name single_active_consumer_tag consumers | grep doc-edits
# Kill the active process for shard 03 and watch the tag change within heartbeat time
kubectl delete pod indexer-7c9f-abcde --grace-period=0 --force
watch -n1 "rabbitmqctl list_queues name single_active_consumer_tag | grep doc-edits.03"
Then run an ordering test: publish 1,000 numbered edits for 50 documents, kill a consumer halfway, and assert each document's applied revisions are strictly increasing.
Gotchas & Edge Cases
Uneven shard load. Consistent hashing is only as even as your key distribution. A few very active documents make their shards hot while others idle; compare messages_ready across shards and consider more shards than consumers.
Waiting consumers still count. Each waiting consumer holds channel resources. A fleet of 40 processes subscribed to 64 shards means 2,560 subscriptions; that is fine for RabbitMQ, but channel limits (channel_max) must allow it.
Priorities do not apply. Consumer priorities interact with SAC — the highest-priority consumer becomes active — which can be used to prefer a local-zone consumer, but message priorities on the queue would reorder the lane and should not be combined with ordering requirements.
Streams are different. RabbitMQ Streams also support single active consumer with named consumers and offset tracking; semantics for offsets differ from queues, so do not mix the two patterns in one lane.
FAQ
Is single active consumer the same as exclusive consumers? No. An exclusive consumer prevents anyone else from subscribing at all, so there is no standby and a reconnect after a crash can race. SAC lets many consumers subscribe and promotes one, which gives automatic failover.
Does SAC work across a RabbitMQ cluster? Yes. The queue's leader node tracks which consumer is active regardless of which node each consumer connects to. For quorum queues, a leader failover keeps the SAC state.
How many shards should I create? Enough that one shard's throughput times the shard count exceeds peak load with margin, and at least two to four times the number of consumer processes so load balances when processes come and go.
Related
- Message Ordering Guarantees in Job Queues — lanes, failure policies, and ordering trade-offs.
- Per-Entity Ordering with Kafka Partition Keys — the log-based equivalent of sharded SAC queues.
- Replaying Dead-Letter Messages in RabbitMQ — recovering edits that hit the delivery limit.
- RabbitMQ Consumer Timeout for Unacked Messages — what happens when an active consumer holds a message too long.