Перейти к основному содержимому
Версия: 1.0.5

Справочник API

Плоский справочник по всем публичным символам в django_rmq. Сигнатуры и описания параметров взяты непосредственно из исходного кода.


Producer​

from django_rmq.producer import Producer

Публикует сообщения в RabbitMQ через потоко-локальный блокирующий канал.

Producer.__init__​

def __init__(
self,
exchange: str = '',
queue: QueueConfig | str = '',
using: str | None = None,
) -> None
ПараметрТипПо умолчаниюОписание
exchangestr''Имя обменника. Пустая строка означает обменник по умолчанию (direct).
queueQueueConfig | str''Конфигурация очереди или её имя. Используется как ключ маршрутизации по умолчанию; объявляется лениво при первой публикации. Пустая строка отключает объявление очереди (режим только-обменник).
usingstr | NoneNoneПсевдоним соединения из RABBITMQ_CONNECTIONS. Можно не указывать, если настроен ровно один псевдоним.

Producer.publish​

def publish(
self,
body: str | bytes,
routing_key: str | None = None,
properties: BasicProperties | None = None,
) -> None
ПараметрТипПо умолчаниюОписание
bodystr | bytesобязательныйТело сообщения. Строки кодируются в байты (UTF-8).
routing_keystr | NoneNoneКлюч маршрутизации. По умолчанию используется self.queue (имя очереди), если не указан.
propertiesBasicProperties | NoneNoneСвойства AMQP-сообщения. Если не указаны, создаются с content_type='application/json'. delivery_mode всегда принудительно устанавливается в 2 (persistence).

Returns: None

Raises:

  • pika.exceptions.UnroutableError — брокер вернул сообщение, так как ни одна очередь не совпала с ключом маршрутизации (требуются подтверждения публикации + mandatory=True, оба всегда включены).
  • pika.exceptions.NackError — брокер отклонил (nack) сообщение.
  • Любое исключение из _RECONNECTABLE_ERRORS, если и первая попытка, и единственный повтор завершились ошибкой.

Producer.__call__​

Позволяет использовать экземпляр Producer как декоратор функции. Декорируемая функция должна возвращать str, bytes или None. Возвращаемое значение публикуется автоматически; None пропускает публикацию.

def __call__(
self,
func: Callable[..., str | bytes | None],
) -> Callable[..., str | bytes | None]

Raises: TypeError, если декорируемая функция возвращает тип, отличный от str, bytes или None.

Пример:

from django_rmq.producer import Producer

producer: Producer = Producer(queue='notifications')


@producer
def build_notification(user_id: int) -> str:
return f'{{"user_id": {user_id}}}'

Consumer​

from django_rmq.consumer import Consumer

Потребляет сообщения из одной очереди RabbitMQ с одним зарегистрированным обработчиком. При транзитных ошибках AMQP переподключается с экспоненциальной задержкой.

Consumer.__init__​

def __init__(
self,
queue: QueueConfig | str,
prefetch_count: int = 1,
reconnect_initial_backoff: float | None = None,
reconnect_max_backoff: float | None = None,
using: str | None = None,
) -> None
ПараметрТипПо умолчаниюОписание
queueQueueConfig | strобязательныйОчередь для потребления. QueueConfig инициирует активное объявление с аргументами; обычная строка инициирует durable-объявление.
prefetch_countint1Максимальное количество неподтверждённых сообщений, доставляемых одновременно (basic_qos).
reconnect_initial_backofffloat | NoneNoneНачальная задержка перед переподключением в секундах. При None берётся из конфигурации псевдонима.
reconnect_max_backofffloat | NoneNoneМаксимальная задержка перед переподключением в секундах (ограничение экспоненциальной задержки). При None берётся из конфигурации псевдонима.
usingstr | NoneNoneПсевдоним соединения. Можно не указывать, если настроен ровно один псевдоним.

Consumer.handler​

def handler(self, func: MessageCallback) -> MessageCallback

