Spec-Zone.ru › Celery

celery.app.task

Реализация задач: контекст запроса и базовый класс задачи.

classcelery.app.task.Context(*args, **kwargs)

Переменные запроса задачи (Task.request).

classcelery.app.task.Task

Базовый класс задачи.

Примечание

При вызове задачи применяют метод run(). Этот метод должен быть определён для всех задач (если только метод __call__() не переопределён).

AsyncResult(task_id, **kwargs)

Получить экземпляр AsyncResult для указанной задачи.

Параметры:

task_id (str) – Идентификатор задачи, для которой нужно получить результат.

exceptionMaxRetriesExceededError(*args, **kwargs)

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

exceptionOperationalError

Устранимая ошибка подключения транспортного уровня сообщений.

Request='celery.worker.request:Request'

Используемый класс Request или его полное имя.

Strategy='celery.worker.strategy:default'

Используемая стратегия выполнения или её полное имя.

abstract=True

Устаревший атрибут abstract оставлен здесь для совместимости.

acks_late=False

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

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

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

acks_on_failure_or_timeout=True

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

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

Значение по умолчанию для приложения можно переопределить с помощью параметра task_acks_on_failure_or_timeout.

add_to_chord(sig, lazy=False)

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

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

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

Параметры:
  • sig (Signature) – Сигнатура, которой нужно расширить аккорд.

  • lazy (bool) – Если параметр включён, новая задача фактически не будет вызвана, и sig.delay() нужно будет вызвать вручную.

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

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

Параметры:
  • status (str) – Текущее состояние задачи.

  • retval (Any) – Возвращаемое значение или исключение задачи.

  • task_id (str) – Уникальный идентификатор задачи.

  • args (Tuple) – Исходные аргументы задачи.

  • kwargs (Dict) – Исходные именованные аргументы задачи.

  • einfo (ExceptionInfo) – Информация об исключении.

Возвращает:

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

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

None

apply(args=None, kwargs=None, link=None, link_error=None, task_id=None, retries=None, throw=None, logfile=None, loglevel=None, headers=None, **options)

Выполнить эту задачу локально, ожидая её завершения.

Параметры:
  • args (Tuple) – Позиционные аргументы, передаваемые задаче.

  • kwargs (Dict) – Именованные аргументы, передаваемые задаче.

  • throw (bool) – Повторно вызвать исключения задачи. По умолчанию используется значение параметра task_eager_propagates.

Возвращает:

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

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

celery.result.EagerResult

apply_async(args=None, kwargs=None, task_id=None, producer=None, link=None, link_error=None, shadow=None, **options)

Асинхронно запустить задачу, отправив сообщение.

