Spec-Zone.ru › Celery

Задачи

Задачи — это строительные блоки приложений Celery.

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

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

Сообщение задачи не удаляется из очереди, пока рабочий процесс не подтвердит его получение (acknowledged). Рабочий процесс может заранее зарезервировать много сообщений, и даже если его остановить — из-за сбоя питания или по какой-либо другой причине, — сообщение будет доставлено повторно другому рабочему процессу.

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

Если ваша задача идемпотентна, можно установить параметр acks_late, чтобы рабочий процесс подтверждал сообщение после завершения задачи. См. также запись в разделе часто задаваемых вопросов: Использовать retry или acks_late?.

Обратите внимание: рабочий процесс подтвердит сообщение, если дочерний процесс, выполняющий задачу, будет завершён (либо в результате вызова задачей sys.exit(), либо сигналом), даже если включён параметр acks_late. Такое поведение предусмотрено намеренно, поскольку…

  1. Мы не хотим повторно запускать задачи, из-за которых ядро отправляет процессу сигнал SIGSEGV (ошибка сегментации) или похожие сигналы.

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

  3. Задача, выделяющая слишком много памяти, может вызвать срабатывание механизма OOM killer ядра; то же самое может произойти снова.

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

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

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

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

Если ваша задача выполняет операции ввода-вывода, обязательно задайте для них тайм-ауты. Например, можно задать тайм-аут для веб-запроса с помощью библиотеки https://pypi.org/project/requests/:

connect_timeout, read_timeout = 5.0, 30.0
response = requests.get(URL, timeout=(connect_timeout, read_timeout))

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

В предыдущих версиях планировщик пула prefork по умолчанию был не приспособлен к длительно выполняющимся задачам, поэтому для задач, выполнявшихся минутами или часами, рекомендовалось включить аргумент командной строки -Ofair для команды celery worker. Однако начиная с версии 4.0 стратегия планирования -Ofair используется по умолчанию. Дополнительную информацию см. в разделе Ограничения предварительной выборки. Для максимальной производительности направляйте длительно и кратковременно выполняющиеся задачи в отдельные рабочие процессы (Автоматическая маршрутизация).

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

–

В этой главе вы узнаете всё о создании задач. Ниже приведено содержание:

Задачу можно легко создать на основе любой вызываемой сущности с помощью декоратора app.task():

from .models import User

@app.task
def create_user(username, password):
    User.objects.create(username=username, password=password)

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

@app.task(serializer='json')
def create_user(username, password):
    User.objects.create(username=username, password=password)

Как импортировать декоратор задачи?

Декоратор задачи доступен в экземпляре приложения Celery. Если вы не знаете, что это такое, прочитайте раздел Первые шаги с Celery.

Если вы используете Django (см. раздел Первые шаги с Django) или являетесь автором библиотеки, вероятно, вам следует использовать декоратор shared_task():

from celery import shared_task

@shared_task
def add(x, y):
    return x + y

Несколько декораторов

При совместном использовании нескольких декораторов с декоратором задачи убедитесь, что декоратор task применяется последним (как ни странно, в Python это означает, что он должен стоять первым в списке):

@app.task
@decorator2
@decorator1
def add(x, y):
    return x + y

Связанные задачи

Связанная задача означает, что первым аргументом задачи всегда будет экземпляр самой задачи (self), как и в связанных методах Python:

logger = get_task_logger(__name__)

@app.task(bind=True)
def add(self, x, y):
    logger.info(self.request.id)

Связанные задачи нужны для повторных попыток (с использованием app.Task.retry()), доступа к сведениям о текущем запросе задачи и любой дополнительной функциональности, добавленной вами в пользовательские базовые классы задач.

Наследование задач

Аргумент base декоратора задачи задаёт базовый класс задачи:

import celery

class MyTask(celery.Task):

    def on_failure(self, exc, task_id, args, kwargs, einfo):
        print('{0!r} failed: {1!r}'.format(task_id, exc))

@app.task(base=MyTask)
def add(x, y):
    raise KeyError()

У каждой задачи должно быть уникальное имя.

Если явное имя не задано, декоратор задачи сгенерирует его автоматически на основе 1) модуля, в котором определена задача, и 2) имени функции задачи.

Пример задания явного имени:

>>> @app.task(name='sum-of-two-numbers')
>>> def add(x, y):
...     return x + y

>>> add.name
'sum-of-two-numbers'

Рекомендуется использовать имя модуля в качестве пространства имён. Так имена не будут конфликтовать, если задача с таким же именем уже определена в другом модуле.

>>> @app.task(name='tasks.add')
>>> def add(x, y):
...     return x + y

Имя задачи можно узнать, проверив её атрибут .name:

>>> add.name
'tasks.add'

Указанное здесь имя (tasks.add) — это именно то имя, которое было бы сгенерировано автоматически, если бы задача была определена в модуле с именем tasks.py:

tasks.py:

@app.task
def add(x, y):
    return x + y
>>> from tasks import add
>>> add.name
'tasks.add'

Примечание

Чтобы просмотреть имена всех зарегистрированных задач, можно использовать команду inspect в рабочем процессе. См. команду inspect registered в разделе Утилиты командной строки для управления (inspect/control) Руководства пользователя.

Изменение поведения автоматического именования

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

В некоторых случаях автоматическое именование по умолчанию не подходит. Представьте, что у вас есть множество задач в разных модулях:

project/
       /__init__.py
       /celery.py
       /moduleA/
               /__init__.py
               /tasks.py
       /moduleB/
               /__init__.py
               /tasks.py

При использовании автоматического именования по умолчанию каждой задаче будет присвоено имя, например moduleA.tasks.taskA, moduleA.tasks.taskB, moduleB.tasks.test и так далее. Возможно, вы захотите убрать tasks из имён всех задач. Как отмечалось выше, можно явно задать имена для всех задач или изменить поведение автоматического именования, переопределив app.gen_task_name(). В продолжение примера в файле celery.py может быть следующий код:

from celery import Celery

