Spec-Zone.ru › Celery

celery.app.control

Клиент удалённого управления рабочими узлами.

Клиент для команд удалённого управления рабочими узлами. Реализация на стороне сервера находится в celery.worker.control. Существует два типа команд удалённого управления:

  • Команды проверки: не имеют побочных эффектов и обычно лишь возвращают какие-либо значения, полученные от рабочего узла, например список зарегистрированных в данный момент задач, список активных задач и т. д. Доступ к командам осуществляется через класс Inspect.

  • Команды управления: выполняют действия, имеющие побочные эффекты, например добавляют новую очередь для обработки. Доступ к командам осуществляется через класс Control.

classcelery.app.control.Control(app=None)

Клиент удалённого управления рабочими узлами.

classMailbox(namespace, type='direct', connection=None, clock=None, accept=None, serializer=None, producer_pool=None, queue_ttl=None, queue_expires=None, queue_durable=False, queue_exclusive=False, reply_queue_ttl=None, reply_queue_expires=10.0)

Обрабатывает почтовый ящик.

Node(hostname=None, state=None, channel=None, handlers=None)
abcast(command, kwargs=None)
accept=['json']

По умолчанию принимает только сообщения в формате json.

call(destination, command, kwargs=None, timeout=None, callback=None, channel=None)
cast(destination, command, kwargs=None)
connection=None

Соединение (если привязано).

exchange=None

Обменник почтового ящика (инициализируется конструктором).

exchange_fmt='%s.pidbox'
get_queue(hostname)
get_reply_queue()
multi_call(command, kwargs=None, timeout=1, limit=None, callback=None, channel=None)
namespace=None

Имя приложения.

node_cls

псевдоним для Node

propertyoid
producer_or_acquire(producer=None, channel=None)
propertyproducer_pool
reply_exchange=None

Обменник для отправки ответов.

reply_exchange_fmt='reply.%s.pidbox'
propertyreply_queue
serializer=None

Сериализатор сообщений

type='direct'

Тип обменника (обычно direct или fanout для широковещательной рассылки).

add_consumer(queue, exchange=None, exchange_type='direct', routing_key=None, options=None, destination=None, **kwargs)

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

Требуется указать только имя очереди: если указана только очередь, обменнику и ключу маршрутизации будет присвоено то же имя (как и для автоматически создаваемых очередей).

Примечание

Эта команда не учитывает параметры очереди и обменника по умолчанию из конфигурации.

Параметры:
  • queue (str) – Имя очереди, из которой следует начать получать сообщения.

  • exchange (str) – Необязательное имя обменника.

  • exchange_type (str) – Тип обменника (по умолчанию «direct»); если значение пустое, команда будет отправлена всем рабочим узлам.

  • routing_key (str) – Необязательный ключ маршрутизации.

  • options (Dict) – Дополнительные параметры, поддерживаемые kombu.entity.Queue.from_dict().

См. также

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

autoscale(max, min, destination=None, **kwargs)

Изменяет настройки автомасштабирования рабочего узла или узлов.

См. также

Поддерживает те же аргументы, что и broadcast().

broadcast(command, arguments=None, destination=None, connection=None, reply=False, timeout=1.0, limit=None, callback=None, channel=None, pattern=None, matcher=None, **extra_kwargs)

Отправляет широковещательную команду управления рабочим узлам Celery.

Параметры:
  • command (str) – Имя команды для отправки.

  • arguments (Dict) – Именованные аргументы команды.

  • destination (List) – Если задано, список узлов, которым нужно отправить команду; если список пуст, команда рассылается всем рабочим узлам.

  • connection (kombu.Connection) – Пользовательское соединение с брокером; если не задано, соединение будет получено из пула.

  • reply (bool) – Ожидать ответ и возвращать его.

  • timeout (float) – Время ожидания ответа в секундах.

  • limit (int) – Ограничение количества ответов.

  • callback (Callable) – Функция обратного вызова, вызываемая сразу при получении каждого ответа.

  • pattern (str) – Пользовательская строка шаблона для сопоставления

  • matcher (Callable) – Пользовательская функция для сопоставления с шаблоном