Регистрирует коллбэк для входящих сообщений. Может использоваться как декоратор.

Raises: RuntimeError, если на этом консьюмере уже зарегистрирован обработчик.

Consumer.__call__​

Сокращение для @consumer.handler. Идентичное поведение.

def __call__(self, func: MessageCallback) -> MessageCallback

Consumer.consume​

def consume(self, stop_event: threading.Event | None = None) -> None

Запускает цикл потребления. При транзитных ошибках переподключается. Завершается при установке stop_event или при возникновении неустранимой ошибки.

ПараметрТипПо умолчаниюОписание
stop_eventthreading.Event | NoneNoneСобытие для сигнала о плановом завершении. При None создаётся внутреннее событие (никогда не устанавливается) — консьюмер работает до ошибки.

Свойства Consumer​

СвойствоТипОписание
prefetch_countintМаксимальное количество неподтверждённых сообщений, доставляемых одновременно.
usingstr | NoneПсевдоним соединения или None, если единственный псевдоним используется неявно.
handler_namestrИмя зарегистрированной функции-обработчика или 'unregistered', если обработчик не зарегистрирован.

MessageCallback​

from django_rmq.consumer import MessageCallback

Псевдоним типа для сигнатуры вызываемого обработчика:

MessageCallback = Callable[
[BlockingChannel, Basic.Deliver, BasicProperties, bytes],
None,
]
АргументТипОписание
chBlockingChannelКанал, по которому доставлено сообщение. Используется для вызова basic_ack / basic_nack.
methodBasic.DeliverМетаданные доставки, включая delivery_tag.
propsBasicPropertiesСвойства AMQP-сообщения.
bodybytesНеобработанное тело сообщения.

QueueConfig​

from django_rmq.queues.queue_config import QueueConfig, QueueType

frozen dataclass для декларативной конфигурации очереди.

@dataclass(frozen=True)
class QueueConfig:
name: str
durable: bool = True
queue_type: QueueType | None = None
dead_letter_exchange: str | None = None
dead_letter_routing_key: str | None = None
ПолеТипПо умолчаниюОписание
namestrобязательноеИмя очереди. Также используется как str(queue_config).
durableboolTrueОчередь переживает перезапуск брокера.
queue_typeQueueType | NoneNoneУстанавливает x-queue-type при объявлении; при None используется default_queue_type брокера.
dead_letter_exchangestr | NoneNoneОбменник, куда маршрутизируются dead-letter сообщения (x-dead-letter-exchange).
dead_letter_routing_keystr | NoneNoneКлюч маршрутизации для dead-letter сообщений (x-dead-letter-routing-key).

QueueType​

Enum поддерживаемых типов очередей RabbitMQ; значение члена сериализуется брокеру как аргумент x-queue-type.

ЧленЗначение
CLASSIC'classic'
QUORUM'quorum'
STREAM'stream'

QueueType наследуется от str (а не от StrEnum, доступного только в 3.11+, чтобы сохранить совместимость с Python 3.10), поэтому каждый член — обычная строка и сериализуется брокеру как есть. Если queue_type не задан на QueueConfig, брокер применяет свой default_queue_type.

Свойство QueueConfig.arguments​

@property
def arguments(self) -> dict[str, Any] | None

Строит словарь AMQP arguments для объявления очереди на основе queue_type (x-queue-type) и полей dead-letter. Возвращает None, если ничего из этого не задано.


RabbitMQConfig​

from django_rmq.dto.rabbitmq_config import RabbitMQConfig

frozen dataclass с разрешённой конфигурацией для одного псевдонима соединения. Создаётся внутри RabbitMQAppConfig.ready() из настройки RABBITMQ_CONNECTIONS. Доступен через RabbitMQConnectionManager.config.

