Spec-Zone.ru › Celery

Руководство по воркерам

Запуск в качестве демона

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

Вы можете запустить воркер в переднем плане, выполнив команду:

$ celery-Aprojworker-lINFO

Полный список доступных параметров командной строки см. в разделе worker или просто выполните:

$ celeryworker--help

На одной машине можно запустить несколько воркеров, но обязательно задайте имя каждому воркеру, указав имя узла с помощью аргумента --hostname:

$ celery-Aprojworker--loglevel=INFO--concurrency=10-nworker1@%h
$ celery-Aprojworker--loglevel=INFO--concurrency=10-nworker2@%h
$ celery-Aprojworker--loglevel=INFO--concurrency=10-nworker3@%h

Аргумент hostname может подставлять следующие переменные:

  • %h: имя хоста, включая доменное имя.

  • %n: только имя хоста.

  • %d: только доменное имя.

Если текущее имя хоста — george.example.com, подстановка даст следующие значения:

Примечание для пользователей https://pypi.org/project/supervisor/

Знак % необходимо экранировать, добавив второй знак: %%h.

Для завершения работы следует отправить сигнал TERM.

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

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

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

$ pkill-9-f'celery worker'

Если в вашей системе нет команды pkill, можно воспользоваться более длинным вариантом:

$ psauxww|awk'/celery worker/ {print $2}'|xargskill-9

Изменено в версии 5.2: В системах Linux Celery теперь поддерживает отправку сигнала KILL всем дочерним процессам после завершения работы воркера. Это реализовано с помощью параметра PR_SET_PDEATHSIG модуля prctl(2).

Завершение работы воркера

Для описания различных этапов завершения работы воркера мы будем использовать термины мягкое, постепенное, холодное и принудительное завершение. Воркер начинает процесс завершения работы, получив сигнал TERM или QUIT. Сигнал INT (Ctrl-C) также обрабатывается во время завершения работы и всегда запускает следующий этап этого процесса.

Мягкое завершение работы

Получив сигнал TERM, воркер начинает мягкое завершение работы. Он дождётся окончания всех выполняющихся задач и только затем завершит работу. При первом получении сигнала INT (Ctrl-C) воркер также начинает мягкое завершение работы.

При мягком завершении работы вызов WorkController.start() будет остановлен, а затем будет вызван WorkController.stop().

  • Дополнительные сигналы TERM будут игнорироваться во время мягкого завершения работы.

  • Следующий сигнал INT запустит следующий этап завершения работы.

Изменено в версии 5.6: В предыдущих версиях Celery при использовании пула prefork во время мягкого завершения работы брокеру не отправлялись сигналы проверки связи. Из-за этого брокер разрывал соединение, и задачи не могли завершиться. Начиная с версии 5.6, при использовании пула prefork сигналы проверки связи продолжают отправляться во время мягкого завершения работы, поэтому задачи могут завершиться до остановки воркера.

Холодное завершение работы

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

Примечание

Если переменной окружения REMAP_SIGTERM присвоено значение SIGQUIT, воркер также начнёт холодное завершение работы при получении сигнала TERM вместо мягкого завершения.

При холодном завершении работы вызов WorkController.start() будет остановлен, а затем будет вызван WorkController.terminate().

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

Постепенное завершение работы

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

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

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

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

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

Например, если задать worker_soft_shutdown_timeout=3, воркер предоставит 3 секунды на завершение всех выполняющихся задач. Если временной лимит будет достигнут, воркер начнёт холодное завершение работы и отменит все выполняющиеся задачи.

[INFO/MainProcess] Task myapp.long_running_task[6f748357-b2c7-456a-95de-f05c00504042] received
[WARNING/ForkPoolWorker-8] long_running_task is running, sleeping 1/2000s
[WARNING/ForkPoolWorker-8] long_running_task is running, sleeping 2/2000s
[WARNING/ForkPoolWorker-8] long_running_task is running, sleeping 3/2000s
^C
worker: Hitting Ctrl+C again will initiate cold shutdown, terminating all running tasks!

