You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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=True ensures the consumer thread shuts down automatically when the Django server stops, preventing orphaned processes.
  • Message Acknowledgment: The basic_ack call 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=True flag ensures the queue survives RabbitMQ restarts (make sure your producers also publish to a durable queue).

内容的提问来源于stack exchange,提问作者Souradipta Roy

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.29 09:02:24