Available for day contractsFrom 21st September I have availability for day and half day contracts. Please contact for more information.

Contact →
mikepreston.org

Python pika (RabbitMQ)

Essential patterns and commands for building message-driven applications with RabbitMQ using the pika library.

Python pika (RabbitMQ)

Essential patterns and commands for building message-driven applications with RabbitMQ using the pika library.

Overview

Pika is a pure-Python implementation of the AMQP 0-9-1 protocol for connecting to RabbitMQ message brokers. It provides both blocking and asynchronous connection adapters, enabling flexible message publishing and consumption patterns for distributed systems.

RabbitMQ BrokerApplicationpublishrouterouterouteconsumeconsumeconsumePublisherConsumerExchangeQueue 1Queue 2Queue 3RabbitMQ BrokerApplicationpublishrouterouterouteconsumeconsumeconsumePublisherConsumerExchangeQueue 1Queue 2Queue 3

Connection and Channel Setup

Connections manage TCP connections to RabbitMQ, while channels are virtual connections that multiplex over a single TCP connection.

Key Concepts

RabbitMQConnectionApplicationChannel 1Channel 2Channel 3TCP SocketBrokerRabbitMQConnectionApplicationChannel 1Channel 2Channel 3TCP SocketBroker
  • Connection: TCP connection to the broker; expensive to create
  • Channel: Lightweight virtual connection; use one per thread
  • Credentials: Authentication with username/password or external mechanisms
  • Virtual Host: Logical grouping of resources within a broker
  • Heartbeat: Keep-alive mechanism to detect dead connections

Basic Connection (BlockingConnection)

import pika

# Simple connection with defaults
connection = pika.BlockingConnection(
    pika.ConnectionParameters('localhost')
)
channel = connection.channel()

# Always close connections when done
channel.close()
connection.close()

Connection with Authentication

import pika

# Connection with credentials
credentials = pika.PlainCredentials('username', 'password')

parameters = pika.ConnectionParameters(
    host='rabbitmq.example.com',
    port=5672,
    virtual_host='/',
    credentials=credentials,
    heartbeat=600,                    # Heartbeat interval in seconds
    blocked_connection_timeout=300,   # Timeout for blocked connections
    connection_attempts=3,            # Number of connection attempts
    retry_delay=5                     # Delay between retries
)

connection = pika.BlockingConnection(parameters)
channel = connection.channel()

Connection with SSL/TLS

import pika
import ssl

# SSL context configuration
ssl_context = ssl.create_default_context(cafile='/path/to/ca_certificate.pem')
ssl_context.load_cert_chain(
    '/path/to/client_certificate.pem',
    '/path/to/client_key.pem'
)

ssl_options = pika.SSLOptions(ssl_context, 'rabbitmq.example.com')

parameters = pika.ConnectionParameters(
    host='rabbitmq.example.com',
    port=5671,                        # Default SSL port
    ssl_options=ssl_options,
    credentials=pika.PlainCredentials('user', 'pass')
)

connection = pika.BlockingConnection(parameters)

URL-Based Connection

import pika

# Connection using AMQP URL
url = 'amqp://user:pass@rabbitmq.example.com:5672/%2F'  # %2F = /
parameters = pika.URLParameters(url)

# With additional options
url = 'amqp://user:pass@host:5672/%2F?heartbeat=600&connection_attempts=3'
parameters = pika.URLParameters(url)

connection = pika.BlockingConnection(parameters)

Context Manager Pattern

import pika
from contextlib import contextmanager

@contextmanager
def rabbitmq_connection(host='localhost'):
    """Context manager for RabbitMQ connections."""
    connection = pika.BlockingConnection(
        pika.ConnectionParameters(host)
    )
    channel = connection.channel()
    try:
        yield channel
    finally:
        channel.close()
        connection.close()

# Usage
with rabbitmq_connection('localhost') as channel:
    channel.queue_declare(queue='my_queue')
    channel.basic_publish(
        exchange='',
        routing_key='my_queue',
        body='Hello World'
    )

SelectConnection (Asynchronous)

import pika

def on_open(connection):
    """Called when connection is opened."""
    connection.channel(on_open_callback=on_channel_open)

def on_channel_open(channel):
    """Called when channel is opened."""
    channel.queue_declare(queue='async_queue', callback=on_queue_declared)

def on_queue_declared(frame):
    """Called when queue is declared."""
    print(f"Queue declared: {frame.method.queue}")

# Create connection with callbacks
parameters = pika.ConnectionParameters('localhost')
connection = pika.SelectConnection(
    parameters,
    on_open_callback=on_open
)

# Start the I/O loop
connection.ioloop.start()

Publishing Messages

Publishers send messages to exchanges, which route them to queues based on bindings and routing keys.

Key Concepts

routing_keybinding_key matchbinding_key matchPublisherExchangeQueue AQueue Brouting_keybinding_key matchbinding_key matchPublisherExchangeQueue AQueue B
  • Exchange: Routes messages to queues based on type and bindings
  • Routing Key: Address used by exchange to route messages
  • Mandatory: Return message if it cannot be routed
  • Persistent: Message survives broker restart
  • Properties: Message metadata (content type, headers, etc.)

Exchange Types

