Spec-Zone.ru › Celery

Руководство по мониторингу и управлению

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

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

celery также можно использовать для проверки узлов-исполнителей и управления ими (а в некоторой степени — и задачами).

Чтобы вывести список всех доступных команд, выполните:

$ celery--help

Чтобы получить справку по конкретной команде, выполните:

$ celery<command>--help
  • shell: перейти в оболочку Python.

    В локальной области видимости будет доступна переменная celery: это текущее приложение. Кроме того, все известные задачи будут автоматически добавлены в локальную область видимости (если не указан флаг --without-tasks).

    Если они установлены, последовательно используются https://pypi.org/project/Ipython/, https://pypi.org/project/bpython/ или обычный python. Вы можете принудительно выбрать реализацию с помощью --ipython, --bpython или --python.

  • status: вывести список активных узлов кластера

    $ celery-Aprojstatus
    
  • result: показать результат задачи

    $ celery-Aprojresult-ttasks.add4e196aa4-0141-4601-8138-7aa33db0f577
    

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

  • purge: удалить сообщения из всех настроенных очередей задач.

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

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

    Операцию нельзя отменить, сообщения будут удалены без возможности восстановления!

    $ celery-Aprojpurge
    

    Очереди для очистки также можно указать с помощью параметра -Q:

    $ celery-Aprojpurge-Qcelery,foo,bar
    

    а очереди, которые не нужно очищать, — с помощью параметра -X:

    $ celery-Aprojpurge-Xcelery
    
  • inspect active: вывести список активных задач

    $ celery-Aprojinspectactive
    

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

  • inspect scheduled: вывести список запланированных задач с ETA

    $ celery-Aprojinspectscheduled
    

    Это задачи, зарезервированные исполнителем, для которых задан аргумент eta или countdown.

  • inspect reserved: вывести список зарезервированных задач

    $ celery-Aprojinspectreserved
    

    Команда выводит все задачи, предварительно полученные исполнителем и ожидающие выполнения (задачи с заданным значением ETA не включаются).

  • inspect revoked: вывести историю отозванных задач

    $ celery-Aprojinspectrevoked
    
  • inspect registered: вывести список зарегистрированных задач

    $ celery-Aprojinspectregistered
    
  • inspect stats: показать статистику исполнителя (см. Статистика)

    $ celery-Aprojinspectstats
    
  • inspect query_task: показать сведения о задаче или задачах по идентификатору.

    Любой исполнитель, у которого задача с одним из этих идентификаторов зарезервирована или выполняется, ответит, указав статус и сведения о ней.

    $ celery-Aprojinspectquery_taske9f6c8f0-fec9-4ae8-a8c6-cf8c8451d4f8
    

    Также можно запросить сведения сразу о нескольких задачах:

    $ celery-Aprojinspectquery_taskid1id2...idN
    
  • control enable_events: включить события

    $ celery-Aprojcontrolenable_events
    
  • control disable_events: отключить события

    $ celery-Aprojcontroldisable_events
    
  • migrate: перенести задачи с одного брокера на другой (ЭКСПЕРИМЕНТАЛЬНО).

    $ celery-Aprojmigrateredis://localhostamqp://localhost
    

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

Примечание

Все команды inspect и control поддерживают аргумент --timeout — количество секунд ожидания ответа. Если ответ не приходит из-за задержки, возможно, потребуется увеличить время ожидания.

По умолчанию команды inspect и control выполняются для всех исполнителей. С помощью аргумента --destination можно указать одного или нескольких исполнителей:

$ celery-Aprojinspect-dw1@e.com,w2@e.comreserved

$ celery-Aprojcontrol-dw1@e.com,w2@e.comenable_events

Flower — это веб-инструмент для мониторинга и администрирования Celery в реальном времени. Он активно разрабатывается, но уже стал незаменимым инструментом. Будучи рекомендуемым средством мониторинга Celery, он заменяет монитор Django-Admin, celerymon и монитор на основе ncurses.

