RabbitMQ Lab¶
A hands-on tour of RabbitMQ 4.3, from a first message to surviving node failure. Every exercise
runs from the rabbitmq-lab project: Docker Compose files for a single node and a three-node
cluster, a Python script per exercise, Prometheus/Grafana monitoring and a few operational
scripts. Concepts are explained in the Overview; this page is about seeing them
happen. Code excerpts below come straight from the project.
Before you start¶
You need Docker with Compose v2, Python 3.10 or newer, make, openssl, and jq for the
health-check script. Most exercises need two or three terminals, labelled T1, T2, T3
below; every Python command runs from py/ with the virtualenv active.
rabbitmq-lab/
├── compose.yaml single node (+ --profile monitoring)
├── compose.cluster.yaml 3 nodes behind HAProxy (+ --profile monitoring)
├── config/ rabbitmq.conf, cluster peer discovery, haproxy, prometheus, grafana
├── py/ one script per exercise; every script has --help
└── scripts/ healthcheck.sh, triage.sh, perftest.sh
All exposed ports bind to 127.0.0.1. Credentials and the Erlang cookie are generated into a
git-ignored .env.
Part 1: Bring up a broker¶
1.1 Start a single node¶
make env # writes .env with random passwords and an Erlang cookie
make up # docker compose up -d --wait: returns once the health check passes
make venv # .venv with pika
make creds # the generated login
source .venv/bin/activate && cd py
The broker runs the official rabbitmq:4.3-management image. Rather than replacing the image's
configuration, the lab layers its own file into /etc/rabbitmq/conf.d/ as 20-lab.conf, after
the image's 10-defaults.conf. Files in conf.d load alphabetically and later files win, which
is a clean way to keep site settings separate from packaged defaults.
# config/rabbitmq.conf
## Lab settings, layered on top of the image's /etc/rabbitmq/conf.d/10-defaults.conf.
## Mounted as /etc/rabbitmq/conf.d/20-lab.conf.
## --- Resource alarms -------------------------------------------------------
## Publishers are blocked when either alarm fires. RabbitMQ reads the container's
## cgroup memory limit, so the relative watermark applies to mem_limit in compose.
vm_memory_high_watermark.relative = 0.6
## The built-in default (50MB) is far too low for anything real.
disk_free_limit.absolute = 1GB
## --- Logging ----------------------------------------------------------------
## Containers log to stdout; let the container runtime handle collection.
log.console = true
log.console.level = info
log.file = false
## --- Connections ------------------------------------------------------------
## Server-suggested heartbeat (seconds). Clients negotiate the lower of the two.
heartbeat = 60
## --- Consumers --------------------------------------------------------------
## Max time a delivery may stay unacknowledged (ms). As of 4.3 this is enforced
## by quorum queues only; classic queues and streams ignore it.
consumer_timeout = 1800000
## --- Prometheus ---------------------------------------------------------------
## Keep /metrics aggregated (cheap). Per-queue series come from /metrics/detailed,
## which Prometheus scrapes with an explicit family filter.
prometheus.return_per_object_metrics = false
The container is limited to 1 GB of memory. RabbitMQ reads the cgroup limit, so the 0.6
watermark means the memory alarm fires at roughly 600 MB. The health check runs
rabbitmq-diagnostics check_running and check_local_alarms, cheap enough to run every 15
seconds without loading the node.
1.2 First look¶
Open http://localhost:15672 and log in with the credentials from make creds. The
Overview tab shows message rates, queue totals and the node's resource use (memory against
its watermark, disk against its limit, file descriptors, Erlang processes). Keep a tab open on
it through the lab.
Then from the CLI, which talks to the node over Erlang distribution using the shared cookie:
docker compose exec rabbitmq rabbitmq-diagnostics status # versions, config files loaded, resources
docker compose exec rabbitmq rabbitmq-diagnostics environment | grep -E 'vm_memory|disk_free'
docker compose exec rabbitmq rabbitmq-plugins list -e # enabled plugins
docker compose exec rabbitmq rabbitmqctl list_feature_flags name state
In status, find the Config files section and confirm both 10-defaults.conf and
20-lab.conf are listed. When a setting doesn't seem to apply, that list is the first place to
look. Installing on a VM instead of in containers is covered in
Operations.
Part 2: Messaging fundamentals¶
2.1 Routing, confirms and acknowledgements¶
basics.py publishes JSON events to a topic exchange, lab.logs, bound to a quorum queue
lab.app-logs with the pattern app.#. Producer and consumer declare the topology through one
shared function, because redeclaring a queue with different arguments is an error:
# py/basics.py
def declare_topology(ch) -> None:
"""Declared identically by producer and consumer. Declarations are idempotent
only if the arguments match: redeclaring a queue with different arguments
fails with 406 PRECONDITION_FAILED and closes the channel."""
ch.exchange_declare(exchange=EXCHANGE, exchange_type="topic", durable=True)
ch.queue_declare(
queue=QUEUE,
durable=True,
arguments={
"x-queue-type": "quorum",
"x-message-ttl": 24 * 60 * 60 * 1000, # drop messages older than 24h
"x-max-length": 10_000,
"x-overflow": "reject-publish", # when full, nack publishers instead of dropping old messages
},
)
# "app.#" matches app.info, app.web.error, ... but not db.error.
ch.queue_bind(queue=QUEUE, exchange=EXCHANGE, routing_key="app.#")
The producer turns on publisher confirms and publishes with mandatory=True. With pika's
BlockingConnection, each basic_publish then waits for the broker's verdict and raises if the
message was returned (unroutable) or nacked:
# py/basics.py
def produce(_args) -> None:
conn = connect("basics-producer")
ch = conn.channel()
declare_topology(ch)
ch.confirm_delivery() # every basic_publish now waits for the broker's ack
events = [
("app.info", "application started"),
("app.warning", "memory usage high"),
("app.error", "database connection failed"),
("app.web.info", "request served"),
("db.error", "this key matches no binding"), # unroutable on purpose
]
for routing_key, message in events:
body = to_json(
{
"ts": datetime.now(timezone.utc).isoformat(),
"key": routing_key,
"message": message,
}
)
try:
# mandatory=True: if no queue is bound for this key, the broker returns
# the message instead of silently dropping it.
ch.basic_publish(EXCHANGE, routing_key, body, json_props(), mandatory=True)
print(f" [x] confirmed {routing_key:13} {message}")
except pika.exceptions.UnroutableError:
print(f" [!] unroutable {routing_key:13} (returned by broker)")
except pika.exceptions.NackError:
print(f" [!] nacked {routing_key:13} (queue full or broker refused)")
conn.close()
python basics.py consume # T1
python basics.py produce # T2
Four events are confirmed and the consumer prints them. db.error comes back as unroutable:
# matches zero or more words, but only after app., so nothing is bound for db.*. Without
mandatory, the broker would have accepted and silently dropped it.
The consumer acknowledges manually, after its (simulated) work, with a prefetch of 10. Look at how it treats a message it can't parse:
# py/basics.py
def consume(args) -> None:
conn = connect("basics-consumer")
ch = conn.channel()
declare_topology(ch)
ch.basic_qos(prefetch_count=args.prefetch)
def on_message(ch, method, props, body):
try:
event = json.loads(body)
except json.JSONDecodeError:
# A message that can never be parsed must not be requeued, or it loops
# forever. With a DLX configured it would be dead-lettered; here it's dropped.
print(f" [!] unparseable message, rejecting: {body[:60]!r}")
ch.basic_reject(method.delivery_tag, requeue=False)
return
print(f" [<] {method.routing_key:13} redelivered={method.redelivered!s:5} {event['message']}")
time.sleep(0.2) # simulated work
ch.basic_ack(method.delivery_tag)
ch.basic_consume(QUEUE, on_message, auto_ack=False)
print(f" [*] consuming {QUEUE} (prefetch={args.prefetch}); Ctrl-C to exit")
run_consumer(ch)
Rejecting with requeue=False matters. Requeuing a message that can never succeed creates an
infinite redelivery loop that burns CPU and blocks the queue head. Lab 3.1 shows the proper
place for such messages.
Try this:
- Stop the consumer (Ctrl-C), run
producethree times and watch Queues → lab.app-logs hold 12 ready messages. Rundocker compose restart rabbitmq, then start the consumer: nothing was lost, because the queue is durable and the messages persistent. - Publish garbage from the UI: Exchanges → lab.logs → Publish message, routing key
app.junk, payloadnot json. The consumer logs it and rejects it. - Start the consumer, and in Queues → lab.app-logs look at the queue's type, its policy and arguments, and the consumer's prefetch under Consumers.
2.2 Work queues and fair dispatch¶
work_queue.py sends tasks that take between 1 and 5 seconds. Several workers consume the same
queue and the broker distributes messages between them.
# py/work_queue.py
def worker(args) -> None:
conn = connect(f"worker-{args.id}")
ch = conn.channel()
declare(ch)
ch.basic_qos(prefetch_count=args.prefetch)
def on_message(ch, method, _props, body):
name, seconds = body.decode().split(":")
print(f" [worker {args.id}] {name} ({seconds}s) ...", end="", flush=True)
time.sleep(int(seconds))
ch.basic_ack(method.delivery_tag)
print(" done")
ch.basic_consume(QUEUE, on_message)
print(f" [worker {args.id}] waiting (prefetch={args.prefetch})")
run_consumer(ch)
python work_queue.py worker 1 # T1
python work_queue.py worker 2 # T2
python work_queue.py send # T3
With prefetch_count=1 a worker gets a new task only after acknowledging its current one, so
the busy worker never sits on a queue of waiting tasks while the other idles. Restart both
workers with --prefetch 50 and send again: the first worker to connect receives most of the
batch up front, and the two finish at noticeably different times. A high prefetch is right for
fast, uniform work (it hides network latency); a low one is right for slow or uneven work.
Try this: while both workers are busy, kill one with Ctrl-C. The task it held is
redelivered to the other worker with redelivered=True, because unacknowledged deliveries
return to the queue when their channel closes.
2.3 Publish/subscribe¶
pubsub.py broadcasts through a fanout exchange. Each subscriber declares its own exclusive,
server-named queue, which exists only while the subscriber is connected:
# py/pubsub.py
def subscribe(args) -> None:
conn = connect(f"pubsub-{args.name}")
ch = conn.channel()
ch.exchange_declare(exchange=EXCHANGE, exchange_type="fanout", durable=True)
# queue="" lets the broker generate a name; exclusive=True ties the queue's
# lifetime to this connection. Exclusive transient queues remain fully supported
# (it's *non-exclusive* transient queues that 4.3 rejects by default).
result = ch.queue_declare(queue="", exclusive=True)
queue = result.method.queue
ch.queue_bind(queue=queue, exchange=EXCHANGE)
def on_message(_ch, _method, _props, body):
print(f" [{args.name}] {body.decode()}")
# auto_ack is acceptable here: notifications are fire-and-forget.
ch.basic_consume(queue, on_message, auto_ack=True)
print(f" [{args.name}] listening on {queue}")
run_consumer(ch)
python pubsub.py subscribe email # T1
python pubsub.py subscribe sms # T2
python pubsub.py publish "maintenance in 10 minutes" # T3
Both subscribers receive the message. In the UI, the subscribers' queues show names like
amq.gen-... with the Exclusive feature. Stop a subscriber and its queue disappears; messages
published while it's away are not kept for it. For durable subscriptions, each subscriber would
declare a named durable queue instead.
RabbitMQ 4.3 refuses to declare non-durable, non-exclusive queues by default. Exclusive ones
like these are fine; a shared temporary queue should be durable with a queue TTL (x-expires)
instead.
2.4 Long-running work¶
A consumer that takes minutes per message needs two things right. First, pika's
BlockingConnection only services heartbeats while its I/O loop runs: a callback that blocks
for longer than about two heartbeat intervals gets the connection closed by the broker, and the
message is redelivered. Second, pika channels aren't thread-safe, so work done in another thread
must hand the acknowledgement back to the connection's thread.
# py/slow_worker.py
def work(args) -> None:
conn = connect("slow-worker", heartbeat=HEARTBEAT)
ch = conn.channel()
declare(ch)
ch.basic_qos(prefetch_count=1)
def ack(tag: int) -> None:
if ch.is_open:
ch.basic_ack(tag)
print(f" [✓] acked delivery {tag}")
def on_message(ch, method, props, body):
seconds = int(body)
headers = props.headers or {}
print(f" [>] job {method.delivery_tag}: {seconds}s (delivery-count={headers.get('x-delivery-count', 0)})")
if args.blocking:
time.sleep(seconds) # starves heartbeats: expect the connection to drop
ch.basic_ack(method.delivery_tag)
return
def job() -> None:
time.sleep(seconds)
conn.add_callback_threadsafe(functools.partial(ack, method.delivery_tag))
threading.Thread(target=job, daemon=True).start()
consumer_args = {}
if args.consumer_timeout:
consumer_args["x-consumer-timeout"] = args.consumer_timeout
ch.add_on_cancel_callback(
lambda _frame: print(" [!] broker cancelled this consumer (consumer timeout); the job was requeued")
)
ch.basic_consume(QUEUE, on_message, arguments=consumer_args or None)
mode = "blocking (broken)" if args.blocking else "threaded"
print(f" [*] worker running, mode={mode}, heartbeat={HEARTBEAT}s")
try:
run_consumer(ch)
except pika.exceptions.AMQPConnectionError as exc:
print(f" [!] connection lost: {exc!r}")
The script uses a 5-second heartbeat so the failure appears quickly:
python slow_worker.py send --seconds 20 # T2
python slow_worker.py work --blocking # T1: the wrong way
The blocking worker sleeps in the callback, misses heartbeats, and loses its connection (the
broker log shows missed heartbeats from client). The job is redelivered with a higher
x-delivery-count. Now the right way:
python slow_worker.py work # T1: threaded
The connection stays healthy for the full 20 seconds and the ack lands normally.
Consumer timeouts. Independently of heartbeats, the broker limits how long a delivery can stay unacknowledged (30 minutes by default). On 4.3 this is enforced by quorum queues, and a timeout cancels just the offending consumer rather than closing the whole channel:
python slow_worker.py send --seconds 20
python slow_worker.py work --consumer-timeout 10000
Once the 10 seconds have passed (timeouts are checked periodically, so allow up to a minute)
the worker reports that the broker cancelled its consumer, and the job goes back to the queue. Long-running workloads should raise the timeout (consumer argument,
queue argument, consumer-timeout policy key, or consumer_timeout in rabbitmq.conf) rather
than disable it.
Part 3: Reliability patterns¶
3.1 Poison messages, retries and dead-lettering¶
retry.py publishes jobs where every third one is "poison": the consumer always fails it. The
work queue is a quorum queue with a delivery limit of 3 and at-least-once dead-lettering into
lab.dlx, which routes to lab.dead:
# py/retry.py
def declare(ch) -> None:
ch.exchange_declare(exchange=DLX, exchange_type="topic", durable=True)
ch.queue_declare(
queue=DLQ,
durable=True,
arguments={
"x-queue-type": "quorum",
# Quorum queues default to a delivery limit of 20; a DLQ you peek at
# repeatedly should never dead-letter (or drop) its own contents.
"x-delivery-limit": -1,
},
)
ch.queue_bind(queue=DLQ, exchange=DLX, routing_key="#")
base = {
"x-queue-type": "quorum",
"x-delivery-limit": DELIVERY_LIMIT,
"x-dead-letter-exchange": DLX,
# at-least-once requires reject-publish overflow; otherwise RabbitMQ
# silently falls back to at-most-once dead-lettering.
"x-dead-letter-strategy": "at-least-once",
"x-overflow": "reject-publish",
}
ch.queue_declare(queue=QUEUE, durable=True, arguments=base)
ch.queue_declare(
queue=DELAYED_QUEUE,
durable=True,
arguments={
**base,
# RabbitMQ 4.3+: hold failed messages aside before redelivery.
# delay = min(min * delivery_count, max) -> 2s, 4s, 6s ...
"x-delayed-retry-type": "failed",
"x-delayed-retry-min": 2000,
"x-delayed-retry-max": 8000,
},
)
python retry.py produce
python retry.py consume # Ctrl-C once the output settles
python retry.py inspect
Healthy jobs are acked. Each poison job is rejected with requeue=True and redelivered with a
rising x-delivery-count header; once the count exceeds the limit of 3 it leaves the work queue.
inspect finds it in lab.dead with x-first-death-reason: delivery_limit and an x-death entry naming
the queue.
Now repeat with nack instead of reject:
python retry.py produce
python retry.py consume --nack
On 4.3, basic.nack means "didn't process this" rather than "processing failed", so it
advances the new acquired-count but not delivery-count, and the delivery limit never
triggers. The script gives up after 10 local attempts and rejects without requeue; a real
consumer would loop forever. (On 4.2 and earlier, both methods counted.) Use reject for
failures.
Delayed retry (4.3). The retries above happen back to back, which is useless if the failure
is a downstream service that needs a moment. lab.retry.delayed adds delayed retry: a failed
message is set aside inside the queue for min(2s × delivery-count, 8s) before redelivery.
python retry.py produce --delayed
python retry.py consume --delayed
The poison jobs now come back after roughly 2, 4 and 6 seconds, and the next failure
dead-letters them.
Before 4.3 the same effect needed a TTL'd retry queue dead-lettering back into the work queue,
with messages rewritten on every hop. python retry.py inspect --drain empties the DLQ.
In production, put the delivery limit, DLX and delayed-retry settings in a policy rather than in queue arguments, so they can be tuned without redeclaring queues (see Operations).
3.2 Scheduled delivery¶
To deliver a message later, delayed.py uses delay tiers: a queue per delay, with a fixed
per-queue TTL and no consumers. Expired messages dead-letter to the real destination.
# py/delayed.py
def declare(ch) -> None:
ch.exchange_declare(exchange=EXCHANGE, exchange_type="direct", durable=True)
ch.queue_declare(queue=TARGET, durable=True, arguments={"x-queue-type": "quorum"})
ch.queue_bind(queue=TARGET, exchange=EXCHANGE, routing_key="due")
for name, ttl in TIERS.items():
ch.queue_declare(
queue=f"lab.delay.{name}",
durable=True,
arguments={
"x-queue-type": "quorum",
"x-message-ttl": ttl,
"x-dead-letter-exchange": EXCHANGE,
"x-dead-letter-routing-key": "due",
"x-dead-letter-strategy": "at-least-once",
"x-overflow": "reject-publish",
},
)
python delayed.py consume # T1
python delayed.py schedule "stand-up in 5" --delay 5s # T2
python delayed.py schedule "coffee" --delay 30s
Each reminder arrives after its delay, give or take a second. One queue per delay is deliberate: a queue only expires messages at its head, so a single queue with per-message TTLs would hold a 5-second message behind a 2-minute one. The once-popular delayed-message-exchange plugin is deprecated and archived, and was never safe for large numbers of messages; for retries, prefer 4.3's delayed retry (3.1).
3.3 Priorities¶
# py/priority.py
def queue_for(ch, classic: bool) -> str:
if classic:
name = "lab.priority.classic"
ch.queue_declare(queue=name, durable=True, arguments={"x-queue-type": "classic", "x-max-priority": 10})
else:
name = "lab.priority"
ch.queue_declare(queue=name, durable=True, arguments={"x-queue-type": "quorum"})
return name
python priority.py publish # with nothing consuming
python priority.py consume
Messages come out highest priority first: CRITICAL, urgent, medium, normal, low, background. On
4.3 quorum queues apply up to 32 priority levels strictly; on 4.0-4.2 they only distinguished
normal from high. Classic queues need x-max-priority declared up front (--classic), and each
level costs an internal sub-queue, so keep it small.
Now start the consumer first and publish again: messages arrive in publish order. Priority only reorders messages that are waiting in the queue, and a large prefetch has the same effect, because whatever is already in flight to the consumer can't be reordered.
3.4 Request/reply¶
rpc.py uses Direct Reply-to: the client consumes from the pseudo-queue
amq.rabbitmq.reply-to and puts that name in reply_to, so replies come straight back to its
channel without a reply queue to declare or clean up.
# py/rpc.py
class RpcClient:
def __init__(self, timeout: float) -> None:
self.timeout = timeout
self.conn = connect("rpc-client")
self.ch = self.conn.channel()
self.pending: dict[str, bytes | None] = {}
# Must consume from the pseudo-queue (auto_ack) BEFORE publishing requests.
self.ch.basic_consume(REPLY_TO, self._on_reply, auto_ack=True)
def _on_reply(self, _ch, _method, props, body) -> None:
if props.correlation_id in self.pending:
self.pending[props.correlation_id] = body
def call(self, n: int) -> int:
corr_id = str(uuid.uuid4())
self.pending[corr_id] = None
self.ch.basic_publish(
"", QUEUE, str(n).encode(),
pika.BasicProperties(
reply_to=REPLY_TO,
correlation_id=corr_id,
expiration=str(int(self.timeout * 1000)),
),
)
deadline = time.monotonic() + self.timeout
while self.pending[corr_id] is None:
if time.monotonic() > deadline:
del self.pending[corr_id]
raise TimeoutError(f"no reply for fib({n}) within {self.timeout}s; is the server running?")
self.conn.process_data_events(time_limit=0.1)
return int(self.pending.pop(corr_id))
def close(self) -> None:
self.conn.close()
python rpc.py server # T1
python rpc.py call 10 50 90 # T2
Each request carries an expiration equal to the client's timeout. Stop the server, run a call
(it times out after 5 seconds), then start the server: the stale request has expired and is
never answered, which is what you want, since nobody is waiting for it.
RPC over a broker is useful when the caller and the service shouldn't know about each other or when requests need queuing during spikes. For plain synchronous calls between services, HTTP or gRPC is usually simpler.
3.5 Streams¶
A stream keeps messages after they're consumed; each consumer chooses where to start reading.
# py/streams.py
def consume(args) -> None:
conn = connect("stream-consumer")
ch = conn.channel()
declare(ch)
ch.basic_qos(prefetch_count=100) # required for stream consumers
received = 0
def on_message(ch, method, props, body):
nonlocal received
offset = (props.headers or {}).get("x-stream-offset")
print(f" [<] offset={offset} {json.loads(body)}")
ch.basic_ack(method.delivery_tag)
received += 1
if args.limit and received >= args.limit:
ch.stop_consuming()
offset = parse_offset(args)
ch.basic_consume(STREAM, on_message, arguments={"x-stream-offset": offset})
print(f" [*] reading {STREAM} from offset={offset!s}")
run_consumer(ch)
python streams.py publish --count 100
python streams.py consume --offset first --limit 5 # replay from the beginning
python streams.py consume --offset 50 --limit 5 # from offset 50
python streams.py consume --offset next # T1: only new messages
python streams.py publish --count 10 # T2
python streams.py consume --since 60 --limit 5 # from a point in time
Run the first consumer twice: it reads the same messages both times, and the stream's message
count in the UI doesn't drop. Retention is by size here (x-max-length-bytes), enforced by
discarding whole segment files once the total exceeds the limit. Over AMQP 0-9-1 a stream
consumer must set a prefetch and must ack (acks only replenish credit); the native stream
protocol on port 5552 is far faster and adds deduplication, filtering and super streams.
Part 4: Observability¶
4.1 Reading the management UI¶
With a consumer and producer from Part 2 running, the numbers that matter are:
- Queues → Messages: Ready / Unacked. Ready is waiting for a consumer; unacked is delivered but not yet acknowledged. Healthy queues have small, stable numbers for both.
- Queues → Message rates. Incoming versus deliver/get versus ack. If incoming exceeds ack for long, the queue grows.
- Queue page → Consumer capacity. The share of time the queue could deliver immediately to a consumer. Well under 100% with a backlog means consumers aren't keeping up.
- Connections → State.
runningis normal,flowmeans the connection is being throttled by internal flow control,blockedmeans a resource alarm. - Channels → Prefetch, Unacked. Unlimited prefetch shows as 0.
4.2 The CLI and a triage script¶
scripts/triage.sh runs the read-only checks from the
troubleshooting playbooks in one pass: cluster status,
alarms, file descriptors and memory, the deepest queues, queues with messages but no consumers,
queues with many unacked messages, connections per user and host, channel states, the
quorum-critical check, deprecated features in use, and stuck processes.
python basics.py produce && python basics.py produce # leave a backlog, no consumer
make triage
lab.app-logs shows up in both the "deepest queues" and "no consumers" sections. The script
works against native installs too: RMQ_EXEC=sudo ./scripts/triage.sh.
4.3 Prometheus, Grafana and alerts¶
make down && make up-monitoring
Prometheus (http://localhost:9090) scrapes the node's built-in metrics endpoint twice: the
cheap aggregated /metrics, and /metrics/detailed filtered to the two per-queue metric
families the alerts use. Scraping /metrics/per-object instead would return every metric for
every queue, connection and channel, which gets expensive on a real broker.
# config/prometheus/single.yml
global:
scrape_interval: 15s
evaluation_interval: 15s
rule_files:
- /etc/prometheus/alerts.yml
scrape_configs:
# Aggregated node/cluster metrics: cheap, safe to scrape often.
- job_name: rabbitmq
static_configs:
- targets: ["rabbitmq:15692"]
# Per-queue metrics, limited to the metric families the alerts actually use.
# Scraping every per-object family on a broker with thousands of queues is expensive.
- job_name: rabbitmq-detailed
metrics_path: /metrics/detailed
params:
family: [queue_coarse_metrics, queue_consumer_count]
static_configs:
- targets: ["rabbitmq:15692"]
The alert rules cover node availability, the three resource alarms, file descriptors, connection churn, unroutable messages, backlogs, queues without consumers and unacknowledged build-up:
# config/prometheus/alerts.yml
groups:
- name: rabbitmq-node
rules:
- alert: RabbitMQNodeDown
expr: up{job="rabbitmq"} == 0
for: 1m
labels: {severity: critical}
annotations:
summary: "RabbitMQ node {{ $labels.instance }} is not reporting metrics"
- alert: RabbitMQMemoryAlarm
expr: rabbitmq_alarms_memory_used_watermark == 1
for: 1m
labels: {severity: critical}
annotations:
summary: "Memory alarm on {{ $labels.instance }}: publishers are blocked"
- alert: RabbitMQDiskAlarm
expr: rabbitmq_alarms_free_disk_space_watermark == 1
for: 1m
labels: {severity: critical}
annotations:
summary: "Disk alarm on {{ $labels.instance }}: publishers are blocked"
- alert: RabbitMQFileDescriptorsHigh
expr: rabbitmq_process_open_fds / rabbitmq_process_max_fds > 0.8
for: 5m
labels: {severity: warning}
annotations:
summary: "{{ $labels.instance }} is using {{ $value | humanizePercentage }} of its file descriptors"
# Opening a connection per message is the most common RabbitMQ anti-pattern.
- alert: RabbitMQConnectionChurn
expr: rate(rabbitmq_connections_opened_total[5m]) > 10
for: 10m
labels: {severity: warning}
annotations:
summary: "{{ $labels.instance }} is opening {{ $value | humanize }} connections/s; clients should reuse connections"
- alert: RabbitMQUnroutableMessages
expr: rate(rabbitmq_global_messages_unroutable_dropped_total[5m]) > 0
for: 5m
labels: {severity: warning}
annotations:
summary: "Messages published to {{ $labels.instance }} match no binding and are being dropped"
- name: rabbitmq-queues
rules:
- alert: RabbitMQQueueBacklog
expr: rabbitmq_detailed_queue_messages_ready > 10000
for: 10m
labels: {severity: warning}
annotations:
summary: "Queue {{ $labels.vhost }}/{{ $labels.queue }} has {{ $value }} ready messages"
- alert: RabbitMQQueueNoConsumers
expr: |
rabbitmq_detailed_queue_messages_ready > 0
and on (instance, vhost, queue)
rabbitmq_detailed_queue_consumers == 0
for: 10m
labels: {severity: warning}
annotations:
summary: "Queue {{ $labels.vhost }}/{{ $labels.queue }} has messages but no consumers"
- alert: RabbitMQUnackedHigh
expr: rabbitmq_detailed_queue_messages_unacked > 1000
for: 10m
labels: {severity: warning}
annotations:
summary: "Queue {{ $labels.vhost }}/{{ $labels.queue }} has {{ $value }} unacknowledged deliveries (slow or stuck consumers)"
In Grafana (http://localhost:3000, user admin, password from make creds) the Prometheus
datasource is already provisioned. Import the RabbitMQ team's dashboards with Dashboards →
New → Import: ID 10991 (RabbitMQ-Overview) and 11352 (Erlang-Distribution).
Try this: leave lab.app-logs with a backlog and no consumer. On Prometheus's Alerts
page, RabbitMQQueueNoConsumers goes pending within a scrape or two and firing after its
10-minute for window. Query rabbitmq_detailed_queue_messages_ready to see the per-queue
series behind it. Alertmanager isn't included; add it when you have somewhere real to route
alerts.
4.4 Health checks¶
scripts/healthcheck.sh is a Nagios/Icinga-style check (exit 0/1/2/3) built on the HTTP API's
/api/health/checks/ endpoints plus queue-level warnings:
make health
# WARNING: lab.app-logs (vhost /): 8 ready, no consumers
Outside the lab, run it with a dedicated monitoring-tagged user rather than an administrator.
Part 5: Performance¶
5.1 Resource alarms¶
Trigger a memory alarm on purpose by lowering the watermark at runtime:
docker compose exec rabbitmq rabbitmqctl set_vm_memory_high_watermark 0.0001
python basics.py produce
The producer hangs: the broker has sent connection.blocked and stopped reading from every
publishing connection. The UI shows the memory bar in red and connections as blocked, and
rabbitmq-diagnostics alarms lists the alarm. After 60 seconds the client's
blocked_connection_timeout gives up. Consumers aren't blocked; they're how a real alarm clears.
Restore it:
docker compose exec rabbitmq rabbitmqctl set_vm_memory_high_watermark 0.6
rabbitmqctl set_disk_free_limit 10000GB triggers the disk alarm the same way (put it back with
set_disk_free_limit 1GB). Runtime changes like these are lost on restart; the config file
holds the real values.
5.2 What confirms cost¶
# py/confirms.py
def run(ch, queue: str, count: int, body: bytes, label: str) -> None:
ch.queue_purge(queue)
props = pika.BasicProperties(delivery_mode=PERSISTENT)
start = time.perf_counter()
for _ in range(count):
ch.basic_publish("", queue, body, props)
elapsed = time.perf_counter() - start
print(f" {label:<26} {count / elapsed:>10,.0f} msg/s ({elapsed:.2f}s)")
python confirms.py --count 5000
python confirms.py --count 5000 --queue-type classic
Waiting for a confirm after every publish is several times slower than fire-and-forget, and the
gap is larger for quorum queues, which confirm only after a majority has written to disk.
That's pika's BlockingConnection limit, not the broker's. Asynchronous clients keep many
publishes in flight and process confirms as they arrive, recovering most of the throughput
without giving up safety.
5.3 PerfTest and prefetch¶
PerfTest is the RabbitMQ team's load generator. scripts/perftest.sh runs it as a container
on the lab network, through a series of 30-second scenarios: confirms in flight, a prefetch
sweep, large messages, and fan-in.
make perftest
./scripts/perftest.sh single -- --help # every PerfTest option
Compare the scenarios rather than the absolute numbers, which mostly measure your laptop. Look at how throughput changes between 1 and 100 confirms in flight, and how a prefetch of 1 starves a fast consumer while very large values stop helping. Then run your own:
./scripts/perftest.sh single -- --quorum-queue --queue perf.qq --producers 4 --consumers 4 \
--size 4096 --confirm 200 --qos 200 --time 60
Part 6: Clustering¶
6.1 Three nodes behind HAProxy¶
make down
make cluster-up
docker compose -f compose.cluster.yaml exec rmq1 rabbitmqctl cluster_status
The three nodes found each other at boot through peer discovery, with no join_cluster
commands:
# config/cluster/rabbitmq.conf
## Cluster-only settings, mounted as /etc/rabbitmq/conf.d/30-cluster.conf
## alongside the shared 20-lab.conf.
cluster_name = rmqlab
## Peer discovery: nodes find each other at boot, no manual join_cluster needed.
## Node names must match the container hostnames (rabbit@<hostname>).
cluster_formation.peer_discovery_backend = classic_config
cluster_formation.classic_config.nodes.1 = rabbit@rmq1
cluster_formation.classic_config.nodes.2 = rabbit@rmq2
cluster_formation.classic_config.nodes.3 = rabbit@rmq3
Clients keep using localhost:5672; HAProxy spreads connections across the nodes. Two settings
in its config matter in any RabbitMQ deployment: client/server timeouts well above the AMQP
heartbeat, and on-marked-down shutdown-sessions so clients of a failed node reconnect
immediately.
# config/haproxy.cfg
global
log stdout format raw local0
maxconn 4096
# Docker's embedded DNS, so backends resolve even if a node starts late or restarts
# with a new IP.
resolvers docker
nameserver dns 127.0.0.11:53
hold valid 10s
defaults
log global
mode tcp
option tcplog
timeout connect 5s
# Client/server timeouts MUST exceed the AMQP heartbeat interval (60s here),
# or HAProxy silently kills idle-but-healthy connections.
timeout client 3h
timeout server 3h
timeout check 5s
default-server check inter 5s rise 2 fall 3 resolvers docker init-addr last,libc,none
frontend amqp
bind :5672
default_backend rabbit_amqp
backend rabbit_amqp
# Long-lived connections: leastconn spreads them better than roundrobin.
balance leastconn
# When a node is marked down, cut its sessions so clients reconnect elsewhere
# immediately instead of waiting for heartbeats to time out.
default-server on-marked-down shutdown-sessions
server rmq1 rmq1:5672
server rmq2 rmq2:5672
server rmq3 rmq3:5672
frontend management
mode http
option httplog
timeout client 60s
bind :15672
default_backend rabbit_management
backend rabbit_management
mode http
timeout server 60s
balance roundrobin
option httpchk
http-check send meth GET uri /
http-check expect status 200
server rmq1 rmq1:15672
server rmq2 rmq2:15672
server rmq3 rmq3:15672
listen stats
mode http
timeout client 60s
timeout server 60s
bind :8404
stats enable
stats uri /
stats refresh 5s
HAProxy's stats page is at http://localhost:8404. Prometheus scrapes each node directly, never through HAProxy, since each node reports only its own metrics.
6.2 Quorum queue replicas¶
Run python failover.py consume for a moment and stop it; that declares lab.ha with three
members. Then look at them:
alias c1='docker compose -f compose.cluster.yaml exec rmq1'
c1 rabbitmq-queues quorum_status lab.ha
c1 rabbitmqctl list_queues name type leader members online
lab.ha has a member on each node, one of them the leader. Every quorum queue declared in the
cluster gets up to three members by default; leaders spread across nodes as queues are
declared.
6.3 Failover¶
python failover.py consume # T1
python failover.py publish # T2
c1 rabbitmqctl list_queues name leader # T3: which node leads lab.ha?
docker compose -f compose.cluster.yaml stop rmq2 # stop that node (rmq2 in this example)
Clients connected through the stopped node lose their connections and reconnect through HAProxy
to a surviving node within a few seconds; the remaining two members elect a new leader. The
publisher resumes from its last confirmed sequence number. Start the node again
(docker compose -f compose.cluster.yaml start rmq2) and it rejoins and catches up.
Stop both scripts with Ctrl-C. The consumer reports unique messages, duplicates and gaps: there should be no gaps (no confirmed message is ever lost) but possibly a few duplicates, from messages that were delivered but whose acks didn't make it before the failover, or published but not yet confirmed when the connection dropped. That's at-least-once delivery in practice, and why consumers need to be idempotent.
6.4 Maintenance mode¶
c1 rabbitmq-queues check_if_node_is_quorum_critical # safe to stop rmq1?
docker compose -f compose.cluster.yaml exec rmq2 rabbitmq-upgrade drain
c1 rabbitmqctl list_queues name leader
docker compose -f compose.cluster.yaml exec rmq2 rabbitmq-upgrade revive
c1 rabbitmq-queues rebalance quorum
Draining moves quorum queue leadership off rmq2, closes its client connections and stops it
from accepting new ones, without stopping the node. This is the first step of every rolling
upgrade. rebalance spreads leaders back out afterwards.
6.5 Losing the majority¶
With the publisher and consumer running, stop two nodes:
docker compose -f compose.cluster.yaml stop rmq2 rmq3
The surviving node is a minority. It can't commit anything to lab.ha or change metadata, so
the publisher stalls and reports failed or unconfirmed publishes. Nothing is inconsistent; the
cluster is unavailable, by design. Start the nodes again and everything resumes. Compare with
3.x, where the same experiment involved choosing a partition-handling strategy and sometimes
losing data on the "losing" side.
Part 7: An order pipeline¶
The capstone is an event-driven order flow: three services communicating only through events on one topic exchange, with a compensating action when a later step fails.
order.created ──▶ [payment] ──▶ payment.completed ──▶ [inventory] ──▶ order.completed
│ │
└─▶ payment.failed inventory.failed
payment.refund.requested ──▶ [payment] ──▶ payment.refunded
# py/orders.py
def declare(ch) -> None:
ch.exchange_declare(exchange=EXCHANGE, exchange_type="topic", durable=True)
ch.exchange_declare(exchange=DLX, exchange_type="fanout", durable=True)
ch.queue_declare(queue=DEAD, durable=True, arguments={"x-queue-type": "quorum", "x-delivery-limit": -1})
ch.queue_bind(queue=DEAD, exchange=DLX)
for queue, keys in QUEUES.items():
ch.queue_declare(
queue=queue,
durable=True,
arguments={
"x-queue-type": "quorum",
"x-delivery-limit": 5,
"x-dead-letter-exchange": DLX,
"x-dead-letter-strategy": "at-least-once",
"x-overflow": "reject-publish",
},
)
for key in keys:
ch.queue_bind(queue=queue, exchange=EXCHANGE, routing_key=key)
Every service follows the same discipline: consume, do the work, publish the outcome on a separate confirmed channel, and only then acknowledge the input. A crash in between causes a redelivery, so handlers are idempotent:
# py/orders.py
class Service:
"""One connection, a consuming channel and a separate confirmed publishing channel."""
def __init__(self, name: str, queue: str, fail_rate: float) -> None:
self.name, self.queue, self.fail_rate = name, queue, fail_rate
self.conn = connect(f"orders-{name}")
self.consume_ch = self.conn.channel()
declare(self.consume_ch)
self.consume_ch.basic_qos(prefetch_count=10)
self.publish_ch = self.conn.channel()
self.publish_ch.confirm_delivery()
self.done: set[tuple[str, str]] = set() # (routing key, order id) already handled
def emit(self, routing_key: str, order: dict) -> None:
order = {**order, "history": [*order.get("history", []), routing_key]}
self.publish_ch.basic_publish(
EXCHANGE, routing_key, to_json(order),
json_props(message_id=str(uuid.uuid4()), correlation_id=order["order_id"]),
mandatory=True,
)
def fails(self) -> bool:
return random.random() < self.fail_rate
def run(self, handler) -> None:
def on_message(ch, method, _props, body):
try:
order = json.loads(body)
key = (method.routing_key, order["order_id"])
except (json.JSONDecodeError, KeyError):
print(f" [{self.name}] malformed message, rejecting to the DLQ")
ch.basic_reject(method.delivery_tag, requeue=False)
return
if key in self.done: # redelivery of something already handled
ch.basic_ack(method.delivery_tag)
return
time.sleep(0.2) # simulated work
try:
handler(method.routing_key, order)
except Exception as exc: # unexpected bug or outage: retry, then dead-letter
print(f" [{self.name}] error on {order['order_id']}: {exc!r}; rejecting for retry")
ch.basic_reject(method.delivery_tag, requeue=True)
return
self.done.add(key)
ch.basic_ack(method.delivery_tag) # only after the outcome is confirmed
self.consume_ch.basic_consume(self.queue, on_message)
print(f" [{self.name}] consuming {self.queue}")
run_consumer(self.consume_ch)
python orders.py payment # T1
python orders.py inventory # T2
python orders.py notify # T3
python orders.py create --count 20
The notifier prints each order's final event and the path it took, for example
order.created -> payment.completed -> order.completed, or for a stock failure
... -> inventory.failed -> payment.refund.requested -> payment.refunded.
Note what's not dead-lettered: a declined card or an empty shelf is a normal business outcome,
published as an event like any other. The DLQ (python orders.py dead) only receives messages a
service couldn't handle at all: malformed input, or an unexpected error that persisted through
the delivery limit.
Try this:
- Stop the inventory service, create orders, and watch
lab.orders.inventorybuffer them. Start it again and they drain. The services never needed to be up at the same time. - Run
python orders.py inventory --fail-rate 0.5and follow the refunds. - Run two payment services. They compete on the same queue and the work splits between them.
- Publish
{"oops": true}tolab.orderswith routing keyorder.createdfrom the UI and find it in the DLQ.
Cleanup¶
make cluster-down # or: make down
make clean # both stacks, including volumes
Where next¶
- Swap
pikafor an asynchronous client (aio-pikain Python,lapinin Rust) and rerun Lab 5.2 with many publishes in flight. - Enable TLS on the listeners with a private CA and point the scripts at
amqps://. - Replace HAProxy and Compose with the RabbitMQ Cluster Operator on Kubernetes.
- Try the native stream protocol (
rabbitmq_streamplugin, port 5552) and a super stream. - Bridge two single-node stacks with the Shovel plugin to see cross-cluster replication.