如何在Django中用Pika实现随服务启动的RabbitMQ消息监听器
Got it, let's work through setting up a RabbitMQ consumer with Pika that automatically starts when your Django server launches. I'll break this down into actionable steps with code examples:
1. Install Required Dependencies
First, make sure you have Pika installed:
pip install pika
2. Create a RabbitMQ Consumer Module
In your Django app (let's say it's named myapp), create a new file called consumers.py. This will hold all your consumer logic, including connection setup, message handling, and reconnection logic for robustness:
import pika import threading import time from django.conf import settings def rabbitmq_consumer(): # Pull RabbitMQ config from Django settings credentials = pika.PlainCredentials( settings.RABBITMQ_USER, settings.RABBITMQ_PASS ) connection_params = pika.ConnectionParameters( host=settings.RABBITMQ_HOST, port=settings.RABBITMQ_PORT, virtual_host=settings.RABBITMQ_VHOST, credentials=credentials, heartbeat=600, # Optional: Add heartbeat for persistent connections blocked_connection_timeout=300 ) def handle_message(ch, method, properties, body): # Replace this with your actual message processing logic try: message_content = body.decode('utf-8') print(f"Processed message: {message_content}") # Acknowledge the message to remove it from the queue ch.basic_ack(delivery_tag=method.delivery_tag) except Exception as e: print(f"Failed to process message: {str(e)}") # Optionally reject the message (requeue or discard) ch.basic_nack(delivery_tag=method.delivery_tag, requeue=False) try: # Establish connection and channel connection = pika.BlockingConnection(connection_params) channel = connection.channel() # Declare a durable queue (survives RabbitMQ restarts) channel.queue_declare(queue='your_target_queue', durable=True) # Limit unacknowledged messages to prevent overwhelming the consumer channel.basic_qos(prefetch_count=1) # Start consuming messages channel.basic_consume( queue='your_target_queue', on_message_callback=handle_message ) print("✅ RabbitMQ consumer initialized. Waiting for messages...") channel.start_consuming() except pika.exceptions.ConnectionClosedByBroker: print("🔌 Connection closed by broker. Restarting consumer...") rabbitmq_consumer() except pika.exceptions.AMQPChannelError as e: print(f"⚠️ Channel error occurred: {str(e)}. Restarting consumer...") rabbitmq_consumer() except pika.exceptions.AMQPConnectionError as e: print(f"⚠️ Connection failed: {str(e)}. Retrying in 5 seconds...") time.sleep(5) rabbitmq_consumer()
3. Add RabbitMQ Config to Django Settings
Update your settings.py to include RabbitMQ connection details (adjust these to match your actual RabbitMQ setup):
# RabbitMQ Configuration RABBITMQ_HOST = 'localhost' # Replace with your RabbitMQ host RABBITMQ_PORT = 5672 RABBITMQ_USER = 'guest' # Replace with your username RABBITMQ_PASS = 'guest' # Replace with your password RABBITMQ_VHOST = '/'
4. Auto-Start Consumer on Django Server Boot
To trigger the consumer when Django starts, we'll use Django's AppConfig system. Update your app's apps.py file:
from django.apps import AppConfig import threading class MyAppConfig(AppConfig): default_auto_field = 'django.db.models.BigAutoField' name = 'myapp' def ready(self): # Start the consumer in a daemon thread so it exits when Django stops from .consumers import rabbitmq_consumer consumer_thread = threading.Thread(target=rabbitmq_consumer, daemon=True) consumer_thread.start()
Then, make sure Django uses this custom config by updating INSTALLED_APPS in settings.py:
INSTALLED_APPS = [ # ... other installed apps 'myapp.apps.MyAppConfig', # Replace 'myapp' with this line ]
Key Notes
- Daemon Thread: Using
daemon=Trueensures the consumer thread shuts down automatically when the Django server stops, preventing orphaned processes. - Message Acknowledgment: The
basic_ackcall tells RabbitMQ the message was processed successfully. If you skip this, RabbitMQ will re-deliver the message when the consumer disconnects. - Reconnection Logic: The consumer will automatically retry connections if RabbitMQ goes down, which is critical for production environments.
- Queue Durability: The
durable=Trueflag ensures the queue survives RabbitMQ restarts (make sure your producers also publish to a durable queue).
内容的提问来源于stack exchange,提问作者Souradipta Roy

