Spec-Zone.ru › Celery

celery.backends.rpc

Бэкенд результатов RPC для брокеров AMQP.

Бэкенд результатов в стиле RPC, использующий reply-to и отдельную очередь для каждого клиента.

exceptioncelery.backends.rpc.BacklogLimitExceeded

Слишком большой объем истории состояния для быстрой перемотки.

classcelery.backends.rpc.RPCBackend(app, connection=None, exchange=None, exchange_type=None, persistent=None, serializer=None, auto_delete=True, **kwargs)

Базовый класс для бэкенда результатов RPC.

exceptionBacklogLimitExceeded

Исключение, возникающее, когда для идентификатора задачи накоплено слишком много сообщений.

classConsumer(channel, queues=None, no_ack=None, auto_declare=None, callbacks=None, on_decode_error=None, on_message=None, accept=None, prefetch_count=None, tag_prefix=None)

Потребитель, для которого требуется вручную объявлять очереди.

auto_declare=False

По умолчанию все сущности объявляются при создании. Если вы хотите управлять этим вручную, задайте значение False.

classExchange(name='', type='', channel=None, **kwargs)

Объявление обменника.

Аргументы:

name (str): см. name. type (str): см. type. channel (kombu.Connection, ChannelT): см. channel. durable (bool): см. durable. auto_delete (bool): см. auto_delete. delivery_mode (enum): см. delivery_mode. arguments (Dict): см. arguments. no_declare (bool): см. no_declare

name(str)

По умолчанию имя не задано (используется обменник по умолчанию).

Тип:

Имя обменника.

type(str)

Это описание типов обменников AMQP беззастенчиво заимствовано из записи в блоге «AMQP за 10 минут: часть 4»_ Раджита Аттапатту. Если вы только начинаете знакомиться с AMQP, рекомендуем прочитать эту статью.

«AMQP определяет четыре типа обменников по умолчанию (алгоритма маршрутизации), охватывающих большинство распространённых сценариев обмена сообщениями. Брокер AMQP может также определять дополнительные типы обменников, поэтому дополнительную информацию о доступных типах см. в документации к своему брокеру.

  • direct (по умолчанию)

    Точное совпадение ключа маршрутизации в сообщении с критерием маршрутизации, заданным при привязке очереди к этому обменнику.

  • topic

    Сопоставление с подстановочными символами ключа маршрутизации и шаблона маршрутизации, указанного в привязке обменника к очереди. Ключ маршрутизации рассматривается как последовательность из нуля или более слов, разделённых символом «.», и поддерживает специальные подстановочные символы. «*» соответствует одному слову, а «#» — нулю или более слов.

  • fanout

    Очереди привязываются к этому обменнику без аргументов. Поэтому любое сообщение, отправленное в этот обменник, пересылается во все привязанные к нему очереди.

  • headers

    Очереди привязываются к этому обменнику с таблицей аргументов, содержащей заголовки и значения (необязательно). Специальный аргумент с именем «x-match» задаёт алгоритм сопоставления: «all» означает AND (должны совпасть все пары), а «any» — OR (должна совпасть хотя бы одна пара).

    Для указания аргументов используется arguments.

channel (ChannelT): канал, к которому привязан обменник (если он привязан).

durable (bool): устойчивые обменники остаются активными после перезапуска сервера.

Неустойчивые обменники (временные обменники) удаляются при перезапуске сервера. Значение по умолчанию — True.

auto_delete (bool): если параметр задан, обменник удаляется, когда все очереди

перестают его использовать. Значение по умолчанию — False.

delivery_mode (enum): режим доставки сообщений по умолчанию.

Значение — целое число или строковый псевдоним.

  • 1 или «transient»

    Сообщение является временным. Это означает, что оно хранится только в памяти и будет потеряно, если сервер выйдет из строя или перезапустится.

  • 2 или «persistent» (по умолчанию)

    Сообщение является постоянным. Это означает, что оно хранится как в памяти, так и на диске, поэтому сохраняется, если сервер выйдет из строя или перезапустится.

Значение по умолчанию — 2 (persistent).

