Skip to content

Using the Blocking Connection with connection recovery with multiple hosts

RabbitMQ nodes can be clustered. In the absence of failure clients can connect to any node and perform any operation. In case a node fails, stops, or becomes unavailable, clients should be able to connect to another node and continue.

To simplify reconnection to a different node, connection recovery mechanism should be combined with connection configuration that specifies multiple hosts.

The BlockingConnection adapter relies on exception handling to check for connection errors:

"""
This module implements a client that connects to multiple RabbitMQ brokers distributed across
different ports (5672, 5673, 5674) and consumes messages from a shared queue.

The process follows these main steps:

1. It randomly selects one of the available RabbitMQ brokers (on ports 5672, 5673, and 5674)
  to establish a connection.
2. Declares a `direct` type exchange called 'test_exchange' and a queue named 'standard'.
3. The 'standard' queue is then bound to the exchange using the routing key 'standard_key'.
4. Configures the `on_message` function as a callback to process messages received from the queue.
5. Starts consuming messages, processing them, and acknowledges receipt via `basic_ack`.
6. If a connection error occurs, the system will attempt to reconnect automatically, except in cases
  of connection closure by the broker or channel errors.


The name of the module, `blocking_consume_recover_multiple_hosts.py`, reflects its key functionalities:
- "blocking_consume": The client uses a blocking connection to RabbitMQ and consumes messages synchronously.
- "recover": The module is designed to recover from connection errors by attempting to reconnect.
- "multiple_hosts": It connects to multiple RabbitMQ brokers distributed across different ports to ensure availability and redundancy.
"""

import functools
import random

import pika
from pika.exchange_type import ExchangeType


def on_message(ch, method_frame, _header_frame, body, userdata=None):
    print(f'Userdata: {userdata} Message body: {body}')
    ch.basic_ack(delivery_tag=method_frame.delivery_tag)


credentials = pika.PlainCredentials('guest', 'guest')

params1 = pika.ConnectionParameters('localhost',
                                    port=5672,
                                    credentials=credentials)
params2 = pika.ConnectionParameters('localhost',
                                    port=5673,
                                    credentials=credentials)
params3 = pika.ConnectionParameters('localhost',
                                    port=5674,
                                    credentials=credentials)
params_all = [params1, params2, params3]

# Infinite loop
while True:
    try:
        random.shuffle(params_all)
        connection = pika.BlockingConnection(params_all)
        channel = connection.channel()
        channel.exchange_declare(exchange='test_exchange',
                                 exchange_type=ExchangeType.direct,
                                 passive=False,
                                 durable=True,
                                 auto_delete=False)
        channel.queue_declare(queue='standard', durable=True, auto_delete=True)
        channel.queue_bind(queue='standard',
                           exchange='test_exchange',
                           routing_key='standard_key')
        channel.basic_qos(prefetch_count=1)

        on_message_callback = functools.partial(on_message,
                                                userdata='on_message_userdata')
        channel.basic_consume('standard', on_message_callback)

        try:
            channel.start_consuming()
        except KeyboardInterrupt:
            channel.stop_consuming()

        connection.close()
        break
    # Do not recover if connection was closed by broker
    except pika.exceptions.ConnectionClosedByBroker:
        break
    # Do not recover on channel errors
    except pika.exceptions.AMQPChannelError:
        break
    # Recover on all other connection errors
    except pika.exceptions.AMQPConnectionError:
        continue
Generic operation retry libraries such as retry can prove useful.

To run the following example, install the library first with pip install retry.

In this example the retry decorator is used to set up recovery with delay:

import random

from retry import retry

import pika


def on_message(channel, method_frame, header_frame, body):
    print(method_frame.delivery_tag)
    print(body)
    print()
    channel.basic_ack(delivery_tag=method_frame.delivery_tag)


## Assuming there are three hosts: host1, host2, and host3
node1 = pika.URLParameters('amqp://node1')
node2 = pika.URLParameters('amqp://node2')
node3 = pika.URLParameters('amqp://node3')
all_endpoints = [node1, node2, node3]


@retry(pika.exceptions.AMQPConnectionError, delay=5, jitter=(1, 3))
def consume():
    random.shuffle(all_endpoints)
    connection = pika.BlockingConnection(all_endpoints)
    channel = connection.channel()
    channel.basic_qos(prefetch_count=1)

    ## This queue is intentionally non-durable. See https://www.rabbitmq.com/ha.html#non-mirrored-queue-behavior-on-node-failure
    ## to learn more.
    channel.queue_declare('recovery-example', durable=False, auto_delete=True)
    channel.basic_consume('recovery-example', on_message)

    try:
        channel.start_consuming()
    except KeyboardInterrupt:
        channel.stop_consuming()
        connection.close()
    except pika.exceptions.ConnectionClosedByBroker:
        # Uncomment this to make the example not attempt recovery
        # from server-initiated connection closure, including
        # when the node is stopped cleanly
        # except pika.exceptions.ConnectionClosedByBroker:
        #     pass
        pass


consume()