cancel_consumer(queue, destination=None, **kwargs)

Указывает всем (или выбранным) рабочим узлам прекратить получение сообщений из queue.

См. также

Поддерживает те же аргументы, что и broadcast().

disable_events(destination=None, **kwargs)

Указывает всем (или выбранным) рабочим узлам отключить события.

См. также

Поддерживает те же аргументы, что и broadcast().

discard_all(connection=None)

Отбрасывает все ожидающие задачи.

Все ожидающие выполнения задачи будут проигнорированы и удалены с сервера сообщений.

Параметры:

connection (kombu.Connection) – Необязательный экземпляр соединения. Если он не указан, соединение будет получено из пула соединений.

Возвращает:

количество отброшенных задач.

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

int

election(id, topic, action=None, connection=None)
enable_events(destination=None, **kwargs)

Указывает всем (или выбранным) рабочим узлам включить события.

См. также

Поддерживает те же аргументы, что и broadcast().

heartbeat(destination=None, **kwargs)

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

См. также

Поддерживает те же аргументы, что и broadcast()

propertyinspect

Создаёт новый экземпляр Inspect.

ping(destination=None, timeout=1.0, **kwargs)

Проверяет связь со всеми (или выбранными) рабочими узлами.

>>> app.control.ping()
[{'celery@node1': {'ok': 'pong'}}, {'celery@node2': {'ok': 'pong'}}]
>>> app.control.ping(destination=['celery@node2'])
[{'celery@node2': {'ok': 'pong'}}]
Возвращает:

Список словарей {HOSTNAME: {'ok': 'pong'}}.

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

List[Dict]

См. также

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

pool_grow(n=1, destination=None, **kwargs)

Указывает всем (или выбранным) рабочим узлам увеличить пул на n.

См. также

Поддерживает те же аргументы, что и broadcast().

pool_restart(modules=None, reload=False, reloader=None, destination=None, **kwargs)

Перезапустить пулы выполнения всех или выбранных рабочих процессов.

Именованные аргументы:
  • modules (Sequence[str]) – Список модулей для перезагрузки.

  • reload (bool) – Флаг включения перезагрузки модулей. По умолчанию — False.

  • reloader (Any) – Функция для перезагрузки модуля.

  • destination (Sequence[str]) – Список имён рабочих процессов, которым нужно отправить эту команду.

См. также

Поддерживаются те же аргументы, что и для broadcast()

pool_shrink(n=1, destination=None, **kwargs)

Указать всем (или выбранным) рабочим процессам уменьшить пул на n.

См. также

Поддерживаются те же аргументы, что и для broadcast().

purge(connection=None)

Удалить все ожидающие задачи.

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

Параметры:

connection (kombu.Connection) – Необязательный экземпляр подключения для использования. Если он не указан, подключение будет получено из пула подключений.

Возвращает:

количество удалённых задач.

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

int

rate_limit(task_name, rate_limit, destination=None, **kwargs)

Указать рабочим процессам установить новое ограничение частоты выполнения задач определённого типа.

Параметры:
  • task_name (str) – Имя задачи, для которой нужно изменить ограничение частоты.

  • rate_limit (int, str) – Ограничение частоты в задачах в секунду или строка ограничения частоты (‘100/m’ и т. д.; дополнительную информацию см. в celery.app.task.Task.rate_limit).

См. также

Сведения о поддерживаемых именованных аргументах см. в broadcast().

revoke(task_id, destination=None, terminate=False, signal='SIGTERM', **kwargs)

Указать всем (или выбранным) рабочим процессам отменить задачу по её идентификатору (или списку идентификаторов).

Если задача отменена, рабочие процессы проигнорируют её и не будут выполнять.

Параметры:
  • task_id (Union(str, list)) – Идентификатор задачи, которую нужно отменить (или список идентификаторов).

  • terminate (bool) – Также завершить процесс, который в данный момент выполняет задачу (если такой есть).

  • signal (str) – Имя сигнала, отправляемого процессу при завершении. По умолчанию — TERM.

См. также

Сведения о поддерживаемых именованных аргументах см. в broadcast().

revoke_by_stamped_headers(headers, destination=None, terminate=False, signal='SIGTERM', **kwargs)