arguments (Dict): дополнительные аргументы, задаваемые при объявлении

обменника.

no_declare (bool): никогда не объявлять этот обменник

(declare() ничего не делает).

Message(body, delivery_mode=None, properties=None, **kwargs)

Создать экземпляр сообщения для отправки с помощью publish().

Аргументы:

body (Any): тело сообщения.

delivery_mode (bool): задать пользовательский режим доставки.

По умолчанию используется delivery_mode.

priority (int): приоритет сообщения от 0 до настроенного брокером

максимального приоритета; чем выше значение, тем выше приоритет.

content_type (str): content_type сообщения. Если content_type

задан, сериализация не выполняется, поскольку предполагается, что это либо двоичный объект, либо вы сериализовали данные самостоятельно. Оставьте поле пустым, если используете встроенную сериализацию: наша библиотека правильно задаёт content_type.

content_encoding (str): кодировка символов, в которой

закодирован этот объект. Используйте «binary» при отправке необработанных двоичных объектов. Оставьте поле пустым, если используете встроенную сериализацию: наша библиотека правильно задаёт content_encoding.

properties (Dict): свойства сообщения.

headers (Dict): заголовки сообщения.

PERSISTENT_DELIVERY_MODE=2
TRANSIENT_DELIVERY_MODE=1
attrs:tuple[tuple[str,Any],...]=(('name', None), ('type', None), ('arguments', None), ('durable', <class 'bool'>), ('passive', <class 'bool'>), ('auto_delete', <class 'bool'>), ('delivery_mode', <function Exchange.<lambda>>), ('no_declare', <class 'bool'>))
auto_delete=False
bind_to(exchange='', routing_key='', arguments=None, nowait=False, channel=None, **kwargs)

Привязать обменник к другому обменнику.

Аргументы:

nowait (bool): если задано, сервер не будет отвечать, а вызов

не будет ожидать ответа. Значение по умолчанию — False.

binding(routing_key='', arguments=None, unbind_arguments=None)
propertycan_cache_declaration

bool(x) -> bool

Возвращает True, если аргумент x истинен, и False в противном случае. Встроенные значения True и False — единственные два экземпляра класса bool. Класс bool является подклассом класса int и не может быть подклассирован.

declare(nowait=False, passive=None, channel=None)

Объявить обменник.

Создаёт обменник на брокере, если только задан параметр passive; в этом случае проверяется лишь наличие обменника.

Аргумент:
nowait (bool): если задано, сервер не будет отвечать, и

ожидание ответа не выполняется. Значение по умолчанию — False.

delete(if_unused=False, nowait=False)

Удалить объявление обменника на сервере.

Аргументы:

if_unused (bool): удалять только в том случае, если у обменника нет привязок.

Значение по умолчанию — False.

nowait (bool): если задано, сервер не будет отвечать, и

ожидание ответа не выполняется. Значение по умолчанию — False.

delivery_mode=None
durable=True
name=''
no_declare=False
passive=False
publish(message, routing_key=None, mandatory=False, immediate=False, exchange=None)

Опубликовать сообщение.

Аргументы:

message (Union[kombu.Message, str, bytes]):

Сообщение для публикации.

routing_key (str): ключ маршрутизации сообщения. mandatory (bool): в настоящее время не поддерживается. immediate (bool): в настоящее время не поддерживается.

type='direct'
unbind_from(source='', routing_key='', nowait=False, arguments=None, channel=None)

Удалить на сервере ранее созданную привязку обменника.

classProducer(channel, exchange=None, routing_key=None, serializer=None, auto_declare=None, compression=None, on_return=None)

Издатель сообщений.

Аргументы:

channel (kombu.Connection, ChannelT): подключение или канал. exchange (kombu.entity.Exchange, str): необязательное значение биржи по умолчанию. routing_key (str): необязательный ключ маршрутизации по умолчанию. serializer (str): сериализатор по умолчанию. По умолчанию — “json”. compression (str): метод сжатия по умолчанию.

По умолчанию сжатие отключено.

auto_declare (bool): автоматически объявлять биржу по умолчанию