Type Routing Behaviour Use Case
direct Exact routing key match Task queues, RPC
fanout Broadcast to all bound queues Pub/sub, broadcast
topic Pattern matching with wildcards Log routing, events
headers Match on message headers Complex routing

Basic Publishing

import pika

connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
channel = connection.channel()

# Declare queue (idempotent)
channel.queue_declare(queue='task_queue')

# Simple publish to default exchange
channel.basic_publish(
    exchange='',              # Default exchange
    routing_key='task_queue', # Queue name for default exchange
    body='Hello World'
)

connection.close()

Publishing with Properties

import pika
import json

connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
channel = connection.channel()

# Message with properties
message = {'task': 'process_data', 'data': [1, 2, 3]}

properties = pika.BasicProperties(
    delivery_mode=pika.DeliveryMode.Persistent,  # Make message persistent
    content_type='application/json',
    content_encoding='utf-8',
    headers={'x-custom-header': 'value'},
    priority=5,                                   # 0-9, higher = more priority
    correlation_id='abc123',                      # For RPC patterns
    reply_to='response_queue',                    # For RPC patterns
    expiration='60000',                           # TTL in milliseconds
    message_id='msg-001',
    timestamp=int(time.time()),
    app_id='my_application'
)

channel.basic_publish(
    exchange='',
    routing_key='task_queue',
    body=json.dumps(message),
    properties=properties
)

connection.close()

Direct Exchange

import pika

connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
channel = connection.channel()

# Declare direct exchange
channel.exchange_declare(
    exchange='direct_logs',
    exchange_type='direct',
    durable=True
)

# Declare and bind queues
for severity in ['info', 'warning', 'error']:
    channel.queue_declare(queue=f'{severity}_queue', durable=True)
    channel.queue_bind(
        exchange='direct_logs',
        queue=f'{severity}_queue',
        routing_key=severity
    )

# Publish to specific severity
channel.basic_publish(
    exchange='direct_logs',
    routing_key='error',              # Routes to error_queue
    body='This is an error message',
    properties=pika.BasicProperties(delivery_mode=2)
)

connection.close()

Fanout Exchange (Broadcast)

import pika

connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
channel = connection.channel()

# Declare fanout exchange
channel.exchange_declare(
    exchange='notifications',
    exchange_type='fanout',
    durable=True
)

# Publish to all bound queues (routing_key ignored)
channel.basic_publish(
    exchange='notifications',
    routing_key='',                   # Ignored for fanout
    body='Broadcast message to all subscribers'
)

connection.close()

Topic Exchange (Pattern Matching)

import pika

connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
channel = connection.channel()

# Declare topic exchange
channel.exchange_declare(
    exchange='topic_logs',
    exchange_type='topic',
    durable=True
)

# Bind queues with patterns
# * matches one word, # matches zero or more words
channel.queue_declare(queue='all_logs', durable=True)
channel.queue_bind(
    exchange='topic_logs',
    queue='all_logs',
    routing_key='#'                   # Match all messages
)

channel.queue_declare(queue='kernel_logs', durable=True)
channel.queue_bind(
    exchange='topic_logs',
    queue='kernel_logs',
    routing_key='*.kernel.*'          # Match *.kernel.*
)

# Publish with routing key
channel.basic_publish(
    exchange='topic_logs',
    routing_key='server1.kernel.error',  # Matches *.kernel.*
    body='Kernel error on server1'
)

connection.close()

Publisher Confirms

import pika

connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
channel = connection.channel()

# Enable publisher confirms
channel.confirm_delivery()

# Publish with confirmation
try:
    channel.basic_publish(
        exchange='',
        routing_key='confirmed_queue',
        body='Message with confirmation',
        properties=pika.BasicProperties(delivery_mode=2),
        mandatory=True                    # Return if unroutable
    )
    print('Message was confirmed')
except pika.exceptions.UnroutableError:
    print('Message was returned - no route to queue')

connection.close()

Mandatory Flag and Returns

import pika

def on_return(channel, method, properties, body):
    """Handle returned messages."""
    print(f"Message returned: {body.decode()}")
    print(f"Reply code: {method.reply_code}, Reply text: {method.reply_text}")

connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
channel = connection.channel()

# Register return callback
channel.add_on_return_callback(on_return)
channel.confirm_delivery()

# This will be returned if queue doesn't exist
channel.basic_publish(
    exchange='',
    routing_key='nonexistent_queue',
    body='This might be returned',
    mandatory=True                        # Enable returns
)

# Process returns
connection.process_data_events()
connection.close()

Consuming Messages

Consumers receive messages from queues using either push (basic_consume) or pull (basic_get) methods.

Key Concepts

DeliverAckNack + RequeueNack + DiscardReject + RequeueReadyUnackedDeliverAckNack + RequeueNack + DiscardReject + RequeueReadyUnacked
  • Consumer Tag: Unique identifier for a consumer
  • Prefetch: Number of unacknowledged messages allowed
  • Auto-Ack: Automatically acknowledge on delivery
  • Exclusive: Only this consumer can access the queue
  • Consumer Priority: Higher priority consumers receive messages first

Basic Consumer

import pika

def callback(channel, method, properties, body):
    """Process received message."""
    print(f"Received: {body.decode()}")
    # Acknowledge after processing
    channel.basic_ack(delivery_tag=method.delivery_tag)

connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
channel = connection.channel()

channel.queue_declare(queue='my_queue')

# Start consuming
channel.basic_consume(
    queue='my_queue',
    on_message_callback=callback,
    auto_ack=False                        # Manual acknowledgement
)

print('Waiting for messages...')
channel.start_consuming()

Consumer with Prefetch (QoS)

import pika

def callback(channel, method, properties, body):
    """Process message with simulated work."""
    print(f"Processing: {body.decode()}")
    import time
    time.sleep(1)                         # Simulate work
    channel.basic_ack(delivery_tag=method.delivery_tag)
    print("Done")

connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
channel = connection.channel()

# Set prefetch count - don't give me more than 10 unacked messages
channel.basic_qos(prefetch_count=10)

channel.queue_declare(queue='work_queue')

channel.basic_consume(
    queue='work_queue',
    on_message_callback=callback,
    auto_ack=False
)

channel.start_consuming()

Multiple Queue Consumer

import pika

def high_priority_callback(channel, method, properties, body):
    """Handle high priority messages."""
    print(f"[HIGH] {body.decode()}")
    channel.basic_ack(delivery_tag=method.delivery_tag)

def low_priority_callback(channel, method, properties, body):
    """Handle low priority messages."""
    print(f"[LOW] {body.decode()}")
    channel.basic_ack(delivery_tag=method.delivery_tag)

connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
channel = connection.channel()

channel.queue_declare(queue='high_priority')
channel.queue_declare(queue='low_priority')

# Consume from multiple queues
channel.basic_consume(queue='high_priority', on_message_callback=high_priority_callback)
channel.basic_consume(queue='low_priority', on_message_callback=low_priority_callback)

channel.start_consuming()

Polling with basic_get

import pika
import time

connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
channel = connection.channel()

channel.queue_declare(queue='poll_queue')

while True:
    # Poll for a single message
    method, properties, body = channel.basic_get(
        queue='poll_queue',
        auto_ack=False
    )

    if method:
        print(f"Got message: {body.decode()}")
        channel.basic_ack(delivery_tag=method.delivery_tag)
    else:
        print("No message available")
        time.sleep(1)

Consumer with Timeout

import pika
import functools

def callback(channel, method, properties, body):
    print(f"Received: {body.decode()}")
    channel.basic_ack(delivery_tag=method.delivery_tag)

connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
channel = connection.channel()

channel.queue_declare(queue='timeout_queue')

# Add timeout to stop consuming
connection.call_later(30, lambda: channel.stop_consuming())

channel.basic_consume(
    queue='timeout_queue',
    on_message_callback=callback,
    auto_ack=False
)

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

connection.close()

Exclusive Consumer

import pika

connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
channel = connection.channel()

# Declare exclusive queue (deleted when consumer disconnects)
result = channel.queue_declare(queue='', exclusive=True)
exclusive_queue = result.method.queue

def callback(channel, method, properties, body):
    print(f"Exclusive consumer received: {body.decode()}")
    channel.basic_ack(delivery_tag=method.delivery_tag)

channel.basic_consume(
    queue=exclusive_queue,
    on_message_callback=callback,
    exclusive=True,                       # Only this consumer
    auto_ack=False
)

channel.start_consuming()

Cancel Consumer

import pika

connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
channel = connection.channel()

channel.queue_declare(queue='cancel_queue')

def callback(channel, method, properties, body):
    print(f"Received: {body.decode()}")
    channel.basic_ack(delivery_tag=method.delivery_tag)

# Start consuming and get consumer tag
consumer_tag = channel.basic_consume(
    queue='cancel_queue',
    on_message_callback=callback,
    auto_ack=False
)

# Later, cancel the consumer
channel.basic_cancel(consumer_tag)
print(f"Cancelled consumer: {consumer_tag}")

connection.close()

Queue Operations

Queues store messages until they are consumed. They can be configured with various properties for durability, TTL, and message limits.

Key Concepts

  • Durable: Queue survives broker restart
  • Exclusive: Queue used by one connection only, deleted when closed
  • Auto-delete: Queue deleted when last consumer unsubscribes
  • Arguments: Additional queue configuration (TTL, limits, etc.)

Declaring Queues

import pika

connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
channel = connection.channel()

# Simple queue declaration
channel.queue_declare(queue='simple_queue')

# Durable queue (survives restart)
channel.queue_declare(queue='durable_queue', durable=True)

# Exclusive queue (auto-generated name, deleted on disconnect)
result = channel.queue_declare(queue='', exclusive=True)
temp_queue = result.method.queue
print(f"Created exclusive queue: {temp_queue}")

# Queue with TTL and max length
channel.queue_declare(
    queue='limited_queue',
    durable=True,
    arguments={
        'x-message-ttl': 60000,           # Message TTL: 60 seconds
        'x-max-length': 1000,             # Max 1000 messages
        'x-max-length-bytes': 10485760,   # Max 10MB
        'x-overflow': 'reject-publish',   # Reject new messages when full
        'x-queue-type': 'classic'         # classic or quorum
    }
)

connection.close()

Queue Arguments Reference