class MyCelery(Celery):

    def gen_task_name(self, name, module):
        if module.endswith('.tasks'):
            module = module[:-6]
        return super().gen_task_name(name, module)

app = MyCelery('main')

Таким образом, имена задач будут выглядеть, например, так: moduleA.taskA, moduleA.taskB и moduleB.test.

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

Убедитесь, что app.gen_task_name() — чистая функция: для одного и того же входного значения она всегда должна возвращать один и тот же результат.

app.Task.request содержит информацию и состояние, связанные с выполняемой в данный момент задачей.

Запрос содержит следующие атрибуты:

id:

Уникальный идентификатор выполняемой задачи.

group:

Уникальный идентификатор группы задачи, если она входит в группу.

chord:

Уникальный идентификатор аккорда, которому принадлежит задача (если задача входит в заголовок).

correlation_id:

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

args:

Позиционные аргументы.

kwargs:

Именованные аргументы.

origin:

Имя узла, отправившего эту задачу.

retries:

Количество повторных попыток выполнения текущей задачи. Целое число, начинающееся с 0.

is_eager:

Имеет значение True, если задача выполняется локально на клиенте, а не рабочим процессом.

eta:

Исходное время ETA задачи (если задано). Указано в формате UTC (в зависимости от параметра enable_utc).

expires:

Исходное время истечения срока действия задачи (если задано). Указано в формате UTC (в зависимости от параметра enable_utc).

hostname:

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

delivery_info:

Дополнительная информация о доставке сообщения. Это отображение, содержащее обменник и ключ маршрутизации, использованные для доставки задачи. Например, оно используется методом app.Task.retry(), чтобы повторно отправить задачу в ту же целевую очередь. Доступность ключей в этом словаре зависит от используемого брокера сообщений.

reply-to:

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

called_directly:

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

timelimit:

Кортеж текущих активных ограничений времени (soft, hard) для этой задачи (если они есть).

callbacks:

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

errbacks:

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

utc:

Имеет значение true, если для вызывающего кода включён UTC (enable_utc).

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

headers:

Отображение заголовков сообщения, отправленных вместе с сообщением задачи (может быть None).

reply_to:

Адрес для отправки ответа (имя очереди).

correlation_id:

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

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

root_id:

Уникальный идентификатор первой задачи в рабочем процессе, частью которого является эта задача (если такая задача есть).

parent_id:

Уникальный идентификатор задачи, вызвавшей эту задачу (если такая задача есть).

chain:

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

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

properties:

Отображение свойств сообщения, полученных вместе с сообщением задачи (может быть None или {})

replaced_task_nesting:

Количество замен задачи, если они выполнялись (может быть 0)

Пример

Пример задачи, обращающейся к информации в контексте:

@app.task(bind=True)
def dump_context(self, x, y):
    print('Executing task id {0.id}, args: {0.args!r} kwargs: {0.kwargs!r}'.format(
            self.request))

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

Рабочий процесс автоматически настроит для вас журналирование; также можно настроить его вручную.

Доступен специальный регистратор с именем «celery.task». От него можно наследовать регистраторы, чтобы автоматически включать в журналы имя задачи и её уникальный идентификатор.

Рекомендуется создать общий регистратор для всех задач в начале модуля:

from celery.utils.log import get_task_logger

logger = get_task_logger(__name__)

@app.task
def add(x, y):
    logger.info('Adding {0} + {1}'.format(x, y))
    return x + y

Celery использует стандартную библиотеку журналирования Python; документацию можно найти here.

Можно также использовать print(): всё, что записывается в стандартный поток вывода или ошибок, будет перенаправлено в систему журналирования (это можно отключить; см. worker_redirect_stdouts).

Примечание

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

Если вы хотите перенаправлять sys.stdout и sys.stderr в пользовательский регистратор, это нужно включить вручную, например:

import sys

logger = get_task_logger(__name__)

@app.task(bind=True)
def add(self, x, y):
    old_outs = sys.stdout, sys.stderr
    rlevel = self.app.conf.worker_redirect_stdouts_level
    try:
        self.app.log.redirect_stdouts_to_logger(logger, rlevel)
        print('Adding {0} + {1}'.format(x, y))
        return x + y
    finally:
        sys.stdout, sys.stderr = old_outs

Примечание

Если нужный регистратор Celery не выводит записи журнала, проверьте, правильно ли настроена передача записей от него. В этом примере включён «celery.app.trace», поэтому выводятся записи «succeeded in»:

import celery
import logging

@celery.signals.after_setup_logger.connect
def on_after_setup_logger(**kwargs):
    logger = logging.getLogger('celery')
    logger.propagate = True
    logger = logging.getLogger('celery.app.trace')
    logger.propagate = True

Примечание

Чтобы полностью отключить настройку журналирования Celery, используйте сигнал setup_logging:

import celery

@celery.signals.setup_logging.connect
def on_setup_logging(**kwargs):
    pass

Проверка аргументов

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

Celery проверит аргументы, переданные при вызове задачи, так же, как Python проверяет их при вызове обычной функции:

>>> @app.task
... def add(x, y):
...     return x + y

# Calling the task with two arguments works:
>>> add.delay(8, 8)
<AsyncResult: f59d71ca-1549-43e0-be41-4e8821a83c0c>

# Calling the task with only one argument fails:
>>> add.delay(8)
Traceback (most recent call last):
  File "<stdin>", line 1, in <module>
  File "celery/app/task.py", line 376, in delay
    return self.apply_async(args, kwargs)
  File "celery/app/task.py", line 485, in apply_async
    check_arguments(*(args or ()), **(kwargs or {}))
TypeError: add() takes exactly 2 arguments (1 given)

Проверку аргументов для любой задачи можно отключить, задав её атрибуту typing значение False:

>>> @app.task(typing=False)
... def add(x, y):
...     return x + y

# Works locally, but the worker receiving the task will raise an error.
>>> add.delay(8)
<AsyncResult: f59d71ca-1549-43e0-be41-4e8821a83c0c>

Скрытие конфиденциальной информации в аргументах

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