при создании экземпляра. По умолчанию — True.

on_return (Callable): функция обратного вызова для недоставленных сообщений,

если используются аргументы mandatory или immediate метода publish(). Эта функция обратного вызова должна иметь следующую сигнатуру: (exception, exchange, routing_key, message). Обратите внимание, что для использования этой возможности издателю необходимо обрабатывать события.

auto_declare=True

По умолчанию, если задана биржа по умолчанию, она объявляется при публикации сообщения.

propertychannel
close()
compression=None

Метод сжатия по умолчанию. По умолчанию отключено.

propertyconnection
declare()

Объявить биржу.

Примечание:

Это происходит автоматически при создании экземпляра, если включён флаг auto_declare.

exchange=None

Биржа по умолчанию

maybe_declare(entity, retry=False, **retry_policy)

Объявить биржу, если она ещё не была объявлена в этой сессии.

on_return=None

Базовая функция обратного вызова при возврате.

publish(body, routing_key=None, delivery_mode=None, mandatory=False, immediate=False, priority=0, content_type=None, content_encoding=None, serializer=None, headers=None, compression=None, exchange=None, retry=False, retry_policy=None, declare=None, expiration=None, timeout=None, confirm_timeout=None, **properties)

Опубликовать сообщение в указанную биржу.

Аргументы:

body (Any): тело сообщения. routing_key (str): ключ маршрутизации сообщения. delivery_mode (enum): см. delivery_mode. mandatory (bool): в настоящее время не поддерживается. immediate (bool): в настоящее время не поддерживается. priority (int): приоритет сообщения. Число от 0 до 9. content_type (str): тип содержимого. По умолчанию определяется автоматически. content_encoding (str): кодировка содержимого. По умолчанию определяется автоматически. serializer (str): используемый сериализатор. По умолчанию определяется автоматически. compression (str): используемый метод сжатия. По умолчанию отсутствует. headers (Dict): соответствие произвольных заголовков, передаваемых вместе

с телом сообщения.

exchange (kombu.entity.Exchange, str): переопределить биржу.

Обратите внимание, что эта биржа должна быть объявлена.

declare (Sequence[EntityT]): необязательный список необходимых сущностей,

которые должны быть объявлены до публикации сообщения. Сущности будут объявлены с помощью maybe_declare().

retry (bool): повторить публикацию или объявление сущностей, если

соединение потеряно.

retry_policy (Dict): конфигурация повторных попыток; поддерживаются ключевые слова,

указанные в ensure().

expiration (float): для каждого сообщения можно задать TTL в секундах.

По умолчанию срок действия не ограничен.

timeout (float): задаёт максимальное время ожидания публикации сообщения

в секундах.

confirm_timeout (float): задаёт максимальное время ожидания подтверждения

публикации сообщения, если для канала включён режим подтверждения публикации.

**properties (Any): дополнительные свойства сообщения; см. спецификацию AMQP.

release()
revive(channel)

Восстановить работу издателя после потери соединения.

routing_key=''

Ключ маршрутизации по умолчанию.

serializer=None

Сериализатор по умолчанию. По умолчанию используется JSON.

classQueue(name='', exchange=None, routing_key='', channel=None, bindings=None, on_declared=None, **kwargs)

Очередь, которая никогда не кэширует сведения об объявлении.

can_cache_declaration=False

Определяет, может ли maybe_declare пропустить повторное объявление этой сущности.

classResultConsumer(*args, **kwargs)
classConsumer(channel, queues=None, no_ack=None, auto_declare=None, callbacks=None, on_decode_error=None, on_message=None, accept=None, prefetch_count=None, tag_prefix=None)

Потребитель сообщений.

Аргументы:

channel (kombu.Connection, ChannelT): см. channel. queues (Sequence[kombu.Queue]): см. queues. no_ack (bool): см. no_ack. auto_declare (bool): см. auto_declare callbacks (Sequence[Callable]): см. callbacks. on_message (Callable): см. on_message on_decode_error (Callable): см. on_decode_error. prefetch_count (int): см. prefetch_count.

exceptionContentDisallowed

