Вызов задач
В этом документе описывается единый «API вызова» Celery, используемый экземплярами задач и canvas.
API определяет стандартный набор параметров выполнения, а также три метода:
-
apply_async(args[, kwargs[, …]])Отправляет сообщение задачи.
-
delay(*args, **kwargs)Сокращённая запись для отправки сообщения задачи, не поддерживающая параметры выполнения.
-
вызов (
__call__)Применение объекта, поддерживающего API вызова (например,
add(2, 2)), означает, что задача будет выполнена не рабочим процессом, а в текущем процессе (сообщение отправлено не будет).
Пример
Метод delay() удобен тем, что выглядит как вызов обычной функции:
task.delay(arg1, arg2, kwarg1='x', kwarg2='y')
При использовании apply_async() вместо него нужно написать:
task.apply_async(args=[arg1, arg2], kwargs={'kwarg1': 'x', 'kwarg2': 'y'})
Таким образом, delay — удобный способ вызова, но для задания дополнительных параметров выполнения нужно использовать apply_async.
Далее в этом документе подробно рассматриваются параметры выполнения задач. Во всех примерах используется задача с именем add, возвращающая сумму двух аргументов:
@app.task
def add(x, y):
return x + y
Celery поддерживает связывание задач, при котором одна задача выполняется вслед за другой. Задаче обратного вызова передаётся результат родительской задачи в качестве частичного аргумента:
add.apply_async((2, 2), link=add.s(16))
Здесь результат первой задачи (4) будет передан новой задаче, которая прибавит 16 к предыдущему результату, образуя выражение
Также можно настроить выполнение обратного вызова при возникновении исключения в задаче (обработчик ошибки). Рабочий процесс не будет вызывать обработчик ошибки как задачу, а вызовет его функцию напрямую, чтобы передать ей исходный запрос, исключение и объекты трассировки стека.
Пример обработчика ошибки:
@app.task
def error_handler(request, exc, traceback):
print('Task {0} raised exception: {1!r}\n{2!r}'.format(
request.id, exc, traceback))
его можно добавить к задаче с помощью параметра выполнения link_error:
add.apply_async((2, 2), link_error=error_handler.s())
Кроме того, параметры link и link_error можно задать в виде списка:
add.apply_async((2, 2), link=[add.s(16), other_task.s()])
После этого обратные вызовы и обработчики ошибок будут вызваны по порядку, а всем обратным вызовам будет передано возвращаемое значение родительской задачи в качестве частичного аргумента.
В случае группы задач (chord) ошибки можно обрабатывать несколькими способами. Подробнее см. в разделе обработка ошибок chord.
Celery позволяет перехватывать все изменения состояния с помощью обратного вызова on_message.
Например, для отправки сведений о ходе выполнения длительных задач можно сделать следующее:
@app.task(bind=True)
def hello(self, a, b):
time.sleep(1)
self.update_state(state="PROGRESS", meta={'progress': 50})
time.sleep(1)
self.update_state(state="PROGRESS", meta={'progress': 90})
time.sleep(1)
return 'hello world: %i' % (a+b)
def on_raw_message(body):
print(body)
a, b = 1, 1
r = hello.apply_async(args=(a, b))
print(r.get(on_message=on_raw_message, propagate=False))
В результате будет выведено примерно следующее:
{'task_id': '5660d3a3-92b8-40df-8ccc-33a5d1d680d7',
'result': {'progress': 50},
'children': [],
'status': 'PROGRESS',
'traceback': None}
{'task_id': '5660d3a3-92b8-40df-8ccc-33a5d1d680d7',
'result': {'progress': 90},
'children': [],
'status': 'PROGRESS',
'traceback': None}
{'task_id': '5660d3a3-92b8-40df-8ccc-33a5d1d680d7',
'result': 'hello world: 10',
'children': [],
'status': 'SUCCESS',
'traceback': None}
hello world: 10
ETA (расчётное время прибытия) позволяет задать конкретные дату и время, раньше которых задача не будет выполнена. Параметр countdown позволяет задать ETA в секундах от текущего момента.
>>> result = add.apply_async((2, 2), countdown=3) >>> result.get() # this takes at least 3 seconds to return 4
Гарантируется, что задача будет выполнена в некоторый момент после указанной даты и времени, но не обязательно точно в этот момент. Возможные причины нарушения сроков — большое количество задач в очереди или высокая задержка сети. Чтобы задачи выполнялись своевременно, следует следить за перегрузкой очереди. Используйте Munin или аналогичные инструменты для получения оповещений, чтобы можно было принять необходимые меры и снизить нагрузку. См. раздел Munin.
Значение countdown — целое число, а eta должно быть объектом datetime, задающим точные дату и время (с точностью до миллисекунд и с указанием часового пояса):
>>> from datetime import datetime, timedelta, timezone >>> tomorrow = datetime.now(timezone.utc) + timedelta(days=1) >>> add.apply_async((2, 2), eta=tomorrow)
Предупреждение
Задачи с параметрами eta или countdown немедленно забираются рабочим процессом и остаются в его памяти до наступления запланированного времени. Если использовать эти параметры для планирования большого количества задач на отдалённое будущее, задачи могут накопиться в рабочем процессе и существенно увеличить потребление оперативной памяти.
Кроме того, задачи не подтверждаются до тех пор, пока рабочий процесс не начнёт их выполнять. При использовании Redis в качестве брокера задача будет доставлена повторно, если значение countdown превышает visibility_timeout (см. раздел Предостережения).
Поэтому не рекомендуется использовать eta и countdown для планирования задач на отдалённое будущее. В идеале следует задавать значения не более нескольких минут. Для более длительных интервалов рассмотрите периодические задачи с хранением в базе данных, например, с помощью https://pypi.org/project/django-celery-beat/ при использовании Django (см. раздел Использование пользовательских классов планировщика).
Предупреждение
При использовании RabbitMQ в качестве брокера сообщений и указании countdown более 15 минут может возникнуть проблема: рабочий процесс завершается, и выдаётся ошибка PreconditionFailed:
amqp.exceptions.PreconditionFailed: (0, 0): (406) PRECONDITION_FAILED - consumer ack timed out on channel
В RabbitMQ начиная с версии 3.8.15 значение consumer_timeout по умолчанию составляет 15 минут. В версии 3.8.17 оно было увеличено до 30 минут. Если потребитель не подтверждает доставку дольше заданного времени ожидания, его канал закрывается с исключением канала PRECONDITION_FAILED. Подробнее см. в разделе Тайм-аут подтверждения доставки.
Чтобы устранить эту проблему, в файле конфигурации RabbitMQ rabbitmq.conf следует задать параметр consumer_timeout, значение которого должно быть не меньше значения countdown. Например, чтобы избежать проблем в будущем, можно указать очень большое значение consumer_timeout = 31622400000, равное одному году в миллисекундах.
Аргумент expires задаёт необязательный срок действия: количество секунд после публикации задачи либо конкретные дату и время с помощью datetime:
>>> # Task expires after one minute from now. >>> add.apply_async((10, 10), expires=60) >>> # Also supports datetime >>> from datetime import datetime, timedelta, timezone >>> add.apply_async((10, 10), kwargs, ... expires=datetime.now(timezone.utc) + timedelta(days=1))
Получив просроченную задачу, рабочий процесс пометит её как REVOKED (TaskRevokedError).
Celery автоматически повторяет отправку сообщений при сбое подключения. Поведение повторных попыток можно настроить — например, задать интервал между попытками или их максимальное количество — либо полностью отключить.
Чтобы отключить повторы, задайте для параметра выполнения retry значение False:
add.apply_async((2, 2), retry=False)
Политика повторных попыток
Политика повторных попыток — это набор параметров, управляющих поведением повторных попыток. Она может содержать следующие ключи:
-
max_retries
Максимальное количество повторных попыток. После его достижения будет вызвано исключение, из-за которого не удалось выполнить повторную попытку.
Значение
Noneозначает, что попытки будут выполняться бесконечно.По умолчанию выполняется 3 повторные попытки.
-
interval_start
Задаёт количество секунд (число с плавающей точкой или целое число), которое нужно подождать перед повторной попыткой. Значение по умолчанию — 0 (первая повторная попытка выполняется немедленно).
-
interval_step
При каждой следующей повторной попытке это число прибавляется к задержке (число с плавающей точкой или целое число). Значение по умолчанию — 0.2.
-
interval_max
Максимальное количество секунд ожидания между повторными попытками (число с плавающей точкой или целое число). Значение по умолчанию — 0.2.
-
retry_errors
retry_errors — это кортеж классов исключений, для которых следует выполнять повторные попытки. Если параметр не задан, он игнорируется. Значение по умолчанию — None (игнорируется).
Например, если нужно повторять попытки только для задач, выполнение которых завершилось по тайм-ауту, можно использовать
TimeoutError:from kombu.exceptions import TimeoutError add.apply_async((2, 2), retry=True, retry_policy={ 'max_retries': 3, 'retry_errors': (TimeoutError, ), })Добавлено в версии 5.3.
Например, стандартная политика соответствует следующим параметрам:
add.apply_async((2, 2), retry=True, retry_policy={
'max_retries': 3,
'interval_start': 0,
'interval_step': 0.2,
'interval_max': 0.2,
'retry_errors': None,
})
максимальное время повторных попыток составит 0.4 секунды. По умолчанию оно сравнительно мало, поскольку сбой подключения может вызвать лавинообразный рост повторных попыток при недоступности подключения к брокеру. Например, множество процессов веб-сервера будут ожидать повтора, блокируя обработку других входящих запросов.
При отправке задачи, если соединение с транспортом сообщений потеряно или его не удаётся установить, будет вызвана ошибка OperationalError:
>>> from proj.tasks import add
>>> add.delay(2, 2)
Traceback (most recent call last):
File "<stdin>", line 1, in <module>
File "celery/app/task.py", line 388, in delay
return self.apply_async(args, kwargs)
File "celery/app/task.py", line 503, in apply_async
**options
File "celery/app/base.py", line 662, in send_task
amqp.send_task_message(P, name, message, **options)
File "celery/backends/rpc.py", line 275, in on_task_call
maybe_declare(self.binding(producer.channel), retry=True)
File "/opt/celery/kombu/kombu/messaging.py", line 204, in _get_channel
channel = self._channel = channel()
File "/opt/celery/py-amqp/amqp/connection.py", line 272, in connect
self.transport.connect()
File "/opt/celery/py-amqp/amqp/transport.py", line 100, in connect
self._connect(self.host, self.port, self.connect_timeout)
File "/opt/celery/py-amqp/amqp/transport.py", line 141, in _connect
self.sock.connect(sa)
kombu.exceptions.OperationalError: [Errno 61] Connection refused
Если повторные попытки включены, ошибка возникнет только после их исчерпания; если они отключены — сразу.
Эту ошибку также можно обработать:
>>> from celery.utils.log import get_logger
>>> logger = get_logger(__name__)
>>> try:
... add.delay(2, 2)
... except add.OperationalError as exc:
... logger.exception('Sending task raised: %r', exc)
Примечание
В RabbitMQ эти ошибки указывают только на недоступность брокера. При достижении брокером пределов ресурсов сообщения могут быть отброшены без уведомления. Чтобы обнаруживать это, включите confirm_publish в broker_transport_options.
Данные, передаваемые между клиентами и рабочими процессами, необходимо сериализовать, поэтому каждое сообщение Celery содержит заголовок content_type, указывающий используемый метод сериализации.
По умолчанию используется сериализатор JSON, но его можно изменить с помощью параметра task_serializer, отдельно для каждой задачи или даже для каждого сообщения.
Встроенная поддержка предусмотрена для JSON, pickle, YAML и msgpack. Также можно добавить собственные сериализаторы, зарегистрировав их в реестре сериализаторов Kombu.
См. также
Раздел Сериализация сообщений в руководстве пользователя Kombu.
У каждого варианта есть свои преимущества и недостатки.
- json — JSON поддерживается многими языками программирования, теперь
-
является стандартной частью Python (начиная с версии 2.6) и достаточно быстро декодируется.
Главный недостаток JSON — ограниченный набор поддерживаемых типов данных: строки, Unicode, числа с плавающей точкой, логические значения, словари и списки. В частности, не поддерживаются десятичные дроби и даты.
Двоичные данные передаются с помощью кодирования Base64, что увеличивает размер передаваемых данных на 34% по сравнению с форматом кодирования, поддерживающим собственные двоичные типы.
Однако если ваши данные соответствуют указанным ограничениям и вам нужна поддержка разных языков, JSON, используемый по умолчанию, вероятно, будет лучшим выбором.
Подробнее см. на сайте http://json.org.
Примечание
(Из официальной документации Python: https://docs.python.org/3.6/library/json.html.) Ключи в парах «ключ — значение» JSON всегда имеют тип
str. При преобразовании словаря в JSON все его ключи приводятся к строкам. Поэтому после преобразования словаря в JSON и обратно полученный словарь может отличаться от исходного. То естьloads(dumps(x)) != x, если в x есть ключи нестрокового типа.Предупреждение
В сложных рабочих процессах, созданных с помощью Canvas: проектирование рабочих процессов, было замечено, что сериализатор JSON может значительно увеличивать размер сообщений из-за рекурсивных ссылок, что приводит к проблемам с ресурсами. Сериализатор pickle не подвержен этой проблеме и в таких случаях может быть предпочтительнее.
- pickle — если вам не требуется поддержка языков, отличных от
-
Python, кодирование pickle обеспечит поддержку всех встроенных типов данных Python (кроме экземпляров классов), уменьшит размер сообщений при передаче двоичных файлов и немного ускорит обработку по сравнению с JSON.
Подробнее см. в
pickle. - yaml — YAML обладает многими свойствами JSON,
-
но изначально поддерживает больше типов данных (в том числе даты, рекурсивные ссылки и т. д.).
Однако библиотеки YAML для Python значительно медленнее библиотек JSON.
Если вам нужен более широкий набор типов данных и совместимость между языками, YAML может подойти лучше, чем описанные выше варианты.
Чтобы использовать его, установите Celery следующей командой:
$ pipinstallcelery[yaml]
Подробнее см. на сайте http://yaml.org/.
- msgpack — msgpack — это двоичный формат сериализации, по возможностям близкий к JSON
-
Этот формат обеспечивает более эффективное сжатие, поэтому его разбор и кодирование выполняются быстрее, чем для JSON.
Чтобы использовать его, установите Celery следующей командой:
$ pipinstallcelery[msgpack]
Подробнее см. на сайте http://msgpack.org/.
Чтобы использовать пользовательский сериализатор, добавьте тип содержимого в accept_content. По умолчанию принимаются только сообщения JSON; задачи с другими заголовками содержимого отклоняются.
Сериализатор для отправки задачи выбирается в следующем порядке:
Параметр выполнения serializer.
Атрибут
Task.serializerПараметр
task_serializer.
Пример задания пользовательского сериализатора для отдельного вызова задачи:
>>> add.apply_async((10, 10), serializer='json')
Celery может сжимать сообщения с помощью следующих встроенных алгоритмов:
-
brotli
brotli оптимизирован для веба, в частности для небольших текстовых документов. Он особенно эффективен при передаче статического содержимого, например шрифтов и HTML-страниц.
Чтобы использовать его, установите Celery следующей командой:
$ pipinstallcelery[brotli]
-
bzip2
bzip2 создаёт файлы меньшего размера, чем gzip, но скорость сжатия и распаковки заметно ниже.
Убедитесь, что ваш исполняемый файл Python собран с поддержкой bzip2.
Если появляется следующая ошибка
ImportError:>>> import bz2 Traceback (most recent call last): File "<stdin>", line 1, in <module> ImportError: No module named 'bz2'
это означает, что версию Python нужно пересобрать с поддержкой bzip2.
-
gzip
gzip подходит для систем с ограниченным объёмом памяти, поскольку потребляет мало памяти. Его часто используют для создания файлов с расширением «.tar.gz».
Убедитесь, что ваш исполняемый файл Python собран с поддержкой gzip.
Если появляется следующая ошибка
ImportError:>>> import gzip Traceback (most recent call last): File "<stdin>", line 1, in <module> ImportError: No module named 'gzip'
это означает, что версию Python нужно пересобрать с поддержкой gzip.
-
lzma
lzma обеспечивает высокий коэффициент сжатия и высокую скорость сжатия и распаковки, но потребляет больше памяти.
Убедитесь, что ваш исполняемый файл Python собран с поддержкой lzma и что используется Python версии 3.3 или новее.
Если появляется следующая ошибка
ImportError:>>> import lzma Traceback (most recent call last): File "<stdin>", line 1, in <module> ImportError: No module named 'lzma'
это означает, что версию Python нужно пересобрать с поддержкой lzma.
Также можно установить обратный порт с помощью команды:
$ pipinstallcelery[lzma]
-
zlib
zlib — это библиотечная реализация алгоритма Deflate, поддерживающая в API как формат файлов gzip, так и облегчённый потоковый формат. Это важный компонент многих программных систем, включая ядро Linux и систему контроля версий Git.
Убедитесь, что ваш исполняемый файл Python собран с поддержкой zlib.
Если появляется следующая ошибка
ImportError:>>> import zlib Traceback (most recent call last): File "<stdin>", line 1, in <module> ImportError: No module named 'zlib'
это означает, что версию Python нужно пересобрать с поддержкой zlib.
-
zstd
zstd предназначен для сценариев сжатия в реальном времени: он обеспечивает скорость на уровне zlib и более высокий коэффициент сжатия. В его основе — очень быстрый этап энтропийного кодирования, реализованный библиотеками Huff0 и FSE.
Чтобы использовать его, установите Celery следующей командой:
$ pipinstallcelery[zstd]
Можно также создавать собственные алгоритмы сжатия и регистрировать их в kombu compression registry.
Алгоритм сжатия для отправки задачи выбирается в следующем порядке:
Параметр выполнения compression.
Атрибут
Task.compression.Атрибут
task_compression.
Пример задания алгоритма сжатия при вызове задачи:
>>> add.apply_async((2, 2), compression='zlib')
Можно управлять соединением вручную, создав издателя:
numbers = [(2, 2), (4, 4), (8, 8), (16, 16)]
results = []
with add.app.pool.acquire(block=True) as connection:
with add.get_publisher(connection) as publisher:
try:
for i, j in numbers:
res = add.apply_async((i, j), publisher=publisher)
results.append(res)
print([res.get() for res in results])
Однако этот пример гораздо удобнее выразить с помощью группы:
>>> from celery import group >>> numbers = [(2, 2), (4, 4), (8, 8), (16, 16)] >>> res = group(add.s(i, j) for i, j in numbers).apply_async() >>> res.get() [4, 8, 16, 32]
Celery может направлять задачи в разные очереди.
Простая маршрутизация (имя <-> имя) выполняется с помощью параметра queue:
add.apply_async(queue='priority.high')
Затем можно назначить рабочие процессы для очереди priority.high, используя аргумент -Q команды запуска рабочих процессов:
$ celery-Aprojworker-lINFO-Qcelery,priority.high
См. также
Не рекомендуется жёстко задавать имена очередей в коде. Лучше использовать маршрутизаторы конфигурации (task_routes).
Подробнее о маршрутизации см. в разделе Маршрутизация задач.
Хранение результатов можно включить или отключить с помощью параметра task_ignore_result или параметра ignore_result:
>>> result = add.apply_async((1, 2), ignore_result=True) >>> result.get() None >>> # Do not ignore result (default) ... >>> result = add.apply_async((1, 2), ignore_result=False) >>> result.get() 3
Чтобы сохранять в бэкенде результатов дополнительные метаданные о задаче, задайте для параметра result_extended значение True.
Примечание
Параметр result_extended определяет, какие расширенные метаданные задачи включает Celery, но автоматически не добавляет метаданные, специфичные для планировщика. Например, некоторые интеграции (например, https://pypi.org/project/django-celery-beat/ вместе с https://pypi.org/project/django-celery-results/) могут сохранять имя периодической задачи в бэкенде результатов только тогда, когда планировщик передаёт его в опубликованном сообщении.
При ручном вызове задач с помощью apply_async/delay контекст периодической задачи обычно отсутствует, если только вы не добавите его явно (например, через заголовки или свойства сообщения в параметрах apply_async). Например:
result = task.apply_async(
headers={"periodic_task_name": "task_name"},
)
См. также
Подробнее о задачах см. в разделе Задачи.
Дополнительные параметры
Эти параметры предназначены для опытных пользователей, которым нужны все возможности маршрутизации AMQP. Заинтересованные читатели могут обратиться к руководству по маршрутизации.
-
exchange
Имя обменника (или
kombu.entity.Exchange), в который нужно отправить сообщение. -
routing_key
Ключ маршрутизации, используемый для определения назначения.
-
priority
Число от 0 до 255, где 255 соответствует наивысшему приоритету.
Поддерживается в RabbitMQ и Redis (в Redis приоритеты расположены в обратном порядке: 0 — наивысший).
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/calling.html