Argument Description Example Value
x-message-ttl Message expiration (ms) 60000
x-max-length Maximum queue length 10000
x-max-length-bytes Maximum queue size in bytes 10485760
x-overflow Behaviour when full: drop-head, reject-publish 'reject-publish'
x-dead-letter-exchange Exchange for rejected messages 'dlx'
x-dead-letter-routing-key Routing key for dead letters 'dead'
x-max-priority Enable priority queue (1-255) 10
x-queue-type Queue type: classic, quorum 'quorum'

Binding Queues to Exchanges

import pika

connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
channel = connection.channel()

# Declare exchange and queue
channel.exchange_declare(exchange='events', exchange_type='topic', durable=True)
channel.queue_declare(queue='user_events', durable=True)

# Bind queue to exchange with routing key
channel.queue_bind(
    exchange='events',
    queue='user_events',
    routing_key='user.*'                  # Matches user.created, user.updated, etc.
)

# Multiple bindings for same queue
channel.queue_bind(exchange='events', queue='user_events', routing_key='account.*')

# Unbind if needed
channel.queue_unbind(
    exchange='events',
    queue='user_events',
    routing_key='account.*'
)

connection.close()

Dead Letter Queue Setup

import pika

connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
channel = connection.channel()

# Declare dead letter exchange and queue
channel.exchange_declare(exchange='dlx', exchange_type='direct', durable=True)
channel.queue_declare(queue='dead_letters', durable=True)
channel.queue_bind(exchange='dlx', queue='dead_letters', routing_key='dead')

# Declare main queue with dead letter configuration
channel.queue_declare(
    queue='main_queue',
    durable=True,
    arguments={
        'x-dead-letter-exchange': 'dlx',
        'x-dead-letter-routing-key': 'dead',
        'x-message-ttl': 30000            # Messages expire after 30s
    }
)

connection.close()

Purging Queues

import pika

connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
channel = connection.channel()

# Purge all messages from queue
message_count = channel.queue_purge(queue='my_queue')
print(f"Purged {message_count.method.message_count} messages")

connection.close()

Deleting Queues

import pika

connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
channel = connection.channel()

# Delete queue unconditionally
channel.queue_delete(queue='old_queue')

# Delete only if empty
channel.queue_delete(queue='maybe_empty', if_empty=True)

# Delete only if no consumers
channel.queue_delete(queue='maybe_unused', if_unused=True)

connection.close()

Getting Queue Information

import pika

connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
channel = connection.channel()

# Passive declare to check if queue exists and get info
try:
    result = channel.queue_declare(queue='my_queue', passive=True)
    print(f"Queue: {result.method.queue}")
    print(f"Messages: {result.method.message_count}")
    print(f"Consumers: {result.method.consumer_count}")
except pika.exceptions.ChannelClosedByBroker as e:
    print(f"Queue does not exist: {e}")

connection.close()

Priority Queues

import pika

connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
channel = connection.channel()

# Declare priority queue
channel.queue_declare(
    queue='priority_queue',
    durable=True,
    arguments={'x-max-priority': 10}      # Priority range 0-10
)

# Publish with priority
channel.basic_publish(
    exchange='',
    routing_key='priority_queue',
    body='High priority message',
    properties=pika.BasicProperties(
        delivery_mode=2,
        priority=9                        # High priority
    )
)

channel.basic_publish(
    exchange='',
    routing_key='priority_queue',
    body='Low priority message',
    properties=pika.BasicProperties(
        delivery_mode=2,
        priority=1                        # Low priority
    )
)

connection.close()

Message Acknowledgments

Acknowledgments ensure reliable message delivery by confirming that messages have been processed successfully.

Key Concepts

SuccessRetryDiscardReject SingleMessage DeliveredProcess Messagebasic_ackbasic_nack + requeuebasic_nack + discardbasic_rejectMessage RemovedMessage RequeuedDead Letter QueueSuccessRetryDiscardReject SingleMessage DeliveredProcess Messagebasic_ackbasic_nack + requeuebasic_nack + discardbasic_rejectMessage RemovedMessage RequeuedDead Letter Queue
  • ack: Positive acknowledgment - message processed successfully
  • nack: Negative acknowledgment - can requeue or discard
  • reject: Similar to nack but for single message
  • delivery_tag: Unique identifier for the delivery
  • multiple: Acknowledge all messages up to delivery_tag

Basic Acknowledgment

import pika

def callback(channel, method, properties, body):
    try:
        # Process message
        print(f"Processing: {body.decode()}")

        # Acknowledge successful processing
        channel.basic_ack(delivery_tag=method.delivery_tag)
        print("Acknowledged")

    except Exception as e:
        print(f"Error: {e}")
        # Reject and requeue on error
        channel.basic_nack(
            delivery_tag=method.delivery_tag,
            requeue=True
        )

connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
channel = connection.channel()

channel.queue_declare(queue='ack_queue')
channel.basic_qos(prefetch_count=1)

channel.basic_consume(
    queue='ack_queue',
    on_message_callback=callback,
    auto_ack=False                        # Important: disable auto-ack
)

channel.start_consuming()

Multiple Acknowledgments

import pika

messages_processed = 0
batch_size = 10

def callback(channel, method, properties, body):
    global messages_processed

    # Process message
    print(f"Processing: {body.decode()}")
    messages_processed += 1

    # Batch acknowledge every 10 messages
    if messages_processed >= batch_size:
        channel.basic_ack(
            delivery_tag=method.delivery_tag,
            multiple=True                 # Ack all up to this tag
        )
        messages_processed = 0
        print(f"Batch acknowledged {batch_size} messages")

connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
channel = connection.channel()

channel.queue_declare(queue='batch_queue')
channel.basic_qos(prefetch_count=batch_size)

channel.basic_consume(
    queue='batch_queue',
    on_message_callback=callback,
    auto_ack=False
)

channel.start_consuming()

Negative Acknowledgment (nack)

import pika

retry_counts = {}

def callback(channel, method, properties, body):
    message_id = properties.message_id or str(method.delivery_tag)

    try:
        # Simulate potential failure
        process_message(body)
        channel.basic_ack(delivery_tag=method.delivery_tag)

    except RecoverableError:
        # Requeue for retry
        channel.basic_nack(
            delivery_tag=method.delivery_tag,
            requeue=True
        )

    except PermanentError:
        # Don't requeue - send to dead letter queue
        channel.basic_nack(
            delivery_tag=method.delivery_tag,
            requeue=False
        )

connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
channel = connection.channel()

channel.basic_consume(queue='nack_queue', on_message_callback=callback, auto_ack=False)
channel.start_consuming()

Reject (Single Message)

import pika

def callback(channel, method, properties, body):
    try:
        process_message(body)
        channel.basic_ack(delivery_tag=method.delivery_tag)

    except Exception as e:
        # Reject single message (requeue=False sends to DLQ if configured)
        channel.basic_reject(
            delivery_tag=method.delivery_tag,
            requeue=False                 # Don't requeue, send to DLQ
        )

connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
channel = connection.channel()

channel.basic_consume(queue='reject_queue', on_message_callback=callback, auto_ack=False)
channel.start_consuming()

Retry with Delay Pattern

import pika
import json

def setup_retry_infrastructure(channel):
    """Set up retry queues with increasing delays."""

    # Main exchange
    channel.exchange_declare(exchange='main', exchange_type='direct', durable=True)

    # Retry exchange
    channel.exchange_declare(exchange='retry', exchange_type='direct', durable=True)

    # Dead letter exchange
    channel.exchange_declare(exchange='dlx', exchange_type='direct', durable=True)

    # Main queue
    channel.queue_declare(
        queue='main_queue',
        durable=True,
        arguments={
            'x-dead-letter-exchange': 'retry',
            'x-dead-letter-routing-key': 'retry.1'
        }
    )
    channel.queue_bind(exchange='main', queue='main_queue', routing_key='main')

    # Retry queues with delays
    delays = [5000, 15000, 60000]  # 5s, 15s, 60s
    for i, delay in enumerate(delays, 1):
        queue_name = f'retry_queue_{i}'
        next_routing = f'retry.{i+1}' if i < len(delays) else 'dead'
        next_exchange = 'retry' if i < len(delays) else 'dlx'

        channel.queue_declare(
            queue=queue_name,
            durable=True,
            arguments={
                'x-message-ttl': delay,
                'x-dead-letter-exchange': 'main' if i < len(delays) else 'dlx',
                'x-dead-letter-routing-key': 'main' if i < len(delays) else 'dead'
            }
        )
        channel.queue_bind(exchange='retry', queue=queue_name, routing_key=f'retry.{i}')

    # Dead letter queue
    channel.queue_declare(queue='dead_letters', durable=True)
    channel.queue_bind(exchange='dlx', queue='dead_letters', routing_key='dead')

def callback(channel, method, properties, body):
    """Process with retry tracking."""
    headers = properties.headers or {}
    retry_count = headers.get('x-retry-count', 0)
    max_retries = 3

    try:
        process_message(body)
        channel.basic_ack(delivery_tag=method.delivery_tag)

    except Exception as e:
        if retry_count < max_retries:
            # Reject to retry queue
            channel.basic_nack(
                delivery_tag=method.delivery_tag,
                requeue=False
            )
        else:
            # Max retries exceeded, send to DLQ
            channel.basic_nack(
                delivery_tag=method.delivery_tag,
                requeue=False
            )

Error Handling and Reconnection Strategies

Robust error handling and automatic reconnection are essential for production message-driven applications.

Key Concepts

Connection LostAuto-reconnectSuccessRetryMax RetriesConnectedDisconnectedReconnectingFailedConnection LostAuto-reconnectSuccessRetryMax RetriesConnectedDisconnectedReconnectingFailed
  • Connection errors: Network issues, broker unavailable
  • Channel errors: Protocol violations, resource limits
  • Heartbeat: Detect dead connections
  • Exponential backoff: Gradually increase retry delays

Basic Error Handling

import pika
from pika.exceptions import (
    AMQPConnectionError,
    AMQPChannelError,
    ChannelClosedByBroker,
    ConnectionClosedByBroker
)

def connect():
    try:
        connection = pika.BlockingConnection(
            pika.ConnectionParameters('localhost')
        )
        channel = connection.channel()
        return connection, channel

    except AMQPConnectionError as e:
        print(f"Connection failed: {e}")
        raise
    except AMQPChannelError as e:
        print(f"Channel error: {e}")
        raise

