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

如何在Django中连接RabbitMQ集群实现故障转移?

Django Dramatiq 配置RabbitMQ集群故障转移

django_dramatiq官方给出的默认配置仅支持单节点RabbitMQ连接:

DEFAULT_BROKER = "dramatiq.brokers.rabbitmq.RabbitmqBroker"
DEFAULT_BROKER_SETTINGS = {
    "BROKER": DEFAULT_BROKER,
    "OPTIONS": {
        "host": "127.0.0.1",
        "port": 5672,
        "heartbeat": 0,
        "connection_attempts": 5,
    },
    "MIDDLEWARE": [
        "dramatiq.middleware.Prometheus",
        "dramatiq.middleware.AgeLimit",
        "dramatiq.middleware.TimeLimit",
        "dramatiq.middleware.Callbacks",
        "dramatiq.middleware.Retries",
        "django_dramatiq.middleware.AdminMiddleware",
        "django_dramatiq.middleware.DbConnectionsMiddleware",
    ]
}

但要实现多节点RabbitMQ集群的故障转移,可利用RabbitmqBroker底层依赖的pika库的多参数连接特性,直接在Django配置中传入parameters列表替代单节点配置即可。

具体配置示例

import os
from pika import PlainCredentials, ConnectionParameters

DEFAULT_BROKER = "dramatiq.brokers.rabbitmq.RabbitmqBroker"
credentials = PlainCredentials(os.getenv('RMQ_USER'), os.getenv('RMQ_PASS'))

# 构造集群节点连接参数列表
rmq_parameters = [
    ConnectionParameters(
        host=os.getenv('RMQ_1'),
        port=int(os.getenv('RMQ_PORT_1')),
        credentials=credentials,
        virtual_host=os.getenv('VIRTUAL_HOST')
    ),
    ConnectionParameters(
        host=os.getenv('RMQ_2'),
        port=int(os.getenv('RMQ_PORT_2')),
        credentials=credentials,
        virtual_host=os.getenv('VIRTUAL_HOST')
    ),
    ConnectionParameters(
        host=os.getenv('RMQ_3'),
        port=int(os.getenv('RMQ_PORT_3')),
        credentials=credentials,
        virtual_host=os.getenv('VIRTUAL_HOST'),
        connection_attempts=5,
        retry_delay=1
    )
]

DEFAULT_BROKER_SETTINGS = {
    "BROKER": DEFAULT_BROKER,
    "OPTIONS": {
        "parameters": rmq_parameters,
        "heartbeat": 0,
    },
    "MIDDLEWARE": [
        "dramatiq.middleware.Prometheus",
        "dramatiq.middleware.AgeLimit",
        "dramatiq.middleware.TimeLimit",
        "dramatiq.middleware.Callbacks",
        "dramatiq.middleware.Retries",
        "django_dramatiq.middleware.AdminMiddleware",
        "django_dramatiq.middleware.DbConnectionsMiddleware",
    ]
}

关键说明

  • 在OPTIONS中传入parameters列表,替代原有的host和port参数,pika会自动按顺序尝试连接集群节点,实现故障转移
  • 每个ConnectionParameters可单独配置连接重试、延迟等参数,也可在全局OPTIONS中设置通用参数
  • 确保所有集群节点的虚拟主机、账号凭证一致,避免连接失败

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 20:40:16