При использовании task_protocol версии 2 или выше (по умолчанию начиная с версии 4.0) можно переопределить представление позиционных и именованных аргументов в журналах и событиях мониторинга с помощью аргументов вызова argsrepr и kwargsrepr:

>>> add.apply_async((2, 3), argsrepr='(<secret-x>, <secret-y>)')

>>> charge.s(account, card='1234 5678 1234 5678').set(
...     kwargsrepr=repr({'card': '**** **** **** 5678'})
... ).delay()

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

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

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

Метод app.Task.retry() можно использовать для повторного выполнения задачи, например при возникновении ошибок, допускающих восстановление.

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

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

Пример использования retry:

@app.task(bind=True)
def send_twitter_status(self, oauth, tweet):
    try:
        twitter = Twitter(oauth)
        twitter.update_status(tweet)
    except (Twitter.FailWhaleError, Twitter.LoginError) as exc:
        raise self.retry(exc=exc)

Примечание

Вызов app.Task.retry() вызовет исключение, поэтому код после повторной попытки не будет выполнен. Это исключение Retry. Оно обрабатывается не как ошибка, а как специальный предикат, сообщающий рабочему процессу о необходимости повторить задачу, чтобы при включённой серверной части результатов можно было сохранить правильное состояние.

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

Аргумент bind декоратора задачи предоставляет доступ к self (экземпляру типа задачи).

Аргумент exc используется для передачи сведений об исключении, которые сохраняются в журналах и при записи результатов задачи. Исключение и трассировка стека будут доступны в состоянии задачи (если включена серверная часть результатов).

Если у задачи задано значение max_retries, при превышении максимального количества повторных попыток будет повторно вызвано текущее исключение, за исключением следующих случаев:

  • Аргумент exc не был указан.

    В этом случае будет вызвано исключение MaxRetriesExceededError.

  • Текущего исключения нет.

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

    self.retry(exc=Twitter.LoginError())
    

    вызовет переданный аргумент exc.

Настройка задержки перед повторной попыткой

Перед повторной попыткой задача может подождать заданное время. Задержка по умолчанию определяется атрибутом default_retry_delay. По умолчанию он равен 3 минутам. Задержка задаётся в секундах (целым числом или числом с плавающей точкой).

Также можно передать аргумент countdown методу retry(), чтобы переопределить значение по умолчанию.

@app.task(bind=True, default_retry_delay=30 * 60)  # retry in 30 minutes.
def add(self, x, y):
    try:
        something_raising()
    except Exception as exc:
        # overrides the default delay to retry after 1 minute
        raise self.retry(exc=exc, countdown=60)

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

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

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

К счастью, можно указать Celery автоматически повторять задачу, используя аргумент autoretry_for в декораторе app.task():

from twitter.exceptions import FailWhaleError

@app.task(autoretry_for=(FailWhaleError,))
def refresh_timeline(user):
    return twitter.refresh_timeline(user)

Чтобы задать пользовательские аргументы для внутреннего вызова retry(), передайте аргумент retry_kwargs декоратору app.task():

@app.task(autoretry_for=(FailWhaleError,),
          retry_kwargs={'max_retries': 5})
def refresh_timeline(user):
    return twitter.refresh_timeline(user)

Это альтернатива ручной обработке исключений. Пример выше работает так же, как если бы тело задачи было обёрнуто в инструкцию try … except:

@app.task
def refresh_timeline(user):
    try:
        twitter.refresh_timeline(user)
    except FailWhaleError as exc:
        raise refresh_timeline.retry(exc=exc, max_retries=5)

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

@app.task(autoretry_for=(Exception,))
def x():
    ...

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

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

from requests.exceptions import RequestException

@app.task(autoretry_for=(RequestException,), retry_backoff=True)
def x():
    ...

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

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

Параметры autoretry_for, max_retries, retry_backoff, retry_backoff_max и retry_jitter можно также задавать в задачах на основе классов:

class BaseTaskWithRetry(Task):
    autoretry_for = (TypeError,)
    max_retries = 5
    retry_backoff = True
    retry_backoff_max = 700
    retry_jitter = False
Task.autoretry_for

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

Task.max_retries

Число. Максимальное количество повторных попыток, после которого выполнение прекращается. Значение None означает, что задача будет повторяться бесконечно. По умолчанию этому параметру присвоено значение 3.

Task.retry_backoff

Логическое значение или число. Если этому параметру присвоено значение True, задержки автоматических повторных попыток будут рассчитываться по правилам экспоненциальной задержки. Первая повторная попытка будет выполнена через 1 секунду, вторая — через 2 секунды, третья — через 4 секунды, четвёртая — через 8 секунд и так далее. (Однако эта задержка изменяется параметром retry_jitter, если он включён.) Если этому параметру присвоено число, оно используется как коэффициент задержки. Например, если задано значение 3, первая повторная попытка будет выполнена через 3 секунды, вторая — через 6 секунд, третья — через 12 секунд, четвёртая — через 24 секунды и так далее. По умолчанию этому параметру присвоено значение False, и автоматические повторные попытки не задерживаются.

Task.retry_backoff_max

Число. Если включён параметр retry_backoff, этот параметр задаёт максимальную задержку в секундах между автоматическими повторными попытками задачи. По умолчанию ему присвоено значение 600, то есть 10 минут.

Task.retry_jitter

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

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

Task.dont_autoretry_for
Список или кортеж классов исключений. Для этих исключений автоматические повторные попытки выполняться не будут.

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

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

Вы можете использовать Pydantic для проверки и преобразования аргументов, а также сериализации результатов на основе подсказок типов, передав pydantic=True.

Примечание

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

Например:

from pydantic import BaseModel

class ArgModel(BaseModel):
    value: int

class ReturnModel(BaseModel):
    value: str

@app.task(pydantic=True)
def x(arg: ArgModel) -> ReturnModel:
    # args/kwargs type hinted as Pydantic model will be converted
    assert isinstance(arg, ArgModel)

    # The returned model will be converted to a dict automatically
    return ReturnModel(value=f"example: {arg.value}")

Затем задачу можно вызвать, передав словарь, соответствующий модели, а в ответ вы получите «выгруженную» модель (сериализованную с помощью BaseModel.model_dump()):