Название Flower произносится как «flow», но при желании можно использовать и ботанический вариант.

  • Мониторинг в реальном времени с использованием событий Celery

    • Ход выполнения задач и история

    • Возможность просматривать подробные сведения о задачах (аргументы, время начала, время выполнения и другое)

    • Графики и статистика

  • Удаленное управление

    • Просмотр состояния и статистики исполнителей

    • Остановка и перезапуск экземпляров исполнителей

    • Управление размером пула исполнителя и настройками автомасштабирования

    • Просмотр и изменение очередей, из которых получает задачи экземпляр исполнителя

    • Просмотр выполняющихся в данный момент задач

    • Просмотр запланированных задач (ETA/countdown)

    • Просмотр зарезервированных и отозванных задач

    • Настройка ограничений по времени и частоте выполнения

    • Просмотр конфигурации

    • Отзыв или принудительное завершение задач

  • HTTP API

    • Вывод списка исполнителей

    • Остановка исполнителя

    • Перезапуск пула исполнителя

    • Увеличение пула исполнителя

    • Уменьшение пула исполнителя

    • Автомасштабирование пула исполнителя

    • Начало получения задач из очереди

    • Прекращение получения задач из очереди

    • Вывод списка задач

    • Вывод списка типов задач (с которыми уже работали)

    • Получение сведений о задаче

    • Запуск задачи

    • Запуск задачи по имени

    • Получение результата задачи

    • Изменение мягкого и жесткого ограничений по времени выполнения задачи

    • Изменение ограничения частоты выполнения задачи

    • Отзыв задачи

  • Аутентификация OpenID

Снимки экрана

../_images/dashboard.png

Другие снимки экрана:

Для установки Flower можно использовать pip:

$ pipinstallflower

Команда flower запускает веб-сервер, который можно открыть в браузере:

$ celery-Aprojflower

По умолчанию используется порт http://localhost:5555, но его можно изменить с помощью аргумента –port:

$ celery-Aprojflower--port=5555

URL брокера также можно передать с помощью аргумента --broker:

$ celery--broker=amqp://guest:guest@localhost:5672//flower
or
$ celery--broker=redis://guest:guest@localhost:6379/0flower

Затем откройте Flower в веб-браузере:

$ openhttp://localhost:5555

У Flower есть множество других возможностей, не описанных здесь, включая параметры авторизации. Дополнительные сведения см. в официальной документации.

Добавлено в версии 2.0.

celery events — это простой монитор на основе curses, отображающий историю задач и исполнителей. С его помощью можно проверять результаты задач и трассировки стека; он также поддерживает некоторые команды управления, например ограничение частоты выполнения и остановку исполнителей. Этот монитор был создан как проверка концепции, поэтому, вероятно, вам лучше использовать Flower.

Запуск:

$ celery-Aprojevents

На экране должно появиться примерно следующее:

../_images/celeryevshotsm.jpg

celery events также используется для запуска камер снимков состояния (см. раздел Снимки состояния):

$ celery-Aprojevents--camera=<camera-class>--frequency=1.0

Кроме того, в него входит инструмент для выгрузки событий в stdout:

$ celery-Aprojevents--dump

Полный список параметров можно получить с помощью --help:

$ celeryevents--help

Для управления кластером Celery важно знать, как выполнять мониторинг RabbitMQ.

В RabbitMQ входит команда rabbitmqctl(1), с помощью которой можно просматривать очереди, обменники, привязки, длину очередей и объем памяти, используемый каждой очередью, а также управлять пользователями, виртуальными хостами и их разрешениями.

Примечание

В этих примерах используется виртуальный хост по умолчанию ("/"). Если используется пользовательский виртуальный хост, к команде нужно добавить аргумент -p, например: rabbitmqctl list_queues -p my_vhost …

Чтобы узнать количество задач в очереди:

$ rabbitmqctllist_queuesnamemessagesmessages_ready\
messages_unacknowledged

Здесь messages_ready — количество сообщений, готовых к доставке (отправленных, но еще не полученных), а messages_unacknowledged — количество сообщений, полученных исполнителем, но еще не подтвержденных (то есть находящихся в обработке или зарезервированных). messages — это сумма готовых и неподтвержденных сообщений.

Чтобы узнать количество исполнителей, которые в данный момент получают задачи из очереди:

$ rabbitmqctllist_queuesnameconsumers

Чтобы узнать объем памяти, выделенной очереди:

$ rabbitmqctllist_queuesnamememory
Совет:

Если добавить параметр -q к команде rabbitmqctl(1), вывод будет проще обработать.

Если в качестве брокера используется Redis, для мониторинга кластера Celery можно использовать команду redis-cli(1), чтобы просматривать длину очередей.

Чтобы узнать количество задач в очереди:

$ redis-cli-hHOST-pPORT-nDATABASE_NUMBERllenQUEUE_NAME

Очередь по умолчанию называется celery. Чтобы получить список всех доступных очередей, выполните:

$ redis-cli-hHOST-pPORT-nDATABASE_NUMBERkeys\*

Примечание

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