worker: Warm shutdown (MainProcess)
[WARNING/ForkPoolWorker-8] long_running_task is running, sleeping 4/2000s
[WARNING/ForkPoolWorker-8] long_running_task is running, sleeping 5/2000s
[WARNING/ForkPoolWorker-8] long_running_task is running, sleeping 6/2000s
^C
worker: Hitting Ctrl+C again will terminate all running tasks!
[WARNING/MainProcess] Initiating Soft Shutdown, terminating in 3 seconds
[WARNING/ForkPoolWorker-8] long_running_task is running, sleeping 7/2000s
[WARNING/ForkPoolWorker-8] long_running_task is running, sleeping 8/2000s
[WARNING/ForkPoolWorker-8] long_running_task is running, sleeping 9/2000s
[WARNING/MainProcess] Restoring 1 unacknowledged message(s)
  • Следующий сигнал QUIT отменит задачи, которые всё ещё выполняются во время постепенного завершения, но воркер продолжит ждать истечения временного лимита, прежде чем завершить работу.

  • Следующий (второй) сигнал QUIT или INT запустит следующий этап завершения работы.

Принудительное завершение работы

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

Принудительное завершение работы предназначено главным образом для локальной разработки или отладки. Оно позволяет многократно отправлять сигнал INT (Ctrl-C), чтобы немедленно завершить воркер. Воркер остановит все выполняющиеся задачи и немедленно завершит работу, вызвав исключение WorkerTerminate в MainProcess.

Обратите внимание на ^C в приведённых ниже журналах (для перехода от одного этапа к другому используется сигнал INT):

[INFO/MainProcess] Task myapp.long_running_task[7235ac16-543d-4fd5-a9e1-2d2bb8ab630a] received
[WARNING/ForkPoolWorker-8] long_running_task is running, sleeping 1/2000s
[WARNING/ForkPoolWorker-8] long_running_task is running, sleeping 2/2000s
^C
worker: Hitting Ctrl+C again will initiate cold shutdown, terminating all running tasks!

worker: Warm shutdown (MainProcess)
[WARNING/ForkPoolWorker-8] long_running_task is running, sleeping 3/2000s
[WARNING/ForkPoolWorker-8] long_running_task is running, sleeping 4/2000s
^C
worker: Hitting Ctrl+C again will terminate all running tasks!
[WARNING/MainProcess] Initiating Soft Shutdown, terminating in 10 seconds
[WARNING/ForkPoolWorker-8] long_running_task is running, sleeping 5/2000s
[WARNING/ForkPoolWorker-8] long_running_task is running, sleeping 6/2000s
^C
Waiting gracefully for cold shutdown to complete...

worker: Cold shutdown (MainProcess)
^C[WARNING/MainProcess] Restoring 1 unacknowledged message(s)

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

Запись журнала Restoring 1 unacknowledged message(s) вводит в заблуждение: нет гарантии, что сообщение будет восстановлено после принудительного завершения работы. Постепенное завершение работы позволяет добавить временное окно между мягким и холодным завершением, чтобы сделать процесс завершения работы более корректным.

Чтобы перезапустить воркер, отправьте сигнал TERM и запустите новый экземпляр. Самый простой способ управлять воркерами во время разработки — использовать celery multi:

$ celerymultistart1-Aproj-lINFO-c4--pidfile=/var/run/celery/%n.pid
$ celerymultirestart1--pidfile=/var/run/celery/%n.pid

В производственной среде следует использовать init-скрипты или систему управления процессами (см. раздел Запуск в качестве демона).

Помимо перезапуска воркера путём остановки и последующего запуска, можно также перезапустить его с помощью сигнала HUP. Учтите, что в этом случае воркер должен перезапустить себя самостоятельно, что может привести к проблемам. Поэтому такой способ не рекомендуется для производственной среды:

$ kill-HUP$pid

Примечание

Перезапуск с помощью HUP работает, только если воркер запущен в фоновом режиме в качестве демона (у него нет управляющего терминала).

HUP отключён в macOS из-за ограничений этой платформы.

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

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

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

Если для параметра worker_cancel_long_running_tasks_on_connection_loss задано значение True, Celery также отменит выполняющиеся в данный момент длительные задачи.

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

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

Эта функция включена по умолчанию, но её можно отключить, задав False для параметра worker_enable_prefetch_count_reduction.

Основной процесс воркера переопределяет следующие сигналы:

Аргументы путей к файлам для --logfile, --pidfile и --statedb могут содержать переменные, которые воркер подставит:

Подстановки имени узла

  • %p: полное имя узла.

  • %h: имя хоста, включая доменное имя.

  • %n: только имя хоста.

  • %d: только доменное имя.

  • %i: индекс процесса пула prefork или 0 для MainProcess.

  • %I: индекс процесса пула prefork с разделителем.