>>> result = x.delay({'value': 1})
>>> result.get(timeout=1)
{'value': 'example: 1'}

Типы объединения, аргументы обобщённых типов

Типы объединения (например, Union[SomeModel, OtherModel]) или аргументы обобщённых типов (например, list[SomeModel]) не поддерживаются.

Если вы хотите поддерживать списки или подобные типы, рекомендуется использовать pydantic.RootModel.

Необязательные параметры и возвращаемые значения

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

from typing import Optional

# models are the same as above

@app.task(pydantic=True)
def x(arg: Optional[ArgModel] = None) -> Optional[ReturnModel]:
    if arg is None:
        return None
    return ReturnModel(value=f"example: {arg.value}")

Поведение будет следующим:

 >>> result = x.delay()
>>> result.get(timeout=1) is None
True
>>> result = x.delay({'value': 1})
>>> result.get(timeout=1)
{'value': 'example: 1'}

Обработка возвращаемых значений

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

Параметры Pydantic

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

Task.pydantic_strict

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

Task.pydantic_context

Передайте дополнительный контекст проверки при проверке модели Pydantic. По умолчанию контекст уже содержит объект приложения под именем celery_app и имя задачи под именем celery_task_name.

Task.pydantic_dump_kwargs

При сериализации результата передайте эти дополнительные аргументы в dump_kwargs(). По умолчанию передаётся только mode='json'.

Декоратор задачи принимает ряд параметров, меняющих поведение задачи. Например, с помощью параметра rate_limit можно задать ограничение частоты выполнения задачи.

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

Общие сведения

Task.name

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

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

См. также Имена.

Task.request

Если задача выполняется, этот атрибут содержит информацию о текущем запросе. Используется локальное хранилище потоков.

См. Запрос задачи.

Task.max_retries

Применяется, только если задача вызывает self.retry или если задача декорирована с аргументом autoretry_for.

Максимальное количество повторных попыток перед завершением. Если количество попыток превысит это значение, будет вызвано исключение MaxRetriesExceededError.

Примечание

Необходимо вызывать retry() вручную, поскольку при возникновении исключения повторная попытка не выполняется автоматически.

Значение по умолчанию — 3. Значение None отключает ограничение числа повторных попыток: задача будет повторяться бесконечно, пока не завершится успешно.

Task.throws

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

Ошибки из этого списка будут переданы в хранилище результатов как неудачное выполнение, но рабочий процесс не запишет это событие как ошибку и не включит трассировку стека.

Пример:

@task(throws=(KeyError, HttpNotFound)):
def get_foo():
    something()

Типы ошибок:

  • Ожидаемые ошибки (в Task.throws)

    Записываются с уровнем важности INFO, трассировка стека не включается.

  • Неожиданные ошибки

    Записываются с уровнем важности ERROR, трассировка стека включается.

Task.default_retry_delay

Время в секундах, которое должно пройти до повторного выполнения задачи по умолчанию. Может иметь тип int или float. По умолчанию задержка составляет три минуты.

Task.rate_limit

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

Если указано None, ограничение частоты не действует. Целое число или число с плавающей точкой интерпретируется как «задач в секунду».

Ограничение частоты можно задать в секундах, минутах или часах, добавив к значению “/s”, “/m” или “/h”. Задачи будут равномерно распределены по указанному промежутку времени.

Пример: “100/m” (сто задач в минуту). Это установит минимальную задержку в 600 мс между запуском двух задач на одном экземпляре рабочего процесса.

По умолчанию используется значение параметра task_default_rate_limit: если оно не задано, ограничение частоты выполнения задач отключено.

Обратите внимание, что это ограничение частоты действует для одного экземпляра рабочего процесса, а не глобально. Чтобы установить глобальное ограничение частоты (например, для API с максимальным количеством запросов в секунду), необходимо привязать задачи к определённой очереди.

Task.time_limit

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

Task.soft_time_limit

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

Task.ignore_result

Не сохранять состояние задачи. Обратите внимание: в этом случае нельзя использовать AsyncResult, чтобы проверить, завершена ли задача, или получить её возвращаемое значение.

Примечание. Некоторые функции не будут работать, если результаты задач отключены. Подробнее см. в документации Canvas.

Task.store_errors_even_if_ignored

Если задано значение True, ошибки будут сохраняться, даже если для задачи настроено игнорирование результатов.

Изменено в версии 5.7: Ранее, если в сообщении запроса отсутствовал ключ ignore_result, для store_errors по умолчанию устанавливалось значение True, и собственная настройка задачи ignore_result игнорировалась. Теперь, если переопределение для запроса отсутствует, рабочий процесс корректно использует значение Task.ignore_result.

Task.serializer

Строка, задающая используемый по умолчанию метод сериализации. По умолчанию используется значение параметра task_serializer. Возможные значения: pickle, json, yaml или любой пользовательский метод сериализации, зарегистрированный с помощью kombu.serialization.registry.

Подробнее см. раздел Сериализаторы.

Task.compression

Строка, задающая используемую по умолчанию схему сжатия.

По умолчанию используется значение параметра task_compression. Возможные значения: gzip, bzip2 или любая пользовательская схема сжатия, зарегистрированная в реестре kombu.compression.

Подробнее см. раздел Сжатие.

Task.backend

Хранилище результатов, используемое для этой задачи. Экземпляр одного из классов хранилищ из celery.backends. По умолчанию используется app.backend, определяемый параметром result_backend.

Task.acks_late

Если задано значение True, сообщения для этой задачи подтверждаются после выполнения задачи, а не непосредственно перед ним (поведение по умолчанию).

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

Глобальное значение по умолчанию можно переопределить с помощью параметра task_acks_late.

Task.track_started

Если задано значение True, при выполнении задачи рабочим процессом её состояние будет отображаться как «запущена». По умолчанию используется значение False, поскольку обычно такая степень детализации не требуется. Задачи либо ожидают выполнения, либо завершены, либо ожидают повторной попытки. Состояние «запущена» может быть полезно для длительных задач, когда необходимо сообщать, какая задача выполняется в данный момент.