Потребителю не разрешён этот тип содержимого.

accept=None

Список принимаемых типов содержимого.

Если потребитель получит сообщение с недоверенным типом содержимого, будет вызвано исключение. По умолчанию принимаются все типы содержимого, но не в случае вызова kombu.disable_untrusted_serializers() — тогда разрешён только json.

add_queue(queue)

Добавить очередь в список очередей для получения сообщений.

Примечание:

Это не запустит получение сообщений из очереди. Для этого нужно после этого вызвать consume().

auto_declare=True

По умолчанию все сущности объявляются при создании экземпляра. Чтобы управлять этим вручную, можно задать значение False.

callbacks=None

Список обратных вызовов, которые вызываются по порядку при получении сообщения.

Сигнатура обратных вызовов должна принимать два аргумента: (body, message), где body — декодированное тело сообщения, а message — экземпляр Message.

cancel()

Завершить работу всех активных потребителей очередей.

Примечание:

Это не влияет на уже доставленные сообщения, но означает, что сервер больше не будет отправлять сообщения этому потребителю.

cancel_by_queue(queue)

Отменить получение сообщений из очереди по её имени.

channel=None

Соединение/канал, используемый этим потребителем.

close()

Завершить работу всех активных потребителей очередей.

Примечание:

Это не влияет на уже доставленные сообщения, но означает, что сервер больше не будет отправлять сообщения этому потребителю.

propertyconnection
consume(no_ack=None)

Начать получение сообщений.

Метод можно вызывать несколько раз. Обратите внимание: при этом будут обрабатываться новые очереди, добавленные после предыдущего вызова, но получение сообщений из удалённых очередей не будет отменено (используйте cancel_by_queue()).

Аргументы:

no_ack (bool): см. no_ack.

consuming_from(queue)

Возвращает True, если в данный момент выполняется получение сообщений из очереди.

declare()

Объявить очереди, обменники и привязки.

Примечание:

Это выполняется автоматически при создании экземпляра, если задан auto_declare.

flow(active)

Включить/отключить передачу данных от узла.

Это простой механизм управления потоком, который узел может использовать, чтобы избежать переполнения очередей или ситуации, когда он получает больше сообщений, чем может обработать.

Узел, получивший запрос прекратить отправку содержимого, завершит отправку текущего содержимого (если оно есть), а затем будет ждать возобновления передачи.

no_ack=None

Флаг автоматического подтверждения сообщений. Если он включён, брокер автоматически подтверждает сообщения. Это может повысить производительность, но лишает вас контроля над тем, когда сообщение удаляется.

По умолчанию отключено.

on_decode_error=None

Обратный вызов, вызываемый, если сообщение не удаётся декодировать.

Сигнатура обратного вызова должна принимать два аргумента: (message, exc), где message — сообщение, которое не удалось декодировать, а exc — исключение, возникшее при попытке его декодирования.

on_message=None

Необязательная функция, вызываемая при получении сообщения.

Если эта функция задана, она будет вызываться вместо метода receive(), а callbacks будут отключены.

Таким образом, её можно использовать вместо callbacks, если тело не нужно декодировать автоматически. Обратите внимание: сообщение всё равно будет распаковано, если в нём задан заголовок compression.

Сигнатура обратного вызова должна принимать один аргумент — объект Message.

Также обратите внимание: атрибут message.body, содержащий исходные данные тела сообщения, в некоторых случаях может быть объектом buffer, доступным только для чтения.

prefetch_count=None

Начальное количество предварительно выбираемых сообщений

Если задано, при запуске потребитель установит значение QoS prefetch_count. Его также можно изменить с помощью qos().

purge()

Удалить сообщения из всех очередей.

Предупреждение:

Это удалит все готовые к обработке сообщения; отменить это действие нельзя.

qos(prefetch_size=0, prefetch_count=0, apply_global=False)

Задать параметры качества обслуживания.

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

Окно предварительной выборки игнорируется, если задан параметр no_ack.

Аргументы:

prefetch_size (int): задаёт размер окна предварительной выборки в октетах.

