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.
graph TB
subgraph Application
PUB[Publisher]
CON[Consumer]
end
subgraph RabbitMQ Broker
EX[Exchange]
Q1[Queue 1]
Q2[Queue 2]
Q3[Queue 3]
end
PUB -->|publish| EX
EX -->|route| Q1
EX -->|route| Q2
EX -->|route| Q3
Q1 -->|consume| CON
Q2 -->|consume| CON
Q3 -->|consume| CON
Connection and Channel Setup
Connections manage TCP connections to RabbitMQ, while channels are virtual connections that multiplex over a single TCP connection.
Key Concepts
flowchart LR
subgraph Application
CH1[Channel 1]
CH2[Channel 2]
CH3[Channel 3]
end
subgraph Connection
TCP[TCP Socket]
end
subgraph RabbitMQ
BROKER[Broker]
end
CH1 --> TCP
CH2 --> TCP
CH3 --> TCP
TCP --> BROKER
- 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
flowchart LR
PUB[Publisher] -->|routing_key| EX[Exchange]
EX -->|binding_key match| Q1[Queue A]
EX -->|binding_key match| Q2[Queue 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
stateDiagram-v2
[*] --> Ready
Ready --> Unacked : Deliver
Unacked --> Ready : Ack
Unacked --> Ready : Nack + Requeue
Unacked --> [*] : Nack + Discard
Unacked --> Ready : Reject + Requeue
- 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
flowchart TB
MSG[Message Delivered] --> PROCESS{Process Message}
PROCESS -->|Success| ACK[basic_ack]
PROCESS -->|Retry| NACK_R[basic_nack + requeue]
PROCESS -->|Discard| NACK_D[basic_nack + discard]
PROCESS -->|Reject Single| REJ[basic_reject]
ACK --> DONE[Message Removed]
NACK_R --> REQUEUE[Message Requeued]
NACK_D --> DLQ[Dead Letter Queue]
REJ --> DLQ
- 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
stateDiagram-v2
[*] --> Connected
Connected --> Disconnected : Connection Lost
Disconnected --> Reconnecting : Auto-reconnect
Reconnecting --> Connected : Success
Reconnecting --> Reconnecting : Retry
Reconnecting --> Failed : Max Retries
Failed --> [*]
- 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:
-
RabbitMQ not running - Start the service
sudo systemctl start rabbitmq-server rabbitmqctl status -
Wrong host/port - Verify connection parameters
# Check parameters params = pika.ConnectionParameters( host='localhost', port=5672 # Default AMQP port ) -
Firewall blocking - Check firewall rules
sudo ufw allow 5672/tcp
Authentication Failures
Symptoms: ProbableAccessDeniedError: (403) ACCESS_REFUSED
Causes and solutions:
-
Wrong credentials - Verify username/password
rabbitmqctl list_users rabbitmqctl set_user_tags myuser administrator -
Missing permissions - Grant permissions
rabbitmqctl set_permissions -p / myuser ".*" ".*" ".*" -
Wrong vhost - Check virtual host
rabbitmqctl list_vhosts rabbitmqctl add_vhost myvhost
Channel Closed by Broker
Symptoms: ChannelClosedByBroker: (404) NOT_FOUND
Causes and solutions:
-
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) -
Exchange doesn't exist - Declare exchange
channel.exchange_declare(exchange='my_exchange', exchange_type='direct') -
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:
-
Wrong routing key - Verify bindings
channel.queue_bind( exchange='my_exchange', queue='my_queue', routing_key='correct.key' # Must match publish routing key ) -
Exchange type mismatch - Check exchange type behaviour
# Fanout ignores routing key # Direct requires exact match # Topic uses pattern matching -
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:
-
Forgot to start consuming - Call start_consuming()
channel.basic_consume(queue='queue', on_message_callback=callback) channel.start_consuming() # Don't forget this! -
Auto-ack without processing - Use manual ack
channel.basic_consume( queue='queue', on_message_callback=callback, auto_ack=False # Manual ack required ) -
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:
-
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 -
High memory usage - Adjust memory watermark
rabbitmqctl set_vm_memory_high_watermark 0.6 -
Too many messages - Purge or consume messages
channel.queue_purge(queue='bloated_queue')
Messages Requeued Indefinitely
Symptoms: Same message processed repeatedly
Causes and solutions:
-
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) -
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