Имя узла и идентификатор процесса рабочего процесса, выполняющего задачу, будут доступны в метаданных состояния (например, result.info[‘pid’])

Глобальное значение по умолчанию можно переопределить с помощью параметра task_track_started.

См. также

Справочник API для Task.

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

На выбор доступны несколько хранилищ результатов, каждое со своими преимуществами и недостатками (см. Хранилища результатов).

За время существования задача переходит через несколько возможных состояний, и к каждому состоянию могут быть добавлены произвольные метаданные. Когда задача переходит в новое состояние, предыдущее состояние забывается, однако некоторые переходы можно вывести (например, из текущего состояния задачи FAILED следует, что в какой-то момент она находилась в состоянии STARTED).

Существуют также наборы состояний, например набор FAILURE_STATES и набор READY_STATES.

Клиент использует принадлежность к этим наборам, чтобы определить, следует ли повторно вызвать исключение (PROPAGATE_STATES) и можно ли кэшировать состояние (это возможно, если задача готова).

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

Хранилища результатов

Если вам нужно отслеживать задачи или получать возвращаемые значения, Celery должна где-либо сохранять состояния или отправлять их, чтобы их можно было получить позднее. Доступны несколько встроенных хранилищ результатов: SQLAlchemy/Django ORM, Memcached, RabbitMQ/QPid (rpc) и Redis. Также можно определить собственное хранилище.

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

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

Хранилища используют ресурсы для хранения и передачи результатов. Чтобы гарантировать освобождение ресурсов, необходимо в конечном итоге вызвать get() или forget() для КАЖДОГО экземпляра AsyncResult, возвращённого после вызова задачи.

См. также

Настройки хранилища результатов задач

Хранилище результатов RPC (RabbitMQ/QPid)