def publish_message(channel, message):
    try:
        channel.basic_publish(
            exchange='',
            routing_key='my_queue',
            body=message
        )
    except ChannelClosedByBroker as e:
        print(f"Channel closed by broker: {e.reply_code} - {e.reply_text}")
        raise
    except ConnectionClosedByBroker as e:
        print(f"Connection closed by broker: {e.reply_code} - {e.reply_text}")
        raise

Reconnecting Consumer

import pika
import time
import logging

logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)

class ReconnectingConsumer:
    """Consumer with automatic reconnection."""

    def __init__(self, host='localhost', queue='my_queue'):
        self.host = host
        self.queue = queue
        self.connection = None
        self.channel = None
        self.should_reconnect = True
        self.reconnect_delay = 0

    def connect(self):
        """Establish connection to RabbitMQ."""
        logger.info(f"Connecting to {self.host}")
        parameters = pika.ConnectionParameters(
            host=self.host,
            heartbeat=600,
            blocked_connection_timeout=300
        )
        self.connection = pika.BlockingConnection(parameters)
        self.channel = self.connection.channel()
        self.channel.queue_declare(queue=self.queue, durable=True)
        self.channel.basic_qos(prefetch_count=1)
        self.reconnect_delay = 0
        logger.info("Connected successfully")

    def on_message(self, channel, method, properties, body):
        """Process received message."""
        try:
            logger.info(f"Received: {body.decode()}")
            # Process message here
            channel.basic_ack(delivery_tag=method.delivery_tag)
        except Exception as e:
            logger.error(f"Error processing message: {e}")
            channel.basic_nack(delivery_tag=method.delivery_tag, requeue=True)

    def run(self):
        """Main consumer loop with reconnection."""
        while self.should_reconnect:
            try:
                self.connect()
                self.channel.basic_consume(
                    queue=self.queue,
                    on_message_callback=self.on_message,
                    auto_ack=False
                )
                logger.info("Starting to consume")
                self.channel.start_consuming()

            except pika.exceptions.ConnectionClosedByBroker:
                logger.warning("Connection closed by broker")
                self._maybe_reconnect()

            except pika.exceptions.AMQPChannelError as e:
                logger.error(f"Channel error: {e}")
                self._maybe_reconnect()

            except pika.exceptions.AMQPConnectionError:
                logger.warning("Connection lost")
                self._maybe_reconnect()

            except KeyboardInterrupt:
                logger.info("Shutting down")
                self.should_reconnect = False
                if self.connection and self.connection.is_open:
                    self.connection.close()
                break

    def _maybe_reconnect(self):
        """Handle reconnection with exponential backoff."""
        if self.should_reconnect:
            self.reconnect_delay = min(self.reconnect_delay + 1, 30)
            logger.info(f"Reconnecting in {self.reconnect_delay} seconds...")
            time.sleep(self.reconnect_delay)

    def stop(self):
        """Stop the consumer."""
        self.should_reconnect = False
        if self.channel:
            self.channel.stop_consuming()

# Usage
if __name__ == '__main__':
    consumer = ReconnectingConsumer(host='localhost', queue='work_queue')
    consumer.run()

Reconnecting Publisher

import pika
import time
import logging

logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)

class ReconnectingPublisher:
    """Publisher with automatic reconnection."""

    def __init__(self, host='localhost'):
        self.host = host
        self.connection = None
        self.channel = None

    def connect(self):
        """Establish connection."""
        for attempt in range(5):
            try:
                logger.info(f"Connection attempt {attempt + 1}")
                parameters = pika.ConnectionParameters(
                    host=self.host,
                    heartbeat=600,
                    connection_attempts=3,
                    retry_delay=2
                )
                self.connection = pika.BlockingConnection(parameters)
                self.channel = self.connection.channel()
                self.channel.confirm_delivery()
                logger.info("Connected successfully")
                return True
            except pika.exceptions.AMQPConnectionError as e:
                logger.error(f"Connection failed: {e}")
                time.sleep(2 ** attempt)  # Exponential backoff
        return False

    def publish(self, exchange, routing_key, body, properties=None):
        """Publish with automatic reconnection."""
        for attempt in range(3):
            try:
                if not self.connection or self.connection.is_closed:
                    if not self.connect():
                        raise Exception("Failed to connect")

                self.channel.basic_publish(
                    exchange=exchange,
                    routing_key=routing_key,
                    body=body,
                    properties=properties or pika.BasicProperties(delivery_mode=2),
                    mandatory=True
                )
                logger.info(f"Published to {routing_key}")
                return True

            except pika.exceptions.UnroutableError:
                logger.error("Message was returned")
                return False

            except (pika.exceptions.AMQPConnectionError,
                    pika.exceptions.AMQPChannelError) as e:
                logger.warning(f"Publish failed: {e}, retrying...")
                self.connection = None
                time.sleep(1)

        return False

    def close(self):
        """Close connection."""
        if self.connection and self.connection.is_open:
            self.connection.close()

# Usage
publisher = ReconnectingPublisher('localhost')
publisher.publish('', 'my_queue', 'Hello World')
publisher.close()

SelectConnection with Reconnection

import pika
import logging
import functools

logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)