Кроме того, если Redis используется и для других целей, вывод команды keys будет содержать посторонние значения, хранящиеся в базе данных. Рекомендуется использовать для Celery выделенный номер базы данных DATABASE_NUMBER. Номера баз данных также можно использовать для разделения приложений Celery (виртуальных хостов), однако это не повлияет на события мониторинга, используемые, например, Flower, поскольку команды Redis pub/sub являются глобальными, а не привязанными к базе данных.

Хотя Prometheus не является встроенной частью Celery, с помощью Flower можно легко отслеживать исполнителей Celery через Prometheus. Flower также предоставляет готовые панели Grafana, позволяющие строить графики количества задач, исполнителей и других показателей.

Инструкции по настройке Prometheus см. в документации Flower: https://flower.readthedocs.io/en/latest/prometheus-integration.html

Ниже приведен список известных подключаемых модулей Munin, полезных при обслуживании кластера Celery.

  • rabbitmq-munin: подключаемые модули Munin для RabbitMQ.

    https://github.com/ask/rabbitmq-munin

  • celery_tasks: отслеживает количество выполнений каждого типа задач (требуется celerymon).

    https://github.com/munin-monitoring/contrib/blob/master/plugins/celery/celery_tasks

  • celery_tasks_states: отслеживает количество задач в каждом состоянии (требуется celerymon).

    https://github.com/munin-monitoring/contrib/blob/master/plugins/celery/celery_tasks_states

Исполнитель может отправлять сообщение при возникновении какого-либо события. Такие события перехватываются инструментами мониторинга кластера, например Flower и celery events.

Добавлено в версии 2.1.

Даже один исполнитель может создавать огромное количество событий, поэтому хранение истории всех событий на диске может обойтись очень дорого.

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

Для создания снимков требуется класс Camera, который позволяет определить, что должно происходить при каждом сохранении состояния: его можно записать в базу данных, отправить по электронной почте или выполнить любое другое действие.

Для создания снимков с помощью камеры используется celery events. Например, чтобы сохранять состояние каждые 2 секунды с помощью камеры myapp.Camera, запустите celery events со следующими аргументами:

$ celery-Aprojevents-cmyapp.Camera--frequency=2.0

Камеры полезны, если нужно перехватывать события и выполнять с ними какие-либо действия с заданным интервалом. Для обработки событий в реальном времени следует напрямую использовать app.events.Receiver, как показано в разделе Обработка в реальном времени.

Ниже приведен пример камеры, которая выводит снимок состояния на экран:

from pprint import pformat

from celery.events.snapshot import Polaroid

class DumpCam(Polaroid):
    clear_after = True  # clear after flush (incl, state.event_count).

    def on_shutter(self, state):
        if not state.event_count:
            # No new events since last snapshot.
            return
        print('Workers: {0}'.format(pformat(state.workers, indent=4)))
        print('Tasks: {0}'.format(pformat(state.tasks, indent=4)))
        print('Total: {0.event_count} events, {0.task_count} tasks'.format(
            state))

Дополнительные сведения об объектах состояния см. в справочнике API для celery.events.state.

Теперь эту камеру можно использовать с celery events, указав ее с помощью параметра -c:

$ celery-Aprojevents-cmyapp.DumpCam--frequency=2.0

Или можно использовать ее программно следующим образом:

from celery import Celery
from myapp import DumpCam

def main(app, freq=1.0):
    state = app.events.State()
    with app.connection() as connection:
        recv = app.events.Receiver(connection, handlers={'*': state.event})
        with DumpCam(state, freq=freq):
            recv.capture(limit=None, timeout=None)

if __name__ == '__main__':
    app = Celery(broker='amqp://guest@localhost//')
    main(app)

Для обработки событий в реальном времени необходимы:

  • Потребитель событий (это Receiver)

  • Набор обработчиков, вызываемых при поступлении событий.

    Для каждого типа события можно задать отдельный обработчик либо использовать обработчик для всех событий («*»).

  • Состояние (необязательно)

    app.events.State — это удобное представление задач и исполнителей кластера в памяти, которое обновляется по мере поступления событий.

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

Все это позволяет легко обрабатывать события в реальном времени:

from celery import Celery


def my_monitor(app):
    state = app.events.State()

    def announce_failed_tasks(event):
        state.event(event)
        # task name is sent only with -received event, and state
        # will keep track of this for us.
        task = state.tasks.get(event['uuid'])

        print('TASK FAILED: %s[%s] %s' % (
            task.name, task.uuid, task.info(),))

    with app.connection() as connection:
        recv = app.events.Receiver(connection, handlers={
                'task-failed': announce_failed_tasks,
                '*': state.event,
        })
        recv.capture(limit=None, timeout=None, wakeup=True)