Указать всем (или выбранным) рабочим процессам отменить задачу по заголовкам.

Если задача отменена, рабочие процессы проигнорируют её и не будут выполнять.

Параметры:
  • headers (dict[str, Union(str, list)]) – Заголовки для сопоставления при отмене задач.

  • terminate (bool) – Также завершить процесс, который в данный момент выполняет задачу (если такой есть).

  • signal (str) – Имя сигнала, отправляемого процессу при завершении. По умолчанию — TERM.

См. также

Сведения о поддерживаемых именованных аргументах см. в broadcast().

shutdown(destination=None, **kwargs)

Завершить работу рабочего процесса (или процессов).

См. также

Поддерживаются те же аргументы, что и для broadcast()

terminate(task_id, destination=None, signal='SIGTERM', **kwargs)

Указать всем (или выбранным) рабочим процессам завершить задачу по её идентификатору (или списку идентификаторов).

См. также

Это просто сокращённая форма вызова revoke() с включённым аргументом terminate.

time_limit(task_name, soft=None, hard=None, destination=None, **kwargs)

Указать рабочим процессам установить ограничения времени выполнения для задач определённого типа.

Параметры:
  • task_name (str) – Имя задачи, для которой нужно изменить ограничения времени выполнения.

  • soft (float) – Новое мягкое ограничение времени выполнения (в секундах).

  • hard (float) – Новое жёсткое ограничение времени выполнения (в секундах).

  • **kwargs (Any) – аргументы, передаваемые в broadcast().

classcelery.app.control.Inspect(destination=None, timeout=1.0, callback=None, connection=None, app=None, limit=None, pattern=None, matcher=None)

API для проверки состояния рабочих процессов.

Этот класс предоставляет прокси для доступа к API Inspect рабочих процессов. API определен в celery.worker.control

active(safe=None)

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

Параметры:

safe (Boolean) – Установите значение True, чтобы отключить десериализацию.

Возвращает:

Словарь {HOSTNAME: [TASK_INFO,...]}.

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

Dict

См. также

Подробности о TASK_INFO см. в возвращаемом значении query_task().

active_queues()

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

Возвращает:

Словарь {HOSTNAME: [QUEUE_INFO, QUEUE_INFO,...]}.

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

Dict

Ниже перечислены поля QUEUE_INFO:

  • name

  • exchange
    • name

    • type

    • arguments

    • durable

    • passive

    • auto_delete

    • delivery_mode

    • no_declare

  • routing_key

  • queue_arguments

  • binding_arguments

  • consumer_arguments

  • durable

  • exclusive

  • auto_delete

  • no_ack

  • alias

  • bindings

  • no_declare

  • expires

  • message_ttl

  • max_length

  • max_length_bytes

  • max_priority

См. также

Дополнительные сведения о полях queue_info см. в документации RabbitMQ/AMQP.

Примечание

Поля queue_info предназначены для RabbitMQ/AMQP. Не все поля применимы к другим транспортам.

app=None
clock()

Получает значение Clock на рабочих процессах.

>>> app.control.inspect().clock()
{'celery@node1': {'clock': 12}}
Возвращает:

Словарь {HOSTNAME: CLOCK_VALUE}.

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

Dict

conf(with_defaults=False)

Возвращает конфигурацию каждого рабочего процесса.

Параметры:

with_defaults (bool) – если задано значение True, метод также возвращает параметры конфигурации со значениями по умолчанию.

Возвращает:

Словарь {HOSTNAME: WORKER_CONFIGURATION}.

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

Dict

См. также

WORKER_CONFIGURATION — это словарь, содержащий текущие параметры конфигурации. Возможные значения см. в разделе Конфигурация и значения по умолчанию.

hello(from_node, revoked=None)
memdump(samples=10)

Выводит статистику предыдущих запросов memsample.

Примечание

Требуется библиотека psutils.

memsample()

Возвращает образец текущего использования памяти RSS.

Примечание

Требуется библиотека psutils.

objgraph(type='Request', n=200, max_depth=10)

Создает граф не освобожденных сборщиком мусора объектов (отладка утечек памяти).