@dataclass(frozen=True)
class RabbitMQConfig:
host: str
port: int
virtual_host: str
user: str
password: str
heartbeat: int
blocked_connection_timeout: int
reconnect_initial_backoff: float
reconnect_max_backoff: float
ПолеТипОписание
hoststrХостнейм или IP брокера.
portintAMQP-порт.
virtual_hoststrВиртуальный хост для подключения.
userstrИмя пользователя для PlainCredentials.
passwordstrПароль для PlainCredentials.
heartbeatintИнтервал heartbeat в секундах.
blocked_connection_timeoutintСекунды ожидания, пока соединение заблокировано брокером.
reconnect_initial_backofffloatНачальная задержка переподключения консьюмера в секундах.
reconnect_max_backofffloatМаксимальная задержка переподключения консьюмера (ограничение экспоненциальной задержки).

RabbitMQConnectionManager​

from django_rmq.connections import RabbitMQConnectionManager

Управляет потоко-локальными соединениями для ролей продюсера и консьюмера. Создаётся по одному экземпляру на псевдоним в RabbitMQAppConfig.ready().

RabbitMQConnectionManager.__init__​

def __init__(self, config: RabbitMQConfig) -> None

RabbitMQConnectionManager.get_producer_connection​

def get_producer_connection(self) -> BlockingConnection

Возвращает потоко-локальное BlockingConnection, используемое продюсерами в текущем потоке. Создаётся при первом обращении; повторно использует кешированный экземпляр, пока соединение открыто.

RabbitMQConnectionManager.get_consumer_connection​

def get_consumer_connection(self) -> BlockingConnection

Возвращает потоко-локальное BlockingConnection, используемое консьюмерами в текущем потоке. Хранится отдельно от соединения продюсера, чтобы публикация из обработчика была безопасной.

RabbitMQConnectionManager.get_producer_channel​

def get_producer_channel(self) -> BlockingChannel

Возвращает потоко-локальный канал продюсера с включёнными подтверждениями публикации (confirm_delivery()). Создаётся при первом обращении; повторно используется, пока открыт.

RabbitMQConnectionManager.reset_producer_channel​

def reset_producer_channel(self) -> None

Закрывает и удаляет кешированный канал продюсера и соединение. Вызывается автоматически Producer.publish после восстанавливаемой ошибки, чтобы следующая публикация открыла новый канал. Безопасен к вызову даже при отсутствии кешированного канала.


get_connection_manager​

from django_rmq.connections import get_connection_manager
def get_connection_manager(using: str | None = None) -> RabbitMQConnectionManager

Возвращает RabbitMQConnectionManager для указанного псевдонима.

ПараметрТипПо умолчаниюОписание
usingstr | NoneNoneПсевдоним для разрешения. Можно не указывать, если настроен ровно один псевдоним.

Raises: ImproperlyConfigured — псевдоним не найден, неоднозначность (несколько псевдонимов, using не указан) или django_rmq не инициализирован.


ConsumersRegistry​

from django_rmq.registries.registry import ConsumersRegistry

Содержит консьюмеры, зарегистрированные для одного псевдонима соединения.

ConsumersRegistry.register​

def register(self, consumer: Consumer) -> None

Добавляет консьюмер в реестр.

ConsumersRegistry.all​

def all(self) -> list[Consumer]

Возвращает копию всех зарегистрированных консьюмеров. Изменение возвращённого списка не влияет на реестр.


get_consumers_registry​

from django_rmq.registries.registry import get_consumers_registry
def get_consumers_registry(using: str | None = None) -> ConsumersRegistry

Возвращает ConsumersRegistry для указанного псевдонима.

Raises: ImproperlyConfigured — те же условия, что и в get_connection_manager.


SetupRegistry​

from django_rmq.registries.setup_registry import SetupRegistry

Содержит идемпотентные функции настройки топологии для одного псевдонима соединения.

SetupRegistry.register​

def register(self, fn: SetupFn) -> None

Добавляет функцию настройки в реестр. Функции вызываются в порядке регистрации.

SetupRegistry.run_all​