Например, если текущее имя хоста — george@foo.example.com, подстановка даст следующие значения:

  • --logfile=%p.log -> george@foo.example.com.log

  • --logfile=%h.log -> foo.example.com.log

  • --logfile=%n.log -> george.log

  • --logfile=%d.log -> example.com.log

Индекс процесса пула prefork

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

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

Обратите внимание: номера останутся в пределах лимита процессов, даже если процессы завершаются или используются автомасштабирование, maxtasksperchild или временные лимиты. Иными словами, это индекс процесса, а не количество процессов и не pid.

  • %i - индекс процесса пула или 0 для MainProcess.

    В этом случае -n worker1@example.com -c2 -f %n-%i.log создаст три файла журнала:

    • worker1-0.log (основной процесс)

    • worker1-1.log (процесс пула 1)

    • worker1-2.log (процесс пула 2)

  • %I - индекс процесса пула с разделителем.

    В этом случае -n worker1@example.com -c2 -f %n%I.log создаст три файла журнала:

    • worker1.log (основной процесс)

    • worker1-1.log (процесс пула 1)

    • worker1-2.log (процесс пула 2)

По умолчанию для параллельного выполнения задач используется multiprocessing, но также можно использовать Eventlet. Количество процессов/потоков воркера можно изменить с помощью аргумента --concurrency. По умолчанию оно равно количеству доступных на машине процессоров.

Количество процессов (пул multiprocessing/prefork)

Как правило, чем больше процессов в пуле, тем лучше, однако существует предел, после которого увеличение их количества отрицательно сказывается на производительности. Есть даже некоторые свидетельства того, что несколько запущенных экземпляров воркера могут работать производительнее, чем один. Например, можно запустить 3 воркера с 10 процессами в пуле каждого. Вам придётся поэкспериментировать и подобрать оптимальные значения, поскольку они зависят от приложения, рабочей нагрузки, времени выполнения задач и других факторов.

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

Команда celery

Программа celery используется для выполнения команд удалённого управления из командной строки. Она поддерживает все перечисленные ниже команды. Дополнительную информацию см. в разделе Утилиты командной строки для управления (inspect/control).

поддержка пула:

prefork, eventlet, gevent, thread, блокирующий:solo (см. примечание)

поддержка брокеров:

amqp, redis

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

Для команд также могут отправляться ответы. Клиент может дождаться ответов и собрать их. Поскольку в кластере нет централизованного источника, который знал бы, сколько доступно воркеров, невозможно оценить, сколько воркеров могут отправить ответ. Поэтому клиенту можно задать тайм-аут — срок в секундах, за который должны поступить ответы. По умолчанию тайм-аут составляет одну секунду. Если воркер не ответил до истечения срока, это не обязательно означает, что он не ответил или, что ещё хуже, не работает. Причиной может быть задержка сети или медленная обработка команд воркером, поэтому при необходимости измените тайм-аут.

Кроме тайм-аута, клиент может указать максимальное количество ожидаемых ответов. Если задан адресат, этот лимит равен количеству целевых хостов.

Примечание

Пул solo поддерживает команды удалённого управления, но выполняющаяся задача блокирует все ожидающие команды управления. Поэтому эта функция мало полезна, если воркер сильно загружен. В этом случае необходимо увеличить тайм-аут ожидания ответов на стороне клиента.

Функция broadcast()

Эта функция клиента используется для отправки команд воркерам. Для некоторых команд удалённого управления также доступны интерфейсы более высокого уровня, использующие в фоновом режиме broadcast(), например rate_limit() и ping().

Отправка команды rate_limit и именованных аргументов:

>>> app.control.broadcast('rate_limit',
...                          arguments={'task_name': 'myapp.mytask',
...                                     'rate_limit': '200/m'})

Команда будет отправлена асинхронно, без ожидания ответа. Чтобы запросить ответ, используйте аргумент reply:

>>> app.control.broadcast('rate_limit', {
...     'task_name': 'myapp.mytask', 'rate_limit': '200/m'}, reply=True)
[{'worker1.example.com': 'New rate limit set successfully'},
 {'worker2.example.com': 'New rate limit set successfully'},
 {'worker3.example.com': 'New rate limit set successfully'}]

С помощью аргумента destination можно указать список воркеров, которым будет отправлена команда:

>>> app.control.broadcast('rate_limit', {
...     'task_name': 'myapp.mytask',
...     'rate_limit': '200/m'}, reply=True,
...                             destination=['worker1@example.com'])
[{'worker1.example.com': 'New rate limit set successfully'}]

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

revoke: отзыв задач

поддержка пула:

все; terminate поддерживается только в prefork, eventlet и gevent

поддержка брокеров:

amqp, redis

команда:

celery -A proj control revoke <task_id>

Все узлы воркеров хранят сведения об отозванных идентификаторах задач — в памяти или на диске (см. раздел Постоянное хранение сведений об отзыве).

Примечание

Максимальное количество отозванных задач, сведения о которых хранятся в памяти, можно задать с помощью переменной окружения CELERY_WORKER_REVOKES_MAX. Значение по умолчанию — 50000. После превышения лимита сведения об отзыве будут оставаться активными в течение 10800 секунд (3 часов), а затем удалятся. Это значение можно изменить с помощью переменной окружения CELERY_WORKER_REVOKE_EXPIRES.

Ограничения памяти для успешно выполненных задач также можно задать с помощью переменных окружения CELERY_WORKER_SUCCESSFUL_MAX и CELERY_WORKER_SUCCESSFUL_EXPIRES. По умолчанию их значения равны 1000 и 10800 соответственно.

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

Примечание

Параметр terminate — крайняя мера для администраторов, если задача зависла. Он предназначен не для завершения задачи, а для завершения процесса, который её выполняет. К моменту отправки сигнала этот процесс уже может начать обрабатывать другую задачу, поэтому никогда не вызывайте этот параметр программно.

Если задан параметр terminate, дочерний процесс воркера, выполняющий задачу, будет завершён. По умолчанию отправляется сигнал TERM, но его можно указать с помощью аргумента signal. В качестве сигнала можно указать имя в верхнем регистре любого сигнала, определённого в модуле signal стандартной библиотеки Python.

Завершение задачи также отзывает её.

Изменено в версии 5.6: При отзыве задачи сервер результатов теперь немедленно обновляется, чтобы отразить состояние REVOKED. Раньше сервер обновлялся только при попытке воркера обработать отозванную задачу. Из-за этого задачи с ETA/обратным отсчётом могли бесконечно оставаться в состоянии PENDING, если воркер завершал работу до запланированного времени.

Пример

>>> result.revoke()

>>> AsyncResult(id).revoke()

>>> app.control.revoke('d9078da5-9915-40a0-bfa1-392c7bde42ed')

>>> app.control.revoke('d9078da5-9915-40a0-bfa1-392c7bde42ed',
...                    terminate=True)

>>> app.control.revoke('d9078da5-9915-40a0-bfa1-392c7bde42ed',
...                    terminate=True, signal='SIGKILL')

Отзыв нескольких задач

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

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

Пример

>>> app.control.revoke([
...    '7993b0aa-1f0b-4780-9af0-c47c0858b3f2',
...    'f565793e-b041-4b2b-9ca4-dca22762a55d',
...    'd9d35e03-2997-42d0-a13e-64a66b88a618',
])

Начиная с версии 3.1, метод GroupResult.revoke использует эту возможность.

Постоянное хранение сведений об отзыве

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

Список отозванных задач хранится в памяти, поэтому при перезапуске всех воркеров список отозванных идентификаторов исчезнет. Чтобы сохранять этот список между перезапусками, укажите файл для его хранения с помощью аргумента –statedb команды celery worker:

$ celery-Aprojworker-lINFO--statedb=/var/run/celery/worker.state

Если вы используете celery multi, создайте отдельный файл для каждого экземпляра воркера и используйте формат %n для подстановки имени текущего узла:

celery multi start 2 -l INFO --statedb=/var/run/celery/%n.state

См. также раздел Переменные в путях к файлам

Для работы отзыва должны функционировать команды удалённого управления. На данный момент их поддерживают только RabbitMQ (amqp) и Redis.

revoke_by_stamped_headers: отзыв задач по помеченным заголовкам

поддержка пула:

все; terminate поддерживается только в prefork и eventlet

поддержка брокеров:

amqp, redis

команда:

celery -A proj control revoke_by_stamped_headers <header=value>

Эта команда похожа на revoke(), но вместо идентификаторов задач указываются помеченные заголовки в виде пар «ключ-значение». Будет отозвана каждая задача, у которой помеченный заголовок совпадает с указанной парой «ключ-значение».

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

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

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

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

