Кластеры
Кластеры
Один алиас соединения может указывать больше чем на один узел брокера. Django-RMQ предоставляет клиентский failover: pika перебирает настроенные узлы по порядку, пока один из них не примет соединение. При каждой попытке подключения библиотека передаёт pika последовательность адресов узлов, и pika сама проходит по этой последовательности.
Поддерживаются две топологии:
- Список нод — алиас перечисляет все ноды кластера через
NODES, аpikaсама выбирает, к какой из них подключиться. - Единая точка входа перед кластером — балансировщик нагрузки или DNS-имя с несколькими A-записями — где алиас по-прежнему использует обычные
HOST/PORT, а распределение и failover происходят на стороне сервера (то есть, например, балансировщика).
Клиентский список узлов (NODES)
Настройте NODES — по одной записи на ноду кластера:
# settings.py
RABBITMQ_CONNECTIONS: dict = {
'default': {
'NODES': [
{'HOST': 'rmq-1.internal', 'PORT': 5672},
{'HOST': 'rmq-2.internal', 'PORT': 5672},
{'HOST': 'rmq-3.internal', 'PORT': 5672},
],
'VIRTUAL_HOST': '/',
'USER': 'guest',
'PASSWORD': 'guest',
'HEARTBEAT': 600,
'BLOCKED_CONNECTION_TIMEOUT': 300,
'RECONNECT_INITIAL_BACKOFF': 1.0,
'RECONNECT_MAX_BACKOFF': 30.0,
},
}_resolve_nodesвdjango_rmq/apps.pyпревращает каждую запись{'HOST': ..., 'PORT': ...}вNodeConfig;- Остальные ключи (
VIRTUAL_HOST,USER,PASSWORD,HEARTBEAT,BLOCKED_CONNECTION_TIMEOUT) применяются ко всем узлам списка.
Далее RabbitMQConnectionManager.__init__ в django_rmq/connections.py строит по одному pika.ConnectionParameters на каждый узел (_node_parameters), все они разделяют одни и те же PlainCredentials, virtual_host, heartbeat и blocked_connection_timeout, и передаёт весь список в BlockingConnection:
# django_rmq/connections.py
setattr(self._local, attr, BlockingConnection(parameters=sequence))BlockingConnection в pika принимает последовательность объектов Parameters и пробует законнектится к каждому по очереди, пока один из них не подключится успешно — именно это и обеспечивает failover по кластеру.
NODES взаимоисключающий со скалярной формой HOST/PORT. См. для полного справочника параметров и случаев ошибок валидации, возникающих, если указаны оба варианта или ни одного.
Как работает failover при переподключении
Последовательность нод строится заново при каждой попытке подключения — не только при первой. RabbitMQConnectionManager._get_or_create_connection вызывает _build_connection_sequence() каждый раз, когда нужно открыть соединение, поэтому переподключение получает свежую копию списка нод, а не переиспользует устаревшую.
Это напрямую связано с уже существующей : когда производитель сталкивается с reconnectable-ошибкой, он вызывает reset_producer_channel() и переоткрывает соединение (повторяя публикацию один раз); потребитель переподключается с экспоненциальным backoff. Поскольку переоткрытие соединения заново прогоняет последовательность узлов, pika автоматически подключается к тому узлу, который сейчас жив — никакого кластерно-осведомлённого кода в приложении не требуется.
Типичный сценарий failover, шаг за шагом:
- Нода
Aвыходит из строя. - Активное соединение с
Aобрывается. - Происходит попытка (повтор продюсера или backoff-цикл потребителя) переоткрыть соединение.
pikaснова перебирает[A, B, C]:Aотказывает в соединении,Bже принимает его.- Публикация/потребление возобновляются на ноде
B.
get_producer_connection и get_consumer_connection логируют debug-запись при каждом открытии нового соединения, включая список узлов и флаг shuffle:
# django_rmq/connections.py
logger.debug(
{
'source': source,
'message': message,
'data': {
'nodes': [{'host': params.host, 'port': params.port} for params in sequence],
'shuffle': self._shuffle_nodes,
},
}
)Распределение клиентов по кластеру (SHUFFLE_NODES)
SHUFFLE_NODES по умолчанию False: каждый клиент проходит список NODES в том порядке, в котором он был объявлен, поэтому все клиенты выбирают ноду №1 первой. При большом числе клиентских процессов это смещает нагрузку на эту ноду.
Установите SHUFFLE_NODES: True, чтобы перемешивать последовательность при каждой попытке подключения:
# django_rmq/connections.py — RabbitMQConnectionManager._build_connection_sequence
sequence: list[ConnectionParameters] = list(self._node_parameters)
if self._shuffle_nodes:
random.shuffle(sequence)
return sequenceshuffle_nodes берётся из RabbitMQConfig.shuffle_nodes. Перемешивается только копия, используемая для конкретной попытки подключения — настроенный в settings порядок NODES никогда не изменяется, меняется лишь порядок перебора для этого конкретного соединения.
Включайте SHUFFLE_NODES для кластеров с большим числом клиентских процессов, чтобы соединения распределялись по нодам, а не концентрировались на первой.
Альтернатива: балансировщик или DNS
Вместо перечисления всех узлов через NODES можно поставить перед кластером единый адрес: load-balancer (например, HAProxy) или DNS-имя с несколькими A-записями, указывающими на узлы кластера. В этом случае достаточно обычных HOST/PORT — pika резолвит адрес и подключается к тому, который отвечает, а распределение и failover происходят на стороне сервера (то есть балансировщика).
Trade-offs:
NODES(клиентская сторона) — не требует дополнительной инфраструктуры; клиент знает обо всех нодах; логика failover находится в приложении (черезpika).- Балансировщик / DNS (серверная сторона) — единая точка входа для настройки; health-check'и централизованы на стороне LB/DNS; требуется разворачивать и поддерживать эту инфраструктуру; конфигурация клиента остаётся простой.
Эти два подхода не исключают друг друга на уровне протокола, но для конкретного алиаса обычно выбирают один из них.
Quorum-очереди
Кворумные очереди ортогональны адресации кластера — django_rmq объявляет их так же, как и любой другой тип очереди, через поле queue_type на QueueConfig (см. в справочнике API):
from django_rmq.queues.queue_config import QueueConfig, QueueType
orders_queue: QueueConfig = QueueConfig(name='orders', queue_type=QueueType.QUORUM)Объявление очереди с этим конфигом ставит x-queue-type: quorum при queue_declare. Определите orders_queue в одном месте (например, в myapp/queues.py) и переиспользуйте её в продюсере, потребителе и в любой setup-функции, которая ссылается на эту очередь — см. .
Либо объявляйте кворумную очередь из raw setup-функции, зарегистрированной через setup-registry (см. ) — это низкоуровневая альтернатива, полезная, если очереди нужны аргументы, которых нет в QueueConfig:
from pika.adapters.blocking_connection import BlockingChannel
def setup_quorum(channel: BlockingChannel) -> None:
channel.queue_declare(
queue='orders',
durable=True,
arguments={'x-queue-type': 'quorum'},
)Зарегистрируйте её так же, как и любую другую setup-функцию:
from django_rmq.registries.setup_registry import get_setup_registry
get_setup_registry().register(fn=setup_quorum)Publisher confirms уже включены на каждом канале продюсера (confirm_delivery() в django_rmq/connections.py — см. ), что хорошо сочетается с кворумными очередями для обеспечения durability в рамках кластера 😃
См. также
- — полный справочник параметров
NODES/SHUFFLE_NODESи случаи ошибок валидации. - — единичный повтор продюсера и переподключение потребителя с backoff.
- — работа с несколькими алиасами.
- — объявление exchanges, очередей и bindings через setup-функции.
- — интеграционные тесты кластера, проверяющие это поведение failover на реальном кластере RabbitMQ с тремя нодами.