Сервер отправит сообщение заранее, если его размер не превышает доступный размер предварительной выборки (и оно также соответствует другим ограничениям предварительной выборки). Можно задать ноль, что означает «без конкретного ограничения», однако другие ограничения предварительной выборки могут по-прежнему применяться.

prefetch_count (int): задаёт размер окна предварительной выборки в

целых сообщениях.

apply_global (bool): применить новые настройки глобально ко всем каналам.

propertyqueues
receive(body, message)

Метод, вызываемый при получении сообщения.

Выполняет зарегистрированные callbacks.

Аргументы:

body (Any): декодированное тело сообщения. message (~kombu.Message): экземпляр сообщения.

вызывает NotImplementedError:

Если обратные вызовы потребителя не были зарегистрированы.

recover(requeue=False)

Повторно доставить неподтверждённые сообщения.

Запрашивает у брокера повторную доставку всех неподтверждённых сообщений указанного канала.

Аргументы:

requeue (bool): по умолчанию сообщения будут доставлены повторно

первоначальному получателю. Если для requeue задано значение true, сервер попытается повторно поместить сообщение в очередь и, возможно, доставит его другому подписчику.

register_callback(callback)

Зарегистрировать новый обратный вызов, который будет выполняться при получении сообщения.

Примечание:

Сигнатура обратного вызова должна принимать два аргумента: (body, message), где body — декодированное тело сообщения, а message — экземпляр Message.

revive(channel)

Возобновить работу потребителя после потери соединения.

cancel_for(task_id)
consume_from(task_id)
drain_events(timeout=None)
on_after_fork()
start(initial_task_id, no_ack=True, **kwargs)
stop()
as_uri(include_password=True)

Возвращает URI бэкенда, скрывая или не скрывая пароль.

propertybinding
delete_group(group_id)
destination_for(task_id, request)

Получить адрес назначения результата по идентификатору задачи.

Возвращает:

кортеж из (reply_to, correlation_id).

Тип возвращаемого значения:

Tuple[str, str]

ensure_chords_allowed()
get_task_meta(task_id, backlog_limit=1000)

Получить метаданные задачи из бэкенда.

Если включён параметр always_retry_backend_operation и возникает восстанавливаемое исключение, операция повторяется с экспоненциальной задержкой до достижения установленного предела.

propertyoid
on_out_of_band_result(task_id, message)
on_reply_declare(task_id)
on_result_fulfilled(result)
on_task_call(producer, task_id)
persistent=False

Установите значение true, если бэкенд по умолчанию является постоянным.

poll(task_id, backlog_limit=1000)

Получить метаданные задачи из бэкенда.

Если включён параметр always_retry_backend_operation и возникает восстанавливаемое исключение, операция повторяется с экспоненциальной задержкой до достижения установленного предела.

reload_group_result(task_id)

Повторно загрузить результат группы, даже если он уже был получен ранее.

reload_task_result(task_id)

Повторно загрузить результат задачи, даже если он уже был получен ранее.

restore_group(group_id, cache=True)

Получить результат группы.

retry_policy={'interval_max': 1, 'interval_start': 0, 'interval_step': 1, 'max_retries': 20}
revive(channel)
save_group(group_id, result)

Сохранить результат выполненной группы.

store_result(task_id, result, state, traceback=None, request=None, **kwargs)

Отправить возвращаемое задачей значение и состояние.

supports_autoexpire=True

Если значение true, бэкенд должен автоматически удалять результаты по истечении срока хранения. В этом случае ежедневная периодическая задача backend_cleanup запускаться не будет.

supports_native_join=True

Если значение true, бэкенд должен реализовывать get_many().

Copyright © 2017-2026 Asif Saif Uddin, core team & contributors. All rights reserved.
Celery is licensed under The BSD License (3 Clause, also known as the new BSD license). The license is an OSI approved Open Source license and is GPL-compatible.
https://docs.celeryq.dev/en/stable/internals/reference/celery.backends.rpc.html

Spec-Zone.ru

Настройки Оффлайн Что нового Помощь О нас
Spec-Zone .ru
спецификации, руководства, описания, API