def run_all(self, channel: BlockingChannel) -> None

Выполняет все зарегистрированные функции настройки на указанном канале в порядке регистрации.


get_setup_registry​

from django_rmq.registries.setup_registry import get_setup_registry
def get_setup_registry(using: str | None = None) -> SetupRegistry

Возвращает SetupRegistry для указанного псевдонима.

Raises: ImproperlyConfigured — те же условия, что и в get_connection_manager.


SetupFn​

from django_rmq.registries.setup_registry import SetupFn

Псевдоним типа для вызываемого объекта настройки топологии:

SetupFn = Callable[[BlockingChannel], None]

SetupFn получает открытый BlockingChannel и должна быть идемпотентной (безопасной для многократного вызова с одинаковым результатом).


RabbitMQAppConfig​

from django_rmq.apps import RabbitMQAppConfig

Django AppConfig для django_rmq. Регистрируется автоматически при наличии 'django_rmq' в INSTALLED_APPS.

RabbitMQAppConfig.ready​

def ready(self) -> None

Читает RABBITMQ_CONNECTIONS из настроек Django и создаёт для каждого псевдонима один RabbitMQConnectionManager, один SetupRegistry и один ConsumersRegistry. Сохраняет их в модуле django_rmq под именами connection_managers, setup_registries и consumers_registries.

Raises: ImproperlyConfigured — RABBITMQ_CONNECTIONS отсутствует или пуст.


Глобальные переменные модуля (django_rmq)​

import django_rmq

django_rmq.connection_managers # dict[str, RabbitMQConnectionManager] | None
django_rmq.setup_registries # dict[str, SetupRegistry] | None
django_rmq.consumers_registries # dict[str, ConsumersRegistry] | None

Все три имеют значение None до выполнения RabbitMQAppConfig.ready(). Используйте функции-аксессоры ( get_connection_manager, get_consumers_registry, get_setup_registry) вместо прямого обращения к этим словарям.


Команды управления​

setup_rabbitmq_topology​

uv run python manage.py setup_rabbitmq_topology [--using ALIAS]

Объявляет все обменники, очереди и привязки, зарегистрированные в SetupRegistry, для одного или всех псевдонимов. Идемпотентна — безопасна для запуска при каждом деплое.

ОпцияОписание
--using ALIASВыполнить только для указанного псевдонима. Без этой опции выполняется для всех псевдонимов.

После настройки выводит отчёт об объявленных обменниках, очередях и привязках, дедуплицированных по имени.

start_consumers​

uv run python manage.py start_consumers [--using ALIAS]

Запускает все консьюмеры, зарегистрированные в ConsumersRegistry, для одного или всех псевдонимов. Каждый консьюмер работает в отдельном потоке, разделяя единый stop_event. SIGTERM и SIGINT устанавливают это событие для планового завершения; команда ожидает завершения всех потоков перед возвратом.

ОпцияОписание
--using ALIASЗапустить консьюмеры только для указанного псевдонима. Без этой опции — для всех псевдонимов.

Справочник настроек​

КлючТипОписание
HOSTstrХостнейм или IP брокера.
PORTintAMQP-порт (обычно 5672).
VIRTUAL_HOSTstrВиртуальный хост (например, '/').
USERstrИмя пользователя для аутентификации.
PASSWORDstrПароль для аутентификации.
HEARTBEATintИнтервал heartbeat в секундах.
BLOCKED_CONNECTION_TIMEOUTintСекунды до таймаута заблокированного соединения.
RECONNECT_INITIAL_BACKOFFfloatНачальная задержка переподключения консьюмера в секундах.
RECONNECT_MAX_BACKOFFfloatМаксимальная задержка переподключения консьюмера (ограничение задержки).

Все девять ключей обязательны для каждого псевдонима. Отсутствующие ключи вызывают KeyError в процессе AppConfig.ready(). Полный пример настроек см. в разделе Конфигурация.