Пример

>>> app.control.revoke_by_stamped_headers({'header': 'value'})

>>> app.control.revoke_by_stamped_headers({'header': 'value'}, terminate=True)

>>> app.control.revoke_by_stamped_headers({'header': 'value'}, terminate=True, signal='SIGKILL')

Отзыв нескольких задач по помеченным заголовкам

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

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

Пример

>> app.control.revoke_by_stamped_headers({
...    'header_A': 'value_1',
...    'header_B': ['value_2', 'value_3'],
})

Будут отозваны все задачи с помеченным заголовком header_A со значением value_1, а также все задачи с помеченным заголовком header_B со значениями value_2 или value_3.

Пример для CLI

$ celery-Aprojcontrolrevoke_by_stamped_headersstamped_header_key_A=stamped_header_value_1stamped_header_key_B=stamped_header_value_2

$ celery-Aprojcontrolrevoke_by_stamped_headersstamped_header_key_A=stamped_header_value_1stamped_header_key_B=stamped_header_value_2--terminate

$ celery-Aprojcontrolrevoke_by_stamped_headersstamped_header_key_A=stamped_header_value_1stamped_header_key_B=stamped_header_value_2--terminate--signal=SIGKILL

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

поддержка пула:

prefork/gevent (см. примечание ниже)

Мягкий или жёсткий?

Ограничение времени задаётся двумя значениями: мягким и жёстким. Мягкое ограничение времени позволяет задаче перехватить исключение и выполнить очистку перед завершением: жёсткое ограничение времени перехватить нельзя, и оно принудительно завершает задачу.

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

Ограничение времени (–time-limit) — это максимальное количество секунд, в течение которого задача может выполняться, прежде чем исполняющий её процесс будет завершён и заменён новым. Вы также можете включить мягкое ограничение времени (–soft-time-limit): оно вызывает исключение, которое задача может перехватить, чтобы выполнить очистку до того, как жёсткое ограничение времени завершит её:

from myapp import app
from celery.exceptions import SoftTimeLimitExceeded

@app.task
def mytask():
    try:
        do_work()
    except SoftTimeLimitExceeded:
        clean_up_in_a_hurry()

Ограничения времени также можно задать с помощью параметров task_time_limit / task_soft_time_limit. Вы также можете указать ограничения времени для клиентских операций, используя аргумент timeout функции AsyncResult.get().

Примечание

Ограничения времени в настоящее время не работают на платформах, не поддерживающих сигнал SIGUSR1.

Примечание

Пул gevent не реализует мягкие ограничения времени. Кроме того, он не обеспечивает соблюдение жёсткого ограничения времени, если задача блокируется.

Изменение ограничений времени во время выполнения

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

поддержка брокеров:

amqp, redis

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

Пример изменения ограничения времени для задачи tasks.crawl_the_web: мягкое ограничение — одна минута, жёсткое — две минуты:

>>> app.control.time_limit('tasks.crawl_the_web',
                           soft=60, hard=120, reply=True)
[{'worker1.example.com': {'ok': 'time limits set successfully'}}]

Изменение повлияет только на задачи, выполнение которых начнётся после изменения ограничения времени.

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

Пример изменения ограничения частоты выполнения для задачи myapp.mytask: не более 200 задач этого типа в минуту:

>>> app.control.rate_limit('myapp.mytask', '200/m')

В примере выше не указано назначение, поэтому запрос на изменение затронет все экземпляры рабочих процессов в кластере. Если нужно затронуть только определённый список рабочих процессов, можно указать аргумент destination:

>>> app.control.rate_limit('myapp.mytask', '200/m',
...            destination=['celery@worker1.example.com'])

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

Это не повлияет на рабочие процессы, в которых включён параметр worker_disable_rate_limits.

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

поддержка пула:

prefork

С помощью этого параметра можно настроить максимальное количество задач, которое рабочий процесс может выполнить до замены новым процессом.

Это полезно, например, если возникают утечки памяти, которые вы не можете контролировать, в том числе в закрытых расширениях на C.

Параметр можно задать с помощью аргумента --max-tasks-per-child команды worker или с помощью параметра worker_max_tasks_per_child.

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

поддержка пула:

prefork

С помощью этого параметра можно настроить максимальный объём резидентной памяти, который может использовать рабочий процесс до замены новым процессом.