Хранилище результатов RPC (rpc://) отличается тем, что фактически не хранит состояния, а отправляет их в виде сообщений. Это важное отличие означает, что результат можно получить только один раз и только клиенту, инициировавшему задачу. Два разных процесса не могут ожидать один и тот же результат.

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

По умолчанию сообщения являются временными (непостоянными), поэтому результаты исчезнут при перезапуске брокера. Можно настроить хранилище результатов на отправку постоянных сообщений с помощью параметра result_persistent.

Хранилище результатов на основе базы данных

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

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

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

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

    Рекомендуется изменить этот уровень на READ-COMMITTED.

Встроенные состояния

PENDING

Задача ожидает выполнения или её состояние неизвестно. Считается, что любой неизвестный идентификатор задачи соответствует состоянию ожидания.

STARTED

Задача запущена. По умолчанию это состояние не сообщается; чтобы включить его, см. app.Task.track_started.

метаданные:

pid и hostname процесса рабочего процесса, выполняющего задачу.

SUCCESS

Задача успешно выполнена.

метаданные:

result содержит возвращаемое значение задачи.

передаёт исключение дальше:

Да

готово:

Да

FAILURE

Выполнение задачи завершилось неудачей.

метаданные:

result содержит возникшее исключение, а traceback — трассировку стека в момент возникновения исключения.

передаёт исключение дальше:

Да

RETRY

Выполняется повторная попытка запуска задачи.

метаданные:

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

передаёт исключение дальше:

Нет

REVOKED

Задача отозвана.

передаёт исключение дальше:

Да

Пользовательские состояния

Определить собственные состояния легко: достаточно выбрать уникальное имя. Обычно имя состояния представляет собой строку в верхнем регистре. Например, можно ознакомиться с abortable tasks, где определено пользовательское состояние ABORTED.

Используйте update_state(), чтобы обновить состояние задачи:

@app.task(bind=True)
def upload_files(self, filenames):
    for i, file in enumerate(filenames):
        if not self.request.called_directly:
            self.update_state(state='PROGRESS',
                meta={'current': i, 'total': len(filenames)})

Здесь было создано состояние “PROGRESS”, которое сообщает любому приложению, поддерживающему это состояние, что задача выполняется. Положение задачи в процессе выполнения задаётся счётчиками current и total, включёнными в метаданные состояния. Например, на их основе можно создавать индикаторы выполнения.

Создание сериализуемых исключений

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

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

Чтобы исключения можно было сериализовать с помощью pickle, атрибут .args исключения ДОЛЖЕН содержать исходные аргументы, с которыми оно было создано. Проще всего обеспечить это, вызвав в исключении Exception.__init__.

Рассмотрим несколько работающих примеров и один неработающий:

# OK:
class HttpError(Exception):
    pass

# BAD:
class HttpError(Exception):

    def __init__(self, status_code):
        self.status_code = status_code

# OK:
class HttpError(Exception):

    def __init__(self, status_code):
        self.status_code = status_code
        Exception.__init__(self, status_code)  # <-- REQUIRED

Итак, правило таково: для любого исключения с пользовательскими аргументами *args необходимо использовать Exception.__init__(self, *args).

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

class HttpError(Exception):

    def __init__(self, status_code, headers=None, body=None):
        self.status_code = status_code
        self.headers = headers
        self.body = body

        super(HttpError, self).__init__(status_code, headers, body)

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

Игнорирование

Задача может вызвать Ignore, чтобы заставить рабочий процесс проигнорировать её. Это означает, что состояние задачи не будет записано, но сообщение всё равно будет подтверждено (удалено из очереди).

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

Пример сохранения отозванных задач в наборе Redis:

from celery.exceptions import Ignore

@app.task(bind=True)
def some_task(self):
    if redis.ismember('tasks.revoked', self.request.id):
        raise Ignore()

Пример ручного сохранения результатов:

from celery import states
from celery.exceptions import Ignore

@app.task(bind=True)
def get_tweets(self, user):
    timeline = twitter.get_timeline(user)
    if not self.request.called_directly:
        self.update_state(state=states.SUCCESS, meta=timeline)
    raise Ignore()

Отклонение

Задача может вызвать Reject, чтобы отклонить сообщение задачи с помощью метода basic_reject AMQP. Это не даст эффекта, если не включён параметр Task.acks_late.

Отклонение сообщения действует так же, как его подтверждение, однако некоторые брокеры могут предоставлять дополнительную функциональность. Например, RabbitMQ поддерживает концепцию обменников недоставленных сообщений, с помощью которых очередь можно настроить на передачу отклонённых сообщений обменнику недоставленных сообщений для повторной доставки.

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

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

import errno
from celery.exceptions import Reject

@app.task(bind=True, acks_late=True)
def render_scene(self, path):
    file = get_file(path)
    try:
        renderer.render_scene(file)

    # if the file is too big to fit in memory
    # we reject it so that it's redelivered to the dead letter exchange
    # and we can manually inspect the situation.
    except MemoryError as exc:
        raise Reject(exc, requeue=False)
    except OSError as exc:
        if exc.errno == errno.ENOMEM:
            raise Reject(exc, requeue=False)

    # For any other error we retry after 10 seconds.
    except Exception as exc:
        raise self.retry(exc, countdown=10)

Пример повторного помещения сообщения в очередь:

from celery.exceptions import Reject

@app.task(bind=True, acks_late=True)
def requeues(self):
    if not self.request.delivery_info['redelivered']:
        raise Reject('no reason', requeue=True)
    print('received two times')

Подробнее о методе basic_reject см. в документации к используемому брокеру.

Повторная попытка

Исключение Retry вызывается методом Task.retry, чтобы сообщить рабочему процессу о повторном выполнении задачи.

Все задачи наследуются от класса app.Task. Метод run() становится телом задачи.

Например, следующий код

@app.task
def add(x, y):
    return x + y

примерно так будет работать за кулисами:

class _AddTask(app.Task):

    def run(self, x, y):
        return x + y
add = app.tasks[_AddTask.name]

Создание экземпляра

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

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

Если у вас есть задача

from celery import Task

class NaiveAuthenticateServer(Task):

    def __init__(self):
        self.users = {'george': 'password'}

    def run(self, username, password):
        try:
            return self.users[username] == password
        except KeyError:
            return False

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

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

from celery import Task

class DatabaseTask(Task):
    _db = None

    @property
    def db(self):
        if self._db is None:
            self._db = Database.connect()
        return self._db

Для отдельной задачи

Приведённое выше можно добавить к каждой задаче следующим образом:

from celery.app import task

@app.task(base=DatabaseTask, bind=True)
def process_rows(self: task):
    for row in self.db.table.all():
        process_row(row)

Тогда атрибут db задачи process_rows будет всегда оставаться неизменным в каждом процессе.

Для всего приложения

Вы также можете использовать свой пользовательский класс во всём приложении Celery, передав его в качестве аргумента task_cls при создании приложения. Этот аргумент должен быть либо строкой с путём Python к вашему классу Task, либо самим классом:

from celery import Celery

app = Celery('tasks', task_cls='your.module.path:DatabaseTask')

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

По умолчанию используется класс, предоставляемый Celery: 'celery.app.task:Task'.

Обработчики

Обработчики задач — это методы, выполняемые в определённые моменты жизненного цикла задачи. Все обработчики выполняются синхронно в том же процессе и потоке рабочего процесса, что и сама задача.

Хронология выполнения

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

Worker Process Timeline
┌───────────────────────────────────────────────────────────────┐
│  1. before_start()      ← Blocks until complete               │
│  2. run()               ← Your task function                  │
│  3. [Result Backend]    ← State + return value persisted      │
│  4. on_success() OR     ← Outcome-specific handler            │
│     on_retry() OR       │                                     │
│     on_failure()        │                                     │
│  5. after_return()      ← Runs last on terminal states        │
│                       (skipped for RETRY/REJECTED/IGNORED)    │
└───────────────────────────────────────────────────────────────┘

Важно

Основные моменты:

  • Все обработчики выполняются в том же рабочем процессе, что и ваша задача

  • before_start блокирует задачу — run() не начнётся, пока он не завершится

  • Бэкенд результатов обновляется до on_success/on_failure — другие клиенты могут видеть задачу завершённой, пока обработчики ещё выполняются

  • after_return выполняется, когда задача достигает конечного состояния. Он не запускается для RETRY, REJECTED или IGNORED. Если вам нужен обработчик, срабатывающий при каждой попытке, используйте сигнал task_postrun.

Доступные обработчики

before_start(self, task_id, args, kwargs)

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

Примечание

Этот обработчик блокирует задачу: метод run() не начнёт выполняться, пока не завершится before_start.

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

Параметры:
  • task_id — уникальный идентификатор выполняемой задачи.

  • args — исходные аргументы выполняемой задачи.

  • kwargs — исходные именованные аргументы выполняемой задачи.

Возвращаемое значение этого обработчика игнорируется.

on_success(self, retval, task_id, args, kwargs)

Обработчик успешного выполнения.

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

Примечание

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

Параметры:
  • retval — возвращаемое значение задачи.

  • task_id — уникальный идентификатор выполненной задачи.

  • args — исходные аргументы выполненной задачи.

  • kwargs — исходные именованные аргументы выполненной задачи.

Возвращаемое значение этого обработчика игнорируется.

on_retry(self, exc, task_id, args, kwargs, einfo)

Обработчик повторной попытки.

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

Примечание

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

Параметры:
  • exc — исключение, переданное в retry().

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

  • args — исходные аргументы задачи, выполняемой повторно.

  • kwargs — исходные именованные аргументы задачи, выполняемой повторно.

  • einfo — экземпляр ExceptionInfo.

Возвращаемое значение этого обработчика игнорируется.

on_failure(self, exc, task_id, args, kwargs, einfo)

Обработчик сбоя.

Вызывается рабочим процессом при сбое задачи.

Примечание

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

Параметры:
  • exc — исключение, вызванное задачей.

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

  • args — исходные аргументы задачи, завершившейся с ошибкой.

  • kwargs — исходные именованные аргументы задачи, завершившейся с ошибкой.

  • einfo — экземпляр ExceptionInfo.

Возвращаемое значение этого обработчика игнорируется.

after_return(self, status, retval, task_id, args, kwargs, einfo)

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

Примечание

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

На практике это означает, что он запускается после on_success или on_failure. Он не выполняется для состояний RETRY, REJECTED или IGNORED. Если обработчик нужен для каждой попытки, рассмотрите возможность использования сигнала task_postrun.

Параметры:
  • status — текущее состояние задачи.

  • retval — возвращаемое значение задачи или исключение.

  • task_id — уникальный идентификатор задачи.

  • args — исходные аргументы завершившейся задачи.

  • kwargs — исходные именованные аргументы завершившейся задачи.

  • einfo — экземпляр ExceptionInfo.

Возвращаемое значение этого обработчика игнорируется.

Пример использования

import time
from celery import Task

class MyTask(Task):

    def before_start(self, task_id, args, kwargs):
        print(f"Task {task_id} starting with args {args}")
        # This blocks - run() won't start until this returns

    def on_success(self, retval, task_id, args, kwargs):
        print(f"Task {task_id} succeeded with result: {retval}")
        # Result is already visible to clients at this point

    def on_failure(self, exc, task_id, args, kwargs, einfo):
        print(f"Task {task_id} failed: {exc}")
        # Task state is already FAILURE in backend

    def after_return(self, status, retval, task_id, args, kwargs, einfo):
        print(f"Task {task_id} finished with status: {status}")
        # Always runs last

@app.task(base=MyTask)
def my_task(x, y):
    return x + y

Запросы и пользовательские запросы

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

Пользовательские классы задач могут переопределить используемый класс запроса, изменив атрибут celery.app.task.Task.Request. Можно присвоить сам пользовательский класс запроса или его полное имя.

У запроса несколько обязанностей. Пользовательские классы запросов должны реализовать их все: именно они фактически запускают и отслеживают задачу. Мы настоятельно рекомендуем наследоваться от celery.worker.request.Request.

При использовании рабочего процесса с предварительным порождением процессов методы on_timeout() и on_failure() выполняются в главном процессе рабочего процесса. Приложение может использовать эту возможность для обнаружения сбоев, которые не обнаруживаются с помощью celery.app.task.Task.on_failure().

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

import logging
from celery import Task
from celery.worker.request import Request

logger = logging.getLogger('my.package')

class MyRequest(Request):
    'A minimal custom request to log failures and hard time limits.'

    def on_timeout(self, soft, timeout):
        super(MyRequest, self).on_timeout(soft, timeout)
        if not soft:
           logger.warning(
               'A hard timeout was enforced for task %s',
               self.task.name
           )

    def on_failure(self, exc_info, send_failed_event=True, return_ok=False):
        super().on_failure(
            exc_info,
            send_failed_event=send_failed_event,
            return_ok=return_ok
        )
        logger.warning(
            'Failure detected for task %s',
            self.task.name
        )

class MyTask(Task):
    Request = MyRequest  # you can use a FQN 'my.package:MyRequest'

@app.task(base=MyTask)
def some_longrunning_task():
    # use your imagination

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

Все определённые задачи перечислены в реестре. Реестр содержит список имён задач и соответствующих им классов. Вы можете самостоятельно изучить этот реестр:

>>> from proj.celery import app
>>> app.tasks
{'celery.chord_unlock':
    <@task: celery.chord_unlock>,
 'celery.backend_cleanup':
    <@task: celery.backend_cleanup>,
 'celery.chord':
    <@task: celery.chord>}

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

Загрузчик по умолчанию импортирует все модули, перечисленные в параметре imports.

Декоратор app.task() отвечает за регистрацию вашей задачи в реестре задач приложения.

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

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

Игнорируйте ненужные результаты

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

@app.task(ignore_result=True)
def mytask():
    something()

Результаты можно также отключить глобально с помощью параметра task_ignore_result.

Результаты можно включать и отключать для отдельного выполнения, передавая булев параметр ignore_result при вызове apply_async.

@app.task
def mytask(x, y):
    return x + y

# No result will be stored
result = mytask.apply_async((1, 2), ignore_result=True)
print(result.get()) # -> None

# Result will be stored
result = mytask.apply_async((1, 2), ignore_result=False)
print(result.get()) # -> 3

По умолчанию задачи не игнорируют результаты (ignore_result=False), если настроен бэкенд результатов.

Порядок приоритета параметров таков:

  1. Глобальный параметр task_ignore_result

  2. Параметр ignore_result

  3. Параметр выполнения задачи ignore_result

Дополнительные советы по оптимизации

Дополнительные советы по оптимизации приведены в руководстве по оптимизации.

Избегайте запуска синхронных подзадач

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

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

Плохо:

@app.task
def update_page_info(url):
    page = fetch_page.delay(url).get()
    info = parse_page.delay(page).get()
    store_page_info.delay(url, info)

@app.task
def fetch_page(url):
    return myhttplib.get(url)

@app.task
def parse_page(page):
    return myparser.parse_document(page)

@app.task
def store_page_info(url, info):
    return PageInfo.objects.create(url, info)

Хорошо:

def update_page_info(url):
    # fetch_page -> parse_page -> store_page
    chain = fetch_page.s(url) | parse_page.s() | store_page_info.s(url)
    chain()

@app.task()
def fetch_page(url):
    return myhttplib.get(url)

@app.task()
def parse_page(page):
    return myparser.parse_document(page)

@app.task(ignore_result=True)
def store_page_info(info, url):
    PageInfo.objects.create(url=url, info=info)

В этом примере я создал цепочку задач, связав несколько объектов signature(). О цепочках и других мощных конструкциях можно прочитать в разделе Canvas: проектирование рабочих процессов.

По умолчанию Celery не позволяет запускать подзадачи синхронно внутри задачи, однако в редких или крайних случаях это может понадобиться. ВНИМАНИЕ: включать синхронный запуск подзадач не рекомендуется!

@app.task
def update_page_info(url):
    page = fetch_page.delay(url).get(disable_sync_subtasks=False)
    info = parse_page.delay(page).get(disable_sync_subtasks=False)
    store_page_info.delay(url, info)

@app.task
def fetch_page(url):
    return myhttplib.get(url)

@app.task
def parse_page(page):
    return myparser.parse_document(page)

@app.task
def store_page_info(url, info):
    return PageInfo.objects.create(url, info)

Детализация

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

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

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

См. также

В книге «Искусство параллелизма» есть раздел, посвящённый детализации задач [AOC1].

[AOC1]

Бреширс, Клей. Раздел 2.2.1, «Искусство параллелизма». O’Reilly Media, Inc. 15 мая 2009 г. ISBN-13 978-0-596-52153-0.

Локальность данных

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

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

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

См. также

Статья Джима Грея «Экономика распределённых вычислений» — отличное введение в тему локальности данных.

Состояние

Поскольку Celery — распределённая система, вы не можете знать, какой процесс или компьютер будет выполнять задачу. Вы даже не можете знать, будет ли задача выполнена своевременно.

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

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

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

class Article(models.Model):
    title = models.CharField()
    body = models.TextField()

@app.task
def expand_abbreviations(article):
    article.body.replace('MyCorp', 'My Corporation')
    article.save()

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

>>> article = Article.objects.get(id=102)
>>> expand_abbreviations.delay(article)

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

Устранить состояние гонки просто: достаточно передавать идентификатор статьи и повторно получать статью в теле задачи:

@app.task
def expand_abbreviations(article_id):
    article = Article.objects.get(id=article_id)
    article.body.replace('MyCorp', 'My Corporation')
    article.save()
>>> expand_abbreviations.delay(article_id)

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

Транзакции базы данных

Рассмотрим ещё один пример:

from django.db import transaction
from django.http import HttpResponseRedirect

@transaction.atomic
def create_article(request):
    article = Article.objects.create()
    expand_abbreviations.delay(article.pk)
    return HttpResponseRedirect('/articles/')

Это представление Django создаёт объект статьи в базе данных, а затем передаёт первичный ключ задаче. Используется декоратор transaction.atomic, который фиксирует транзакцию при возврате представления или откатывает её, если представление вызывает исключение.

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

Решение — использовать вместо этого delay_on_commit():

from django.db import transaction
from django.http import HttpResponseRedirect

@transaction.atomic
def create_article(request):
    article = Article.objects.create()
    expand_abbreviations.delay_on_commit(article.pk)
    return HttpResponseRedirect('/articles/')

Этот метод появился в Celery 5.4. Это сокращённая запись, использующая функцию обратного вызова on_commit в Django для запуска задачи Celery после успешной фиксации всех транзакций.

Для Celery <5.4

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

import functools
from django.db import transaction
from django.http import HttpResponseRedirect

@transaction.atomic
def create_article(request):
    article = Article.objects.create()
    transaction.on_commit(
        functools.partial(expand_abbreviations.delay, article.pk)
    )
    return HttpResponseRedirect('/articles/')

Примечание

on_commit доступен в Django 1.9 и более поздних версиях. Для более ранних версий библиотека django-transaction-hooks добавляет эту возможность.

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

У меня есть приложение блога на Django, позволяющее оставлять комментарии к записям. Я расскажу о моделях, представлениях и задачах этого приложения.

blog/models.py

Модель комментария выглядит так:

from django.db import models
from django.utils.translation import ugettext_lazy as _


class Comment(models.Model):
    name = models.CharField(_('name'), max_length=64)
    email_address = models.EmailField(_('email address'))
    homepage = models.URLField(_('home page'),
                               blank=True, verify_exists=False)
    comment = models.TextField(_('comment'))
    pub_date = models.DateTimeField(_('Published date'),
                                    editable=False, auto_add_now=True)
    is_spam = models.BooleanField(_('spam?'),
                                  default=False, editable=False)

    class Meta:
        verbose_name = _('comment')
        verbose_name_plural = _('comments')

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

blog/views.py

from django import forms
from django.http import HttpResponseRedirect
from django.template.context import RequestContext
from django.shortcuts import get_object_or_404, render_to_response

from blog import tasks
from blog.models import Comment


class CommentForm(forms.ModelForm):

    class Meta:
        model = Comment


def add_comment(request, slug, template_name='comments/create.html'):
    post = get_object_or_404(Entry, slug=slug)
    remote_addr = request.META.get('REMOTE_ADDR')

    if request.method == 'post':
        form = CommentForm(request.POST, request.FILES)
        if form.is_valid():
            comment = form.save()
            # Check spam asynchronously.
            tasks.spam_filter.delay(comment_id=comment.id,
                                    remote_addr=remote_addr)
            return HttpResponseRedirect(post.get_absolute_url())
    else:
        form = CommentForm()

    context = RequestContext(request, {'form': form})
    return render_to_response(template_name, context_instance=context)

Для проверки комментариев на спам я использую Akismet — сервис, применяемый для фильтрации спама в комментариях на бесплатной платформе для блогов Wordpress. Akismet бесплатен для личного использования, но за коммерческое использование нужно платить. Чтобы получить ключ API, необходимо зарегистрироваться в сервисе.

Для вызовов API Akismet я использую библиотеку akismet.py, написанную Майклом Фурдом.

blog/tasks.py

from celery import Celery

from akismet import Akismet

from django.core.exceptions import ImproperlyConfigured
from django.contrib.sites.models import Site

from blog.models import Comment


app = Celery(broker='amqp://')


@app.task
def spam_filter(comment_id, remote_addr=None):
    logger = spam_filter.get_logger()
    logger.info('Running spam filter for comment %s', comment_id)

    comment = Comment.objects.get(pk=comment_id)
    current_domain = Site.objects.get_current().domain
    akismet = Akismet(settings.AKISMET_KEY, 'http://{0}'.format(domain))
    if not akismet.verify_key():
        raise ImproperlyConfigured('Invalid AKISMET_KEY')


    is_spam = akismet.comment_check(user_ip=remote_addr,
                        comment_content=comment.comment,
                        comment_author=comment.name,
                        comment_author_email=comment.email_address)
    if is_spam:
        comment.is_spam = True
        comment.save()

    return is_spam

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

Spec-Zone.ru

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