Параметры:
  • n (int) – Максимальное число объектов для построения графа.

  • max_depth (int) – Обход не глубже n уровней.

  • type (str) – Имя объекта для построения графа. По умолчанию — "Request".

Возвращает:

Словарь {'filename': FILENAME}

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

Dict

Примечание

Требуется библиотека objgraph.

ping(destination=None)

Проверяет доступность всех (или указанных) рабочих процессов.

>>> app.control.inspect().ping()
{'celery@node1': {'ok': 'pong'}, 'celery@node2': {'ok': 'pong'}}
>>> app.control.inspect().ping(destination=['celery@node1'])
{'celery@node1': {'ok': 'pong'}}
Параметры:

destination (List) – Если задано, список узлов, которым отправляется команда; если список пуст, команда рассылается всем рабочим процессам.

Возвращает:

Словарь {HOSTNAME: {'ok': 'pong'}}.

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

Dict

См. также

broadcast() — поддерживаемые именованные аргументы.

query_task(*ids)

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

Параметры:

*ids (str) – Идентификаторы задач, сведения о которых нужно запросить.

Возвращает:

Словарь {HOSTNAME: {TASK_ID: [STATE, TASK_INFO]}}.

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

Dict

Ниже перечислены поля TASK_INFO:
  • id — идентификатор задачи

  • name — имя задачи

  • args — позиционные аргументы, переданные задаче

  • kwargs — именованные аргументы, переданные задаче

  • type — тип задачи

  • hostname — имя узла рабочего процесса, обрабатывающего задачу

  • time_start — время начала обработки

  • acknowledged — True, если задача была подтверждена брокеру

  • delivery_info — словарь со сведениями о доставке
    • exchange — имя обменника, в который была опубликована задача

    • routing_key — ключ маршрутизации, использованный при публикации задачи

    • priority — приоритет, использованный при публикации задачи

    • redelivered — True, если задача была доставлена повторно

  • worker_pid — PID рабочего процесса, обрабатывающего задачу

registered(*taskinfoitems)

Возвращает все зарегистрированные задачи для каждого рабочего процесса.

>>> app.control.inspect().registered()
{'celery@node1': ['task1', 'task1']}
>>> app.control.inspect().registered('serializer', 'max_retries')
{'celery@node1': ['task_foo [serializer=json max_retries=3]', 'tasb_bar [serializer=json max_retries=3]']}
Параметры:

taskinfoitems (Sequence[str]) – Список атрибутов Task, которые нужно включить.

Возвращает:

Словарь {HOSTNAME: [TASK1_INFO, ...]}.

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

Dict

registered_tasks(*taskinfoitems)

Возвращает все зарегистрированные задачи для каждого рабочего процесса.

>>> app.control.inspect().registered()
{'celery@node1': ['task1', 'task1']}
>>> app.control.inspect().registered('serializer', 'max_retries')
{'celery@node1': ['task_foo [serializer=json max_retries=3]', 'tasb_bar [serializer=json max_retries=3]']}
Параметры:

taskinfoitems (Sequence[str]) – Список атрибутов Task, которые нужно включить.

Возвращает:

Словарь {HOSTNAME: [TASK1_INFO, ...]}.

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

Dict

report()

Возвращает понятный человеку отчет для каждого рабочего процесса.

Возвращает:

Словарь {HOSTNAME: {'ok': REPORT_STRING}}.

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

Dict

reserved(safe=None)

Возвращает список зарезервированных в данный момент задач, за исключением запланированных и активных.

Возвращает:

Словарь {HOSTNAME: [TASK_INFO,...]}.

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

Dict

См. также

Подробности о TASK_INFO см. в возвращаемом значении query_task().

revoked()

Возвращает список отозванных задач.

>>> app.control.inspect().revoked()
{'celery@node1': ['16f527de-1c72-47a6-b477-c472b92fef7a']}
Возвращает:

Словарь {HOSTNAME: [TASK_ID, ...]}.

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

Dict

scheduled(safe=None)

Возвращает список запланированных задач с подробной информацией.

Возвращает:

Словарь {HOSTNAME: [TASK_SCHEDULED_INFO,...]}.

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

Dict