Это полезно, например, если возникают утечки памяти, которые вы не можете контролировать, в том числе в закрытых расширениях на C.

Параметр можно задать с помощью аргумента --max-memory-per-child команды worker или с помощью параметра worker_max_memory_per_child.

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

поддержка пула:

prefork, gevent

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

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

Автомасштабирование включается параметром --autoscale, которому требуются два числа: максимальное и минимальное количество процессов в пуле:

--autoscale=AUTOSCALE
     Enable autoscaling by providing
     max_concurrency,min_concurrency.  Example:
       --autoscale=10,3 (always keep 3 processes, but grow to
      10 if necessary).

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

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

При запуске можно указать очереди, из которых нужно получать сообщения, передав параметру -Q список очередей, разделённых запятыми:

$ celery-Aprojworker-lINFO-Qfoo,bar,baz

Если имя очереди указано в task_queues, будет использована её конфигурация. Если же очередь не указана в списке, Celery автоматически создаст для вас новую очередь (в зависимости от параметра task_create_missing_queues).

Вы также можете указать рабочему процессу начать или прекратить получать сообщения из очереди во время выполнения с помощью команд удалённого управления add_consumer и cancel_consumer.

Очереди: добавление получателей

Команда управления add_consumer указывает одному или нескольким рабочим процессам начать получать сообщения из очереди. Эта операция идемпотентна.

Чтобы указать всем рабочим процессам кластера начать получать сообщения из очереди с именем «foo», используйте программу celery control:

$ celery-Aprojcontroladd_consumerfoo
-> worker1.local: OK
    started consuming from u'foo'

Чтобы указать конкретный рабочий процесс, используйте аргумент --destination:

$ celery-Aprojcontroladd_consumerfoo-dcelery@worker1.local

То же самое можно выполнить динамически с помощью метода app.control.add_consumer():

>>> app.control.add_consumer('foo', reply=True)
[{u'worker1.local': {u'ok': u"already consuming from u'foo'"}}]

>>> app.control.add_consumer('foo', reply=True,
...                          destination=['worker1@example.com'])
[{u'worker1.local': {u'ok': u"already consuming from u'foo'"}}]

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

>>> app.control.add_consumer(
...     queue='baz',
...     exchange='ex',
...     exchange_type='topic',
...     routing_key='media.*',
...     options={
...         'queue_durable': False,
...         'exchange_durable': False,
...     },
...     reply=True,
...     destination=['w1@example.com', 'w2@example.com'])

Очереди: отмена получателей

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

Чтобы указать всем рабочим процессам кластера прекратить получать сообщения из очереди, используйте программу celery control:

$ celery-Aprojcontrolcancel_consumerfoo

Аргумент --destination позволяет указать рабочий процесс или список рабочих процессов, к которым будет применена команда:

$ celery-Aprojcontrolcancel_consumerfoo-dcelery@worker1.local

Также можно программно отменить получение сообщений с помощью метода app.control.cancel_consumer():

>>> app.control.cancel_consumer('foo', reply=True)
[{u'worker1.local': {u'ok': u"no longer consuming from u'foo'"}}]

Очереди: список активных очередей

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

$ celery-Aprojinspectactive_queues
[...]

Как и все остальные команды удалённого управления, эта команда поддерживает аргумент --destination, который позволяет указать рабочие процессы, от которых нужно получить ответ:

$ celery-Aprojinspectactive_queues-dcelery@worker1.local
[...]

То же самое можно сделать программно с помощью метода active_queues():

>>> app.control.inspect().active_queues()
[...]

>>> app.control.inspect(['worker1.local']).active_queues()
[...]

app.control.inspect позволяет проверять работающие процессы. Внутри он использует команды удалённого управления.

Для проверки рабочих процессов также можно использовать команду celery. Она поддерживает те же команды, что и интерфейс app.control.

>>> # Inspect all nodes.
>>> i = app.control.inspect()

>>> # Specify multiple nodes to inspect.
>>> i = app.control.inspect(['worker1.example.com',
                            'worker2.example.com'])

>>> # Specify a single node to inspect.
>>> i = app.control.inspect('worker1.example.com')

Список зарегистрированных задач

Получить список задач, зарегистрированных в рабочем процессе, можно с помощью registered():

>>> i.registered()
[{'worker1.example.com': ['tasks.add',
                          'tasks.sleeptask']}]

Список выполняющихся задач