Параметры:
  • args (Tuple) – Позиционные аргументы, передаваемые задаче.

  • kwargs (Dict) – Именованные аргументы, передаваемые задаче.

  • countdown (float) – Через сколько секунд должна выполниться задача. По умолчанию задача выполняется немедленно.

  • eta (datetime) – Абсолютные дата и время выполнения задачи. Нельзя указывать, если также задан countdown.

  • expires (float, datetime) – Дата и время или количество секунд до истечения срока действия задачи. По истечении этого срока задача выполняться не будет.

  • shadow (str) – Переопределяет имя задачи, используемое в журналах и системах мониторинга. По умолчанию берётся из shadow_name().

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

  • retry (bool) – Если параметр включён, отправка сообщения задачи будет повторяться при потере или сбое подключения. По умолчанию используется значение параметра task_publish_retry. Обратите внимание: чтобы это работало, необходимо вручную управлять производителем и подключением.

  • retry_policy (Mapping) – Переопределяет используемую политику повторных попыток. См. параметр task_publish_retry_policy.

  • time_limit (int) – Если задан, переопределяет ограничение времени по умолчанию.

  • soft_time_limit (int) – Если задан, переопределяет мягкое ограничение времени по умолчанию.

  • queue (str, kombu.Queue) – Очередь, в которую направляется задача. Она должна присутствовать в качестве ключа в task_queues либо должен быть включён параметр task_create_missing_queues. Подробнее см. в разделе Маршрутизация задач.

  • exchange (str, kombu.Exchange) – Именованный пользовательский обменник, в который отправляется задача. Обычно не используется вместе с аргументом queue.

  • routing_key (str) – Пользовательский ключ маршрутизации для направления задачи на сервер рабочего процесса. Если задан вместе с аргументом queue, используется только для указания пользовательских ключей маршрутизации для тематических обменников.

  • priority (int) – Приоритет задачи — число от 0 до 9. По умолчанию используется значение атрибута priority.

  • serializer (str) – Метод сериализации. Допустимы pickle, json, yaml, msgpack или любой зарегистрированный пользовательский метод сериализации из kombu.serialization.registry. По умолчанию используется значение атрибута serializer.

  • compression (str) – Необязательный метод сжатия. Допустимы zlib, bzip2 или любые пользовательские методы сжатия, зарегистрированные с помощью kombu.compression.register(). По умолчанию используется значение параметра task_compression.

  • link (Signature) – Одна сигнатура задачи или список сигнатур, применяемых при успешном возврате задачи.

  • link_error (Signature) – Одна сигнатура задачи или список сигнатур, применяемых при возникновении ошибки во время выполнения задачи.

  • producer (kombu.Producer) – Пользовательский производитель, используемый при публикации задачи.

  • add_to_parent (bool) – Если установлено значение True (по умолчанию) и задача запускается во время выполнения другой задачи, её результат будет добавлен в атрибут request.children родительской задачи. Добавление в журнал выполнения также можно отключить по умолчанию с помощью атрибута trail.

  • ignore_result (bool) – Если задано значение False (по умолчанию), результат задачи будет сохранён в бэкенде. Если задано значение True, результат сохранён не будет. Это также можно настроить с помощью параметра ignore_result в декораторе app.task.

  • publisher (kombu.Producer) – Устаревший псевдоним для producer.

  • headers (Dict) – Заголовки, включаемые в сообщение. Их можно использовать как дополнительную маркировку с помощью функции Маркировка.

  • task_id (str) – Необязательный аргумент для переопределения идентификатора задачи по умолчанию. По умолчанию Celery генерирует уникальный идентификатор (UUID4) для каждой отправленной задачи. Вместо этого можно задать собственный строковый идентификатор. Если он указан, это значение будет использоваться в качестве идентификатора задачи вместо автоматической генерации. При переопределении идентификаторов задач следите за тем, чтобы избежать совпадений.

Возвращает:

Обещание будущего вычисления.

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

celery.result.AsyncResult

Вызывает исключения:
  • TypeError – Если передано недостаточно или слишком много аргументов. Обратите внимание, что проверку сигнатуры можно отключить, указав @task(typing=False).

  • ValueError – Если заданы оба параметра soft_time_limit и time_limit, но значение soft_time_limit больше значения time_limit.

  • kombu.exceptions.OperationalError – Если не удаётся установить подключение к транспортному уровню или если подключение потеряно.

Примечание

Также поддерживаются все именованные аргументы, поддерживаемые методом kombu.Producer.publish().

propertybackend

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

before_start(task_id, args, kwargs)

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

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

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

  • args (Tuple) – Исходные аргументы задачи, которую нужно выполнить.

  • kwargs (Dict) – Исходные именованные аргументы задачи, которую нужно выполнить.

Возвращает:

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

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

None

chunks(it, n)

Создать задачу chunks для этой задачи.

default_retry_delay=180

Задержка в секундах перед повторным запуском задачи по умолчанию. По умолчанию — 3 минуты.

delay(*args, **kwargs)

Версия метода apply_async() с аргументами со звёздочкой.

Не поддерживает дополнительные параметры, доступные в методе apply_async().

Параметры:
  • *args (Any) – Позиционные аргументы, передаваемые задаче.

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

Возвращает:

Обещание будущего результата.

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

celery.result.AsyncResult

expires=None

Срок действия задачи по умолчанию.

ignore_result=False

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

map(it)

Создать задачу xmap из it.

max_retries=3

Максимальное число повторных попыток до прекращения. Если задано значение None, повторные попытки никогда не прекратятся.

name=None

Имя задачи.

classmethodon_bound(app)

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

Примечание

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

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

Обработчик ошибок.

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

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

  • task_id (str) – Уникальный идентификатор задачи, завершившейся с ошибкой.

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

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

  • einfo (ExceptionInfo) – Информация об исключении.

Возвращает:

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

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

None

on_replace(sig)

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

При переопределении необходимо возвращать super().on_replace(sig), чтобы замена задачи обрабатывалась правильно.

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

Параметры:

sig (Signature) – сигнатура, которой нужно заменить задачу.

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

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

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

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

  • task_id (str) – Уникальный идентификатор задачи, для которой выполняется повторная попытка.

  • args (Tuple) – Исходные аргументы задачи, для которой выполняется повторная попытка.

  • kwargs (Dict) – Исходные именованные аргументы задачи, для которой выполняется повторная попытка.

  • einfo (ExceptionInfo) – Информация об исключении.

Возвращает:

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

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

None

on_success(retval, task_id, args, kwargs)

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

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

Параметры:
  • retval (Any) – Возвращаемое значение задачи.

  • task_id (str) – Уникальный идентификатор выполненной задачи.

  • args (Tuple) – Исходные аргументы выполненной задачи.

  • kwargs (Dict) – Исходные именованные аргументы выполненной задачи.

Возвращает:

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

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

None

priority=None

Приоритет задачи по умолчанию.

rate_limit=None

None (без ограничения частоты), ‘100/s’ (сто задач в секунду), ‘100/m’ (сто задач в минуту),`’100/h’` (сто задач в час)

Тип:

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

reject_on_worker_lost=None

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

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

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

replace(sig)

Заменить эту задачу новой задачей, унаследовавшей ее идентификатор.

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

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

Параметры:
  • sig (Signature) – сигнатура, которой нужно заменить задачу.

  • visitor (StampingVisitor) – объект API посетителя.

Вызывает исключения:
  • ~@Ignore – Всегда вызывается при обращении в асинхронном контексте.

  • Лучше всегда использовать return self.replace(...) , чтобы сообщить –

  • читателю, что после замены задача не продолжит выполнение. –

propertyrequest

Получить текущий объект запроса.

request_stack=<celery.utils.threads._LocalStack object>

Стек запросов задач; текущий запрос находится на вершине стека.

resultrepr_maxsize=1024

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

retry(args=None, kwargs=None, exc=None, throw=True, eta=None, countdown=None, max_retries=None, **options)

Повторить задачу, добавив ее в конец очереди.

Пример

>>> from imaginary_twitter_lib import Twitter
>>> from proj.celery import app
>>> @app.task(bind=True)
... def tweet(self, auth, message):
...     twitter = Twitter(oauth=auth)
...     try:
...         twitter.post_status_update(message)
...     except twitter.FailWhale as exc:
...         # Retry in 5 minutes.
...         raise self.retry(countdown=60 * 5, exc=exc)

Примечание

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