Ниже перечислены поля TASK_SCHEDULED_INFO:

  • eta — запланированное время выполнения задачи в виде строки в формате ISO 8601

  • priority — приоритет задачи

  • request — поле, содержащее значение TASK_INFO.

См. также

Дополнительные сведения о TASK_INFO см. в возвращаемом значении query_task().

stats()

Возвращает статистику рабочего процесса.

Возвращает:

Словарь {HOSTNAME: STAT_INFO}.

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

Dict

Ниже перечислены поля STAT_INFO:

  • broker — раздел со сведениями о брокере.
    • connect_timeout — время ожидания в секундах (int/float) при установлении нового соединения.

    • heartbeat — текущее значение heartbeat (задается клиентом).

    • hostname — имя узла удаленного брокера.

    • insist — больше не используется.

    • login_method — метод входа, используемый для подключения к брокеру.

    • port — порт удаленного брокера.

    • ssl — состояние SSL (включен/отключен).

    • transport — имя используемого транспорта (например, amqp или redis)

    • transport_options — параметры, переданные транспорту.

    • uri_prefix — в некоторых транспортах имя хоста должно быть URL-адресом. Например, redis+socket:///tmp/redis.sock. В этом примере префиксом URI будет redis.

    • userid — идентификатор пользователя, используемый для подключения к брокеру.

    • virtual_host — используемый виртуальный хост.

  • clock — значение логических часов рабочего процесса. Это положительное целое число, которое должно увеличиваться при каждом получении статистики.

  • uptime — количество секунд с момента запуска контроллера рабочего процесса

  • pid — идентификатор процесса экземпляра рабочего процесса (основного процесса).

  • pool — раздел, относящийся к пулу.
    • max-concurrency — максимальное количество процессов/потоков/зеленых потоков.

    • max-tasks-per-child — максимальное количество задач, которые поток может выполнить до повторного использования.

    • processes — список PID (или идентификаторов потоков).

    • put-guarded-by-semaphore — внутреннее поле

    • timeouts — значения по умолчанию для ограничений времени.

    • writes — относится к пулу prefork; показывает распределение операций записи между процессами пула при использовании асинхронного ввода-вывода.

  • prefetch_count — текущее значение счетчика предварительной выборки для потребителя задач.

  • rusage — статистика использования системы. Доступные поля могут различаться в зависимости от платформы. Из getrusage(2):

    • stime — время, затраченное на выполнение кода операционной системы от имени этого процесса.

    • utime — время, затраченное на выполнение пользовательских инструкций.

    • maxrss — максимальный размер резидентной памяти, используемой этим процессом (в килобайтах).

    • idrss — объем несвязанной с другими процессами памяти, использованной для данных (в килобайтах, умноженных на число тактов выполнения)

    • isrss — объем несвязанной с другими процессами памяти, использованной для стека (в килобайтах, умноженных на число тактов выполнения)

    • ixrss — объем памяти, совместно используемой с другими процессами (в килобайтах, умноженных на число тактов выполнения).

    • inblock — количество обращений файловой системы к диску от имени этого процесса.

    • oublock — количество обращений файловой системы для записи на диск от имени этого процесса.

    • majflt — количество страничных ошибок, обработанных с выполнением операций ввода-вывода.

    • minflt — количество страничных ошибок, обработанных без выполнения операций ввода-вывода.

    • msgrcv — количество полученных сообщений IPC.

    • msgsnd — количество отправленных сообщений IPC.

    • nvcsw — количество добровольных переключений контекста, выполненных этим процессом.

    • nivcsw — количество принудительных переключений контекста.

    • nsignals — количество полученных сигналов.

    • nswap — количество случаев, когда этот процесс был полностью выгружен из памяти.

  • total — отображение имен типов задач и общего количества задач каждого типа, принятых рабочим процессом с момента запуска.

celery.app.control.flatten_reply(reply)

Объединяет ответы узлов.

Преобразует список ответов следующего формата:

[{'a@example.com': reply},
 {'b@example.com': reply}]

в следующий формат:

{'a@example.com': reply,
 'b@example.com': reply}

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/reference/celery.app.control.html

Spec-Zone.ru

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