class AsyncReconnectingConsumer:
    """Asynchronous consumer with automatic reconnection."""

    def __init__(self, host='localhost', queue='my_queue'):
        self.host = host
        self.queue = queue
        self.connection = None
        self.channel = None
        self.closing = False
        self.consumer_tag = None
        self.reconnect_delay = 0

    def connect(self):
        """Connect to RabbitMQ."""
        logger.info(f"Connecting to {self.host}")
        return pika.SelectConnection(
            pika.ConnectionParameters(self.host),
            on_open_callback=self.on_connection_open,
            on_open_error_callback=self.on_connection_open_error,
            on_close_callback=self.on_connection_closed
        )

    def on_connection_open(self, connection):
        """Called when connection is established."""
        logger.info("Connection opened")
        self.reconnect_delay = 0
        self.connection = connection
        self.connection.channel(on_open_callback=self.on_channel_open)

    def on_connection_open_error(self, connection, error):
        """Called when connection fails."""
        logger.error(f"Connection open failed: {error}")
        self.reconnect()

    def on_connection_closed(self, connection, reason):
        """Called when connection is closed."""
        self.channel = None
        if self.closing:
            connection.ioloop.stop()
        else:
            logger.warning(f"Connection closed: {reason}")
            self.reconnect()

    def reconnect(self):
        """Reconnect with exponential backoff."""
        if self.closing:
            return

        self.reconnect_delay = min(self.reconnect_delay + 1, 30)
        logger.info(f"Reconnecting in {self.reconnect_delay} seconds")
        self.connection.ioloop.call_later(
            self.reconnect_delay,
            self._reconnect
        )

    def _reconnect(self):
        """Start new connection."""
        self.connection = self.connect()

    def on_channel_open(self, channel):
        """Called when channel is opened."""
        logger.info("Channel opened")
        self.channel = channel
        self.channel.add_on_close_callback(self.on_channel_closed)
        self.channel.queue_declare(
            queue=self.queue,
            durable=True,
            callback=self.on_queue_declared
        )

    def on_channel_closed(self, channel, reason):
        """Called when channel is closed."""
        logger.warning(f"Channel closed: {reason}")
        if not self.closing:
            self.connection.close()

    def on_queue_declared(self, frame):
        """Called when queue is declared."""
        logger.info(f"Queue declared: {self.queue}")
        self.channel.basic_qos(
            prefetch_count=10,
            callback=self.on_qos_set
        )

    def on_qos_set(self, frame):
        """Called when QoS is set."""
        logger.info("QoS set")
        self.consumer_tag = self.channel.basic_consume(
            queue=self.queue,
            on_message_callback=self.on_message
        )

    def on_message(self, channel, method, properties, body):
        """Process message."""
        logger.info(f"Received: {body.decode()}")
        channel.basic_ack(delivery_tag=method.delivery_tag)

    def run(self):
        """Run the consumer."""
        self.connection = self.connect()
        self.connection.ioloop.start()

    def stop(self):
        """Stop the consumer."""
        logger.info("Stopping")
        self.closing = True
        if self.channel:
            self.channel.basic_cancel(self.consumer_tag)
        self.connection.ioloop.stop()

# Usage
if __name__ == '__main__':
    consumer = AsyncReconnectingConsumer('localhost', 'async_queue')
    try:
        consumer.run()
    except KeyboardInterrupt:
        consumer.stop()

Handling Blocked Connections

import pika
import logging

logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)

def on_connection_blocked(connection, method):
    """Called when connection is blocked by broker."""
    logger.warning(f"Connection blocked: {method.reason}")

def on_connection_unblocked(connection, method):
    """Called when connection is unblocked."""
    logger.info("Connection unblocked")

# Set up connection with blocked callback
parameters = pika.ConnectionParameters(
    host='localhost',
    blocked_connection_timeout=300        # Timeout for blocked state
)

connection = pika.BlockingConnection(parameters)

# Add callbacks for blocked/unblocked events
connection.add_on_connection_blocked_callback(on_connection_blocked)
connection.add_on_connection_unblocked_callback(on_connection_unblocked)

channel = connection.channel()

Quick Reference

Connection Setup

# BlockingConnection
conn = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
channel = conn.channel()

# With credentials
creds = pika.PlainCredentials('user', 'pass')
params = pika.ConnectionParameters('host', credentials=creds)
conn = pika.BlockingConnection(params)

# From URL
params = pika.URLParameters('amqp://user:pass@host:5672/%2F')
conn = pika.BlockingConnection(params)

Publishing

# Simple publish
channel.basic_publish(exchange='', routing_key='queue', body='message')

# With properties
props = pika.BasicProperties(delivery_mode=2, content_type='application/json')
channel.basic_publish(exchange='', routing_key='queue', body='msg', properties=props)

# With confirmation
channel.confirm_delivery()
channel.basic_publish(exchange='', routing_key='queue', body='msg', mandatory=True)

Consuming

# Basic consume
channel.basic_consume(queue='queue', on_message_callback=callback, auto_ack=False)
channel.start_consuming()

# With QoS
channel.basic_qos(prefetch_count=10)

# Polling
method, props, body = channel.basic_get(queue='queue', auto_ack=False)

Queue Operations

# Declare
channel.queue_declare(queue='queue', durable=True)

# Bind
channel.queue_bind(exchange='ex', queue='queue', routing_key='key')

# Purge
channel.queue_purge(queue='queue')

# Delete
channel.queue_delete(queue='queue')

