Andrew Mercer
on this page

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 produce three times and watch Queues → lab.app-logs hold 12 ready messages. Run docker 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, payload not 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. running is normal, flow means the connection is being throttled by internal flow control, blocked means 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.inventory buffer 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.5 and follow the refunds.
  • Run two payment services. They compete on the same queue and the work splits between them.
  • Publish {"oops": true} to lab.orders with routing key order.created from the UI and find it in the DLQ.

Cleanup

make cluster-down      # or: make down
make clean             # both stacks, including volumes

Where next

  • Swap pika for an asynchronous client (aio-pika in Python, lapin in 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_stream plugin, port 5552) and a super stream.
  • Bridge two single-node stacks with the Shovel plugin to see cross-cluster replication.