如何在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
相关产品推荐
相关产品推荐