Получить список активных задач можно с помощью active():

>>> i.active()
[{'worker1.example.com':
    [{'name': 'tasks.sleeptask',
      'id': '32666e9b-809c-41fa-8e93-5ae0c80afbbf',
      'args': '(8,)',
      'kwargs': '{}'}]}]

Список запланированных задач (ETA)

Получить список задач, ожидающих запланированного выполнения, можно с помощью scheduled():

>>> i.scheduled()
[{'worker1.example.com':
    [{'eta': '2010-06-07 09:07:52', 'priority': 0,
      'request': {
        'name': 'tasks.sleeptask',
        'id': '1a7980ea-8b19-413e-91d2-0b74f3844c4d',
        'args': '[1]',
        'kwargs': '{}'}},
     {'eta': '2010-06-07 09:07:53', 'priority': 0,
      'request': {
        'name': 'tasks.sleeptask',
        'id': '49661b9a-aa22-4120-94b7-9ee8031d219d',
        'args': '[2]',
        'kwargs': '{}'}}]}]

Примечание

Это задачи с аргументом ETA/countdown, а не периодические задачи.

Список зарезервированных задач

Зарезервированные задачи — это задачи, которые уже получены, но ещё ожидают выполнения.

Получить их список можно с помощью reserved():

>>> i.reserved()
[{'worker1.example.com':
    [{'name': 'tasks.sleeptask',
      'id': '32666e9b-809c-41fa-8e93-5ae0c80afbbf',
      'args': '(8,)',
      'kwargs': '{}'}]}]

Статистика

Команда удалённого управления inspect stats (или stats()) выводит длинный список полезных (или не очень полезных) статистических данных о рабочем процессе:

$ celery-Aprojinspectstats

Подробное описание выходных данных см. в справочной документации по stats().

Удалённое завершение работы

Эта команда корректно завершит работу удалённого рабочего процесса:

>>> app.control.broadcast('shutdown') # shutdown all workers
>>> app.control.broadcast('shutdown', destination='worker1@example.com')

Проверка связи

Эта команда отправляет запрос проверки связи работающим процессам. В ответ они отправляют строку «pong» — и на этом всё. Если не указать собственное время ожидания, для ответа будет использоваться значение по умолчанию — одна секунда:

>>> app.control.ping(timeout=0.5)
[{'worker1.example.com': 'pong'},
 {'worker2.example.com': 'pong'},
 {'worker3.example.com': 'pong'}]

ping() также поддерживает аргумент destination, позволяющий указать рабочие процессы для проверки связи:

>>> ping(['worker2.example.com', 'worker3.example.com'])
[{'worker2.example.com': 'pong'},
 {'worker3.example.com': 'pong'}]

Включение и отключение событий

События можно включать и отключать с помощью команд enable_events и disable_events. Это удобно для временного мониторинга рабочего процесса с помощью celery events/celerymon.

>>> app.control.enable_events()
>>> app.control.disable_events()

Существует два типа команд удалённого управления:

  • Команда инспектирования

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

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

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

Команды удалённого управления регистрируются в панели управления и принимают один аргумент: текущий экземпляр celery.worker.control.ControlDispatch. При необходимости через него можно получить доступ к активному объекту Consumer.

Ниже приведён пример команды управления, увеличивающей количество предварительно получаемых задач:

from celery.worker.control import control_command

@control_command(
    args=[('n', int)],
    signature='[N=1]',  # <- used for help on the command-line.
)
def increase_prefetch_count(state, n=1):
    state.consumer.qos.increment_eventually(n)
    return {'ok': 'prefetch count incremented'}

Добавьте этот код в модуль, импортируемый рабочим процессом: это может быть тот же модуль, в котором определено ваше приложение Celery, или модуль, добавленный в параметр imports.

Перезапустите рабочий процесс, чтобы зарегистрировать команду управления. После этого её можно вызвать с помощью утилиты celery control:

$ celery-Aprojcontrolincrease_prefetch_count3

В программу celery inspect также можно добавить действия, например действие для чтения текущего количества предварительно получаемых задач:

from celery.worker.control import inspect_command

@inspect_command()
def current_prefetch_count(state):
    return {'prefetch_count': state.consumer.qos.value}

После перезапуска рабочего процесса это значение можно запросить с помощью программы celery inspect:

$ celery-Aprojinspectcurrent_prefetch_count

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/workers.html

Spec-Zone.ru

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