if __name__ == '__main__':
    app = Celery(broker='amqp://guest@localhost//')
    my_monitor(app)

Примечание

Аргумент wakeup команды capture отправляет всем исполнителям сигнал с требованием передать контрольный сигнал. Благодаря этому монитор сразу обнаруживает исполнителей при запуске.

Чтобы прослушивать определенные события, укажите обработчики:

from celery import Celery

def my_monitor(app):
    state = app.events.State()

    def announce_failed_tasks(event):
        state.event(event)
        # task name is sent only with -received event, and state
        # will keep track of this for us.
        task = state.tasks.get(event['uuid'])

        print('TASK FAILED: %s[%s] %s' % (
            task.name, task.uuid, task.info(),))

    with app.connection() as connection:
        recv = app.events.Receiver(connection, handlers={
                'task-failed': announce_failed_tasks,
        })
        recv.capture(limit=None, timeout=None, wakeup=True)

if __name__ == '__main__':
    app = Celery(broker='amqp://guest@localhost//')
    my_monitor(app)

В этом списке перечислены события, отправляемые исполнителем, и их аргументы.

сигнатура:

task-sent(uuid, name, args, kwargs, retries, eta, expires, queue, exchange, routing_key, root_id, parent_id)

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

сигнатура:

task-received(uuid, name, args, kwargs, retries, eta, hostname, timestamp, root_id, parent_id)

Отправляется при получении задачи исполнителем.

сигнатура:

task-started(uuid, hostname, timestamp, pid)

Отправляется непосредственно перед выполнением задачи исполнителем.

сигнатура:

task-succeeded(uuid, result, runtime, hostname, timestamp)

Отправляется, если задача выполнена успешно.

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

сигнатура:

task-failed(uuid, exception, traceback, hostname, timestamp)

Отправляется, если выполнить задачу не удалось.

сигнатура:

task-rejected(uuid, requeue)

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

сигнатура:

task-revoked(uuid, terminated, signum, expired)

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

  • terminated принимает значение true, если процесс задачи был завершен принудительно,

    а в поле signum указывается использованный сигнал.

  • expired принимает значение true, если срок действия задачи истек.

сигнатура:

task-retried(uuid, exception, traceback, hostname, timestamp)

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

сигнатура:

worker-online(hostname, timestamp, freq, sw_ident, sw_ver, sw_sys)

Исполнитель подключился к брокеру и находится в сети.

  • hostname: имя узла исполнителя.

  • timestamp: временная метка события.

  • freq: частота контрольных сигналов в секундах (число с плавающей точкой).

  • sw_ident: название программного обеспечения исполнителя (например, py-celery).

  • sw_ver: версия программного обеспечения (например, 2.2.0).

  • sw_sys: операционная система (например, Linux/Darwin).

сигнатура:

worker-heartbeat(hostname, timestamp, freq, sw_ident, sw_ver, sw_sys, active, processed)

Отправляется каждую минуту. Если исполнитель не отправлял контрольный сигнал в течение 2 минут, он считается отключенным.

  • hostname: имя узла исполнителя.

  • timestamp: временная метка события.

  • freq: частота контрольных сигналов в секундах (число с плавающей точкой).

  • sw_ident: название программного обеспечения исполнителя (например, py-celery).

  • sw_ver: версия программного обеспечения (например, 2.2.0).

  • sw_sys: операционная система (например, Linux/Darwin).

  • active: количество выполняемых в данный момент задач.

  • processed: общее количество задач, обработанных этим исполнителем.

сигнатура:

worker-offline(hostname, timestamp, freq, sw_ident, sw_ver, sw_sys)

Исполнитель отключился от брокера.

Celery использует kombu.pidbox.Mailbox внутри системы для отправки исполнителям команд управления и широковещательных команд.

Добавлено в версии Kombu: 5.6.0

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

  • durable (по умолчанию: False): если установлено значение True, обменники управления сохраняются после перезапуска брокера.

  • exclusive (по умолчанию: False): если установлено значение True, обменники будут доступны только одному подключению.

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

Одновременная установка durable=True и exclusive=True запрещена и приведет к ошибке, поскольку эти параметры взаимоисключают друг друга в AMQP.

Дополнительные сведения о расширенной настройке см. в разделах event_queue_durable и event_queue_exclusive.

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/userguide/monitoring.html

Spec-Zone.ru

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