Acknowledgments

# Positive ack
channel.basic_ack(delivery_tag=method.delivery_tag)

# Multiple ack
channel.basic_ack(delivery_tag=method.delivery_tag, multiple=True)

# Negative ack with requeue
channel.basic_nack(delivery_tag=method.delivery_tag, requeue=True)

# Reject without requeue
channel.basic_reject(delivery_tag=method.delivery_tag, requeue=False)

Common Issues and Solutions

Connection Refused

Symptoms: AMQPConnectionError: Connection refused

Causes and solutions:

  1. RabbitMQ not running - Start the service

    sudo systemctl start rabbitmq-server
    rabbitmqctl status
    
  2. Wrong host/port - Verify connection parameters

    # Check parameters
    params = pika.ConnectionParameters(
        host='localhost',
        port=5672  # Default AMQP port
    )
    
  3. Firewall blocking - Check firewall rules

    sudo ufw allow 5672/tcp
    

Authentication Failures

Symptoms: ProbableAccessDeniedError: (403) ACCESS_REFUSED

Causes and solutions:

  1. Wrong credentials - Verify username/password

    rabbitmqctl list_users
    rabbitmqctl set_user_tags myuser administrator
    
  2. Missing permissions - Grant permissions

    rabbitmqctl set_permissions -p / myuser ".*" ".*" ".*"
    
  3. Wrong vhost - Check virtual host

    rabbitmqctl list_vhosts
    rabbitmqctl add_vhost myvhost
    

Channel Closed by Broker

Symptoms: ChannelClosedByBroker: (404) NOT_FOUND

Causes and solutions:

  1. Queue doesn't exist - Declare queue first or use passive=False

    # Will create if doesn't exist
    channel.queue_declare(queue='my_queue')
    
    # Check if exists (raises exception if not)
    channel.queue_declare(queue='my_queue', passive=True)
    
  2. Exchange doesn't exist - Declare exchange

    channel.exchange_declare(exchange='my_exchange', exchange_type='direct')
    
  3. Mismatched declaration - Use same parameters

    # Must match existing queue properties
    channel.queue_declare(queue='my_queue', durable=True)  # If originally durable
    

Messages Not Being Delivered

Symptoms: Published messages not appearing in queue

Causes and solutions:

  1. Wrong routing key - Verify bindings

    channel.queue_bind(
        exchange='my_exchange',
        queue='my_queue',
        routing_key='correct.key'  # Must match publish routing key
    )
    
  2. Exchange type mismatch - Check exchange type behaviour

    # Fanout ignores routing key
    # Direct requires exact match
    # Topic uses pattern matching
    
  3. Message expired - Check TTL settings

    # Message TTL
    props = pika.BasicProperties(expiration='60000')
    
    # Queue TTL
    channel.queue_declare(
        queue='ttl_queue',
        arguments={'x-message-ttl': 60000}
    )
    

Consumer Not Receiving Messages

Symptoms: Consumer connected but not receiving

Causes and solutions:

  1. Forgot to start consuming - Call start_consuming()

    channel.basic_consume(queue='queue', on_message_callback=callback)
    channel.start_consuming()  # Don't forget this!
    
  2. Auto-ack without processing - Use manual ack

    channel.basic_consume(
        queue='queue',
        on_message_callback=callback,
        auto_ack=False  # Manual ack required
    )
    
  3. Prefetch too low - Increase prefetch count

    channel.basic_qos(prefetch_count=10)  # Allow 10 unacked messages
    

Memory/Disk Alarms

Symptoms: Connection blocked, publishing fails

Causes and solutions:

  1. Low disk space - Free disk space or adjust watermark

    # Check status
    rabbitmqctl status | grep -A 5 "Alarms"
    
    # Adjust watermark
    rabbitmqctl set_disk_free_limit 1GB
    
  2. High memory usage - Adjust memory watermark

    rabbitmqctl set_vm_memory_high_watermark 0.6
    
  3. Too many messages - Purge or consume messages

    channel.queue_purge(queue='bloated_queue')
    

Messages Requeued Indefinitely

Symptoms: Same message processed repeatedly

Causes and solutions:

  1. Always requeuing on error - Use dead letter queue

    # Track retries
    retry_count = (properties.headers or {}).get('x-death', [{}])[0].get('count', 0)
    if retry_count > 3:
        channel.basic_nack(delivery_tag=method.delivery_tag, requeue=False)
    
  2. No dead letter queue - Configure DLQ

    channel.queue_declare(
        queue='main_queue',
        arguments={
            'x-dead-letter-exchange': 'dlx',
            'x-dead-letter-routing-key': 'dead'
        }
    )
    

Related Topics

The following topics complement Python pika knowledge and are commonly used together:

  • Message Queue Patterns - At-least-once delivery, exactly-once processing, dead letter queues; design patterns for reliable messaging
  • Python - redis + streams - Alternative messaging with Redis streams; pub/sub and consumer groups
  • Kafka - High-throughput distributed event streaming; comparison with RabbitMQ for different use cases
  • Python - FastAPI - Building HTTP APIs that integrate with message queues; background task processing
  • Docker - Containerising RabbitMQ and Python applications; Docker Compose for local development
  • Kubernetes - Deploying RabbitMQ clusters and consumers; scaling message-driven applications