Параметры:
  • args (Tuple) – Позиционные аргументы для повторного запуска.

  • kwargs (Dict) – Именованные аргументы для повторного запуска.

  • exc (Exception) –

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

    Если этот аргумент задан и retry вызывается при возникшем исключении (установлен sys.exc_info()), будет предпринята попытка повторно вызвать текущее исключение.

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

  • countdown (float) – Задержка перед повторной попыткой в секундах.

  • eta (datetime) – Точное время и дата повторного запуска.

  • max_retries (int) – Если задано, переопределяет ограничение числа повторных попыток по умолчанию для этого выполнения. Изменения этого параметра не распространяются на последующие попытки повторного запуска задачи. Значение None означает «использовать значение по умолчанию», поэтому для бесконечных повторных попыток сначала нужно установить атрибут задачи max_retries в None.

  • time_limit (int) – Если задано, переопределяет ограничение времени по умолчанию.

  • soft_time_limit (int) – Если задано, переопределяет мягкое ограничение времени по умолчанию.

  • throw (bool) – Если установлено значение False, исключение Retry не вызывается; оно сообщает рабочему процессу, что задача повторно запускается. Это означает, что задача будет помечена как завершившаяся с ошибкой, если вызовет исключение, или как успешно завершившаяся, если вернет результат после вызова retry.

  • **options (Any) – Дополнительные параметры, передаваемые в apply_async().

Вызывает исключения:

celery.exceptions.Retry – Сообщает рабочему процессу, что задача отправлена повторно для нового запуска. Это происходит всегда, если только для именованного аргумента throw явно не задано значение False; такое поведение считается штатным.

run(*args, **kwargs)

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

s(*args, **kwargs)

Создать сигнатуру.

Сокращение для .s(*a, **k) -> .signature(a, k).

send_event(type_, retry=True, retry_policy=None, **fields)

Отправить сообщение с событием мониторинга.

Можно использовать для добавления пользовательских типов событий в https://pypi.org/project/Flower/ и другие системы мониторинга.

Параметры:

type (str) – Тип события, например "task-failed".

Именованные аргументы:
  • retry (bool) – Повторить отправку сообщения при потере соединения. Значение по умолчанию берется из настройки task_publish_retry.

  • retry_policy (Mapping) – Параметры повторных попыток. Значение по умолчанию берется из настройки task_publish_retry_policy.

  • **fields (Any) – Словарь со сведениями о событии. Данные должны быть сериализуемыми в JSON.

send_events=True

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

serializer='json'

Имя сериализатора, зарегистрированного в kombu.serialization.registry. Значение по умолчанию — ‘json’.

shadow_name(args, kwargs, options)

Переопределение имени задачи в журналах и средствах мониторинга рабочего процесса.

Пример

from celery.utils.imports import qualname

def shadow_name(task, args, kwargs, options):
    return qualname(args[0])

@app.task(shadow_name=shadow_name, serializer='pickle')
def apply_function_async(fun, *args, **kwargs):
    return fun(*args, **kwargs)
Параметры:
  • args (Tuple) – Позиционные аргументы задачи.

  • kwargs (Dict) – Именованные аргументы задачи.

  • options (Dict) – Параметры выполнения задачи.

si(*args, **kwargs)

Создать неизменяемую сигнатуру.

Сокращение для .si(*a, **k) -> .signature(a, k, immutable=True).

signature(args=None, *starargs, **starkwargs)

Создать сигнатуру.

Возвращает:
объект для

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

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

signature

soft_time_limit=None

Мягкое ограничение времени. По умолчанию используется настройка task_soft_time_limit.

starmap(it)

Создать задачу xstarmap из it.

store_errors_even_if_ignored=False

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

subtask(args=None, *starargs, **starkwargs)

Создать сигнатуру.

Возвращает:
объект для

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

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

signature

throws=()

Кортеж ожидаемых исключений.

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

time_limit=None

Жесткое ограничение времени. По умолчанию используется настройка task_time_limit.

track_started=False

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

Состояние «started» может быть полезно для длительных задач, когда необходимо сообщать, какая задача выполняется в данный момент.

Значение по умолчанию для приложения можно переопределить с помощью настройки task_track_started.

trail=True

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

typing=True

Включить проверку аргументов. Установите значение false, если не нужно проверять сигнатуру при вызове задачи. По умолчанию используется значение Celery.strict_typing.

update_state(task_id=None, state=None, meta=None, **kwargs)

Обновить состояние задачи.

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

  • state (str) – Новое состояние.

  • meta (Dict) – Метаданные состояния.

celery.app.task.TaskType

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

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

Spec-Zone.ru

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