Spec-Zone.ru › Celery

Canvas: проектирование рабочих процессов

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

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

signature() объединяет аргументы, именованные аргументы и параметры выполнения одного вызова задачи таким образом, что её можно передавать функциям, сериализовать и отправлять по сети.

  • Создать сигнатуру задачи add можно по её имени следующим образом:

    >>> from celery import signature
    >>> signature('tasks.add', args=(2, 2), countdown=10)
    tasks.add(2, 2)
    

    Арность сигнатуры этой задачи равна 2 (два аргумента): (2, 2); для параметра выполнения countdown задано значение 10.

  • или создать её с помощью метода signature задачи:

    >>> add.signature((2, 2), countdown=10)
    tasks.add(2, 2)
    
  • Есть также сокращённая запись с использованием позиционных аргументов:

    >>> add.s(2, 2)
    tasks.add(2, 2)
    
  • Поддерживаются и именованные аргументы:

    >>> add.s(2, 2, debug=True)
    tasks.add(2, 2, debug=True)
    
  • У любого экземпляра сигнатуры можно проверить значения различных полей:

    >>> s = add.signature((2, 2), {'debug': True}, countdown=10)
    >>> s.args
    (2, 2)
    >>> s.kwargs
    {'debug': True}
    >>> s.options
    {'countdown': 10}
    
  • Поддерживается «API вызова» для delay, apply_async и т. д., включая непосредственный вызов (__call__).

    Вызов сигнатуры выполнит задачу непосредственно в текущем процессе:

    >>> add(2, 2)
    4
    >>> add.s(2, 2)()
    4
    

    delay — наш любимый сокращённый способ вызвать apply_async с позиционными аргументами:

    >>> result = add.delay(2, 2)
    >>> result.get()
    4
    

    apply_async принимает те же аргументы, что и метод app.Task.apply_async():

    >>> add.apply_async(args, kwargs, **options)
    >>> add.signature(args, kwargs, **options).apply_async()
    
    >>> add.apply_async((2, 2), countdown=1)
    >>> add.signature((2, 2), countdown=1).apply_async()
    
  • С помощью s() нельзя задавать параметры, но это можно сделать при помощи цепочки вызовов set:

    >>> add.s(2, 2).set(countdown=1)
    proj.tasks.add(2, 2)
    

Сигнатура позволяет запустить задачу в рабочем процессе:

>>> add.s(2, 2).delay()
>>> add.s(2, 2).apply_async(countdown=1)

Или вызвать её непосредственно в текущем процессе:

>>> add.s(2, 2)()
4

Передача дополнительных аргументов, именованных аргументов или параметров в apply_async/delay создаёт частичные вызовы:

  • Все добавленные аргументы будут поставлены перед аргументами сигнатуры:

    >>> partial = add.s(2)          # incomplete signature
    >>> partial.delay(4)            # 4 + 2
    >>> partial.apply_async((4,))  # same
    

    Примечание

    Дополнительные аргументы, переданные в delay/apply_async, добавляются в начало аргументов сигнатуры. Поскольку add коммутативна, порядок может быть неочевиден. Некоммутативная задача, например subtract(x, y) -> x - y, наглядно это демонстрирует:

    @app.task
    def subtract(x, y):
        return x - y
    
    partial = subtract.s(10)    # incomplete: second arg only
    partial.delay(30)           # -> subtract(30, 10) = 20
    

    Здесь delay(30) добавляет 30 первым аргументом, в результате получается subtract(30, 10), а не subtract(10, 30).

  • Все добавленные именованные аргументы будут объединены с именованными аргументами сигнатуры; при совпадении приоритет имеют новые аргументы:

    >>> s = add.s(2, 2)
    >>> s.delay(debug=True)                    # -> add(2, 2, debug=True)
    >>> s.apply_async(kwargs={'debug': True})  # same
    
  • Все добавленные параметры будут объединены с параметрами сигнатуры; при совпадении приоритет имеют новые параметры:

    >>> s = add.signature((2, 2), countdown=10)
    >>> s.apply_async(countdown=1)  # countdown is now 1
    

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

>>> s = add.s(2)
proj.tasks.add(2)

>>> s.clone(args=(4,), kwargs={'debug': True})
proj.tasks.add(4, 2, debug=True)

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

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

>>> add.apply_async((2, 2), link=reset_buffers.signature(immutable=True))

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

>>> add.apply_async((2, 2), link=reset_buffers.si())

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

Примечание

В этом руководстве я иногда использую для сигнатур оператор-префикс ~. В рабочем коде использовать его, вероятно, не следует, но это удобное сокращение для экспериментов в оболочке Python:

>>> ~sig

>>> # is the same as
>>> sig.delay().get()

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

К любой задаче можно добавить обратный вызов с помощью аргумента link метода apply_async:

add.apply_async((2, 2), link=other_task.s())

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

Как я уже упоминал, любые аргументы, добавленные к сигнатуре, будут поставлены перед аргументами, заданными самой сигнатурой!

Если имеется сигнатура:

>>> sig = add.s(10)

то вызов sig.delay(result) превращается в:

>>> add.apply_async(args=(result, 10))

…

Теперь вызовем задачу add с обратным вызовом, используя частичные аргументы:

>>> add.apply_async((2, 2), link=add.s(8))

Как и ожидалось, сначала будет запущена задача, вычисляющая 2 + 2, а затем другая задача, вычисляющая 4 + 8.

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

Обзор

  • group

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

  • chain

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

  • chord

    Примитив chord похож на group, но с обратным вызовом. Chord состоит из группы заголовка и тела, где тело — это задача, которая должна выполниться после завершения всех задач заголовка.

  • map

    Примитив map работает подобно встроенной функции map, но создаёт временную задачу, которой передаётся список аргументов. Например, task.map([1, 2]) — в результате вызывается одна задача, которой аргументы передаются по порядку, так что результат будет следующим:

    res = [task(1), task(2)]
    
  • starmap

    Работает точно так же, как map, но аргументы передаются как *args. Например, add.starmap([(2, 2), (4, 4)]) приводит к вызову одной задачи:

    res = [add(2, 2), add(4, 4)]
    
  • chunks

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

    >>> items = zip(range(1000), range(1000))  # 1000 items
    >>> add.chunks(items, 10)
    

    разделит список элементов на фрагменты по 10 элементов, создав 100 задач (каждая последовательно обработает 10 элементов).

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

Вот несколько примеров:

  • Простая цепочка

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

    >>> from celery import chain
    
    >>> # 2 + 2 + 4 + 8
    >>> res = chain(add.s(2, 2), add.s(4), add.s(8))()
    >>> res.get()
    16
    

    То же самое можно записать с помощью операторов конвейера:

    >>> (add.s(2, 2) | add.s(4) | add.s(8))().get()
    16
    
  • Неизменяемые сигнатуры

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

    В этом случае можно пометить сигнатуру как неизменяемую, чтобы её аргументы нельзя было изменить:

    >>> add.signature((2, 2), immutable=True)
    

    Для этого есть также сокращённая запись .si(), и это предпочтительный способ создания сигнатур:

    >>> add.si(2, 2)
    

    Теперь вместо этого можно создать цепочку независимых задач:

    >>> res = (add.si(2, 2) | add.si(4, 4) | add.si(8, 8))()
    >>> res.get()
    16
    
    >>> res.parent.get()
    8
    
    >>> res.parent.parent.get()
    4
    
  • Простая группа

    Можно легко создать группу задач для параллельного выполнения:

    >>> from celery import group
    >>> res = group(add.s(i, i) for i in range(10))()
    >>> res.get(timeout=1)
    [0, 2, 4, 6, 8, 10, 12, 14, 16, 18]
    
  • Простой chord

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

    >>> from celery import chord
    >>> res = chord((add.s(i, i) for i in range(10)), tsum.s())()
    >>> res.get()
    90
    

    В приведённом выше примере создаются 10 задач, которые запускаются параллельно. Когда все они завершатся, возвращённые значения объединяются в список и передаются задаче tsum.

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

    >>> chord((import_contact.s(c) for c in contacts),
    ...       notify_complete.si(import_id)).apply_async()
    

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

  • Комбинируем и поражаемся

    Цепочки тоже могут быть частичными:

    >>> c1 = (add.s(4) | mul.s(8))
    
    # (16 + 4) * 8
    >>> res = c1(16)
    >>> res.get()
    160
    

    это позволяет объединять цепочки:

    # ((4 + 16) * 2 + 4) * 8
    >>> c2 = (add.s(4, 16) | mul.s(2) | (add.s(4) | mul.s(8)))
    
    >>> res = c2()
    >>> res.get()
    352
    

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

    >>> c3 = (group(add.s(i, i) for i in range(10)) | tsum.s())
    >>> res = c3()
    >>> res.get()
    90
    

    Группы и chord также принимают частичные аргументы, поэтому в цепочке возвращённое значение предыдущей задачи передаётся всем задачам группы:

    >>> new_user_workflow = (create_user.s() | group(
    ...                      import_contacts.s(),
    ...                      send_welcome_email.s()))
    ... new_user_workflow.delay(username='artv',
    ...                         first='Art',
    ...                         last='Vandelay',
    ...                         email='art@vandelay.com')
    

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

    >>> res = (add.s(4, 4) | group(add.si(i, i) for i in range(10)))()
    >>> res.get()
    [0, 2, 4, 6, 8, 10, 12, 14, 16, 18]
    
    >>> res.parent.get()
    8
    

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

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

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

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

>>> res = add.apply_async((2, 2), link=mul.s(16))
>>> res.get()
4

Связанная задача получит результат родительской задачи в качестве первого аргумента. В приведённом выше случае, когда результат равен 4, получится mul(4, 16).

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

>>> res.children
[<AsyncResult: 8c350acf-519d-4553-8a53-4ad3a5c5aeb4>]

>>> res.children[0].get()
64

У экземпляра результата есть также метод collect(), который рассматривает результат как граф и позволяет перебирать результаты:

>>> list(res.collect())
[(<AsyncResult: 7b720856-dc5f-4415-9134-5c89def5664e>, 4),
 (<AsyncResult: 8c350acf-519d-4553-8a53-4ad3a5c5aeb4>, 64)]

По умолчанию collect() вызывает исключение IncompleteStream, если граф сформирован не полностью (одна из задач ещё не завершилась). Однако можно получить и промежуточное представление графа:

>>> for result, value in res.collect(intermediate=True):
....

Можно связать сколько угодно задач; связывать можно и сигнатуры:

>>> s = add.s(2, 2)
>>> s.link(mul.s(4))
>>> s.link(log_result.s())

Также можно добавлять обратные вызовы при ошибке с помощью метода on_error:

>>> add.s(2, 2).on_error(log_error.s()).delay()

В результате при выполнении сигнатуры будет вызван .apply_async:

>>> add.apply_async((2, 2), link_error=log_error.s())

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

Вот пример errback:

import os

from proj.celery import app

@app.task
def log_error(request, exc, traceback):
    with open(os.path.join('/var/errors', request.id), 'a') as fh:
        print('--\n\n{0} {1} {2}'.format(
            request.id, exc, traceback), file=fh)

Чтобы ещё проще связывать задачи, существует специальная сигнатура chain, позволяющая объединять задачи в цепочки:

>>> from celery import chain
>>> from proj.tasks import add, mul

>>> # (4 + 4) * 8 * 10
>>> res = chain(add.s(4, 4), mul.s(8), mul.s(10))
proj.tasks.add(4, 4) | proj.tasks.mul(8) | proj.tasks.mul(10)

Вызов chain выполнит задачи в текущем процессе и вернёт результат последней задачи в цепочке:

>>> res = chain(add.s(4, 4), mul.s(8), mul.s(10))()
>>> res.get()
640

Она также задаёт атрибуты parent, поэтому можно пройти по цепочке вверх и получить промежуточные результаты:

>>> res.parent.get()
64

>>> res.parent.parent.get()
8

>>> res.parent.parent
<AsyncResult: eeaad925-6778-4ad1-88c8-b2a63d017933>

Цепочки также можно создавать с помощью оператора | (конвейера):

>>> (add.s(2, 2) | mul.s(8) | mul.s(10)).apply_async()

Идентификатор задачи

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

Цепочка наследует идентификатор задачи последней задачи в цепочке.

Графы

Кроме того, с графом результатов можно работать как с объектом DependencyGraph:

>>> res = chain(add.s(4, 4), mul.s(8), mul.s(10))()

>>> res.parent.parent.graph
285fa253-fcf8-42ef-8b95-0078897e83e6(1)
    463afec2-5ed4-4036-b22d-ba067ec64f52(0)
872c3995-6fa0-46ca-98c2-5a19155afcf0(2)
    285fa253-fcf8-42ef-8b95-0078897e83e6(1)
        463afec2-5ed4-4036-b22d-ba067ec64f52(0)

Можно даже преобразовать эти графы в формат dot:

>>> with open('graph.dot', 'w') as fh:
...     res.parent.parent.graph.to_dot(fh)

и создать изображения:

$ dot-Tpnggraph.dot-ograph.png
../_images/result_graph.png

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

Примечание

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

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

Функция group принимает список сигнатур:

>>> from celery import group
>>> from proj.tasks import add

>>> group(add.s(2, 2), add.s(4, 4))
(proj.tasks.add(2, 2), proj.tasks.add(4, 4))

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

>>> g = group(add.s(2, 2), add.s(4, 4))
>>> res = g()
>>> res.get()
[4, 8]

Группа также поддерживает итераторы:

>>> group(add.s(i, i) for i in range(100))()

Группа является объектом-сигнатурой, поэтому её можно комбинировать с другими сигнатурами.

Обратные вызовы группы и обработка ошибок

К группам также можно привязывать сигнатуры обратных вызовов и errback. Однако такое поведение может оказаться неожиданным, поскольку группы не являются настоящими задачами и просто передают связанные задачи своим вложенным сигнатурам. Это значит, что возвращаемые значения группы не собираются для передачи связанной сигнатуре обратного вызова. Кроме того, связь с задачей не гарантирует, что она запустится только после завершения всех задач группы. Например, следующий фрагмент с простой задачей add(a, b) содержит ошибку: связанная сигнатура add.s() не получит итоговый результат группы, как можно было бы ожидать.

>>> g = group(add.s(2, 2), add.s(4, 4))
>>> g.link(add.s())
>>> res = g()
[4, 8]

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

Errback группы также передаются вложенным сигнатурам. Поэтому errback, привязанный один раз, может быть вызван несколько раз, если завершатся с ошибкой несколько задач группы. Например, в следующем фрагменте с задачей fail(), вызывающей исключение, сигнатура log_error() может вызываться для каждой завершившейся с ошибкой задачи группы.

>>> g = group(fail.s(), fail.s())
>>> g.link_error(log_error.s())
>>> res = g()

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

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

Результаты группы

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

>>> from celery import group
>>> from tasks import add

>>> job = group([
...             add.s(2, 2),
...             add.s(4, 4),
...             add.s(8, 8),
...             add.s(16, 16),
...             add.s(32, 32),
... ])

>>> result = job.apply_async()

>>> result.ready()  # have all subtasks completed?
True
>>> result.successful() # were all subtasks successful?
True
>>> result.get()
[4, 8, 16, 32, 64]

GroupResult принимает список экземпляров AsyncResult и работает с ними так, как если бы это была одна задача.

Поддерживаются следующие операции:

  • successful()

    Возвращает True, если все подзадачи успешно завершились (например, не вызвали исключение).

  • failed()

    Возвращает True, если какая-либо подзадача завершилась с ошибкой.

  • waiting()

    Возвращает True, если какая-либо подзадача ещё не готова.

  • ready()

    Возвращает True, если все подзадачи готовы.

  • completed_count()

    Возвращает количество завершённых подзадач. Здесь complete означает successful. Иными словами, этот метод возвращает количество задач successful.

  • revoke()

    Отзывает все подзадачи.

  • join()

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

Разворачивание группы

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

>>> from celery import chain, group
>>> from tasks import add
>>> chain(add.s(2, 2), group(add.s(1)), add.s(1))
add(2, 2) | add(1) | add(1)
>>> chain(add.s(2, 2), group(add.s(1), add.s(2)), add.s(1))
add(2, 2) | %add((add(1), add(2)), 1)

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

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

В Celery 4.x из-за ошибки следующая группа не разворачивалась в цепочку, а вместо этого canvas преобразовывался в chord.

>>> from celery import chain, group
>>> from tasks import add
>>> chain(group(add.s(1, 1)), add.s(2))
%add([add(1, 1)], 2)

В Celery 5.x эта ошибка исправлена, и группа корректно разворачивается в отдельную сигнатуру.

>>> from celery import chain, group
>>> from tasks import add
>>> chain(group(add.s(1, 1)), add.s(2))
add(1, 1) | add(2)

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

Примечание

Задачи в chord не должны игнорировать свои результаты. Если для любой задачи (в заголовке или теле) в chord отключён бэкенд результатов, прочитайте раздел «Важные замечания». Chord пока не поддерживается с бэкендом результатов RPC.

Chord — это задача, которая выполняется только после завершения всех задач группы.

Вычислим сумму выражения 1 + 1 + 2 + 2 + 3 + 3 ... n + n с точностью до ста цифр.

Сначала понадобятся две задачи: add() и tsum() (sum() — стандартная функция):

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

@app.task
def tsum(numbers):
    return sum(numbers)

Теперь можно использовать chord, чтобы параллельно вычислить каждое слагаемое, а затем сложить полученные числа:

>>> from celery import chord
>>> from tasks import add, tsum

>>> chord(add.s(i, i)
...       for i in range(100))(tsum.s()).get()
9900

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

>>> sum(i + i for i in range(100))

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

Разберём выражение chord:

>>> callback = tsum.s()
>>> header = [add.s(i, i) for i in range(100)]
>>> result = chord(header)(callback)
>>> result.get()
9900

Помните, что обратный вызов можно выполнить только после возврата всех задач заголовка. Каждый этап заголовка выполняется как отдельная задача параллельно, возможно, на разных узлах. Затем обратному вызову передаётся возвращённое значение каждой задачи заголовка. Идентификатор задачи, возвращаемый chord(), — это идентификатор обратного вызова; по нему можно дождаться его завершения и получить итоговое возвращаемое значение (но помните: задача никогда не должна ожидать завершения других задач).

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

Что произойдёт, если одна из задач вызовет исключение?

Результат обратного вызова chord перейдёт в состояние ошибки, а в качестве ошибки будет задано исключение ChordError:

>>> c = chord([add.s(4, 4), raising_task.s(), add.s(8, 8)])
>>> result = c()
>>> result.get()
Traceback (most recent call last):
  File "<stdin>", line 1, in <module>
  File "*/celery/result.py", line 120, in get
interval=interval)
  File "*/celery/backends/amqp.py", line 150, in wait_for
raise meta['result']
celery.exceptions.ChordError: Dependency 97de6f3f-ea67-4517-a21c-d867c61fcb47
    raised ValueError('something something',)

Трассировка стека может различаться в зависимости от используемого бэкенда результатов, но описание ошибки содержит идентификатор задачи, завершившейся с ошибкой, и строковое представление исходного исключения. Исходную трассировку стека также можно найти в result.traceback.

Обратите внимание, что остальные задачи всё равно будут выполнены: третья задача (add.s(8, 8)) запустится, несмотря на ошибку в средней задаче. Кроме того, ChordError указывает только на задачу, которая первой завершилась с ошибкой по времени, а не по порядку в группе заголовка.

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

@app.task
def on_chord_error(request, exc, traceback):
    print('Task {0!r} raised error: {1!r}'.format(request.id, exc))
>>> c = (group(add.s(i, i) for i in range(10)) |
...      tsum.s().on_error(on_chord_error.s())).delay()

К chord можно привязывать сигнатуры обратных вызовов и errback, что позволяет решить некоторые проблемы, возникающие при привязке сигнатур к группам. В этом случае указанная сигнатура привязывается к телу chord; обратные вызовы будут корректно вызваны один раз после завершения тела, а errback — один раз, если завершится с ошибкой какая-либо задача заголовка или тела chord.

Это поведение можно изменить с помощью флага task_allow_error_cb_on_chord_header, чтобы обрабатывать ошибки заголовка chord. При включении этого флага errback тела будет вызван для заголовка chord (поведение по умолчанию) и для каждой задачи заголовка chord, завершившейся с ошибкой.

Важные замечания

Задачи в chord не должны игнорировать свои результаты. На практике это означает, что для использования chord необходимо включить result_backend. Кроме того, если в конфигурации для task_ignore_result задано значение True, обязательно определите задачи, используемые в chord, с параметром ignore_result=False. Это относится как к подклассам Task, так и к задачам, объявленным с помощью декоратора.

Пример подкласса Task:

class MyTask(Task):
    ignore_result = False

Пример задачи с декоратором:

@app.task(ignore_result=False)
def another_task(project):
    do_something()

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

Пример реализации:

from celery import maybe_signature

@app.task(bind=True)
def unlock_chord(self, group, callback, interval=1, max_retries=None):
    if group.ready():
        return maybe_signature(callback).delay(group.join())
    raise self.retry(countdown=interval, max_retries=max_retries)

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

Подход Redis, Memcached и DynamoDB значительно лучше, но его непросто реализовать в других бэкендах (предложения приветствуются!).

Примечание

Chord некорректно работает с Redis до версии 2.2; для его использования необходимо обновиться как минимум до redis-server 2.2.

Примечание

Если вы используете chord с бэкендом результатов Redis и переопределяете метод Task.after_return(), обязательно вызовите метод родительского класса, иначе обратный вызов chord не будет выполнен.

def after_return(self, *args, **kwargs):
    do_something()
    super().after_return(*args, **kwargs)

map и starmap — это встроенные задачи, которые вызывают указанную задачу для каждого элемента последовательности.

От group они отличаются тем, что:

  • отправляется только одно сообщение задачи.

  • операция выполняется последовательно.

Например, использование map:

>>> from proj.tasks import add

>>> ~tsum.map([list(range(10)), list(range(100))])
[45, 4950]

эквивалентно задаче, выполняющей:

@app.task
def temp():
    return [tsum(range(10)), tsum(range(100))]

а использование starmap:

>>> ~add.starmap(zip(range(10), range(10)))
[0, 2, 4, 6, 8, 10, 12, 14, 16, 18]

эквивалентно задаче, выполняющей:

@app.task
def temp():
    return [add(i, i) for i in range(10)]

И map, и starmap являются объектами-сигнатурами, поэтому их можно использовать как другие сигнатуры и объединять в группы и т. д. Например, чтобы запустить starmap через 10 секунд:

>>> add.starmap(zip(range(10), range(10))).apply_async(countdown=10)

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

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

Чтобы создать сигнатуру chunks, можно использовать app.Task.chunks():

>>> add.chunks(zip(range(100), range(100)), 10)

Как и в случае с group, сообщения для фрагментов отправляются из текущего процесса при вызове:

>>> from proj.tasks import add

>>> res = add.chunks(zip(range(100), range(100)), 10)()
>>> res.get()
[[0, 2, 4, 6, 8, 10, 12, 14, 16, 18],
 [20, 22, 24, 26, 28, 30, 32, 34, 36, 38],
 [40, 42, 44, 46, 48, 50, 52, 54, 56, 58],
 [60, 62, 64, 66, 68, 70, 72, 74, 76, 78],
 [80, 82, 84, 86, 88, 90, 92, 94, 96, 98],
 [100, 102, 104, 106, 108, 110, 112, 114, 116, 118],
 [120, 122, 124, 126, 128, 130, 132, 134, 136, 138],
 [140, 142, 144, 146, 148, 150, 152, 154, 156, 158],
 [160, 162, 164, 166, 168, 170, 172, 174, 176, 178],
 [180, 182, 184, 186, 188, 190, 192, 194, 196, 198]]

А вызов .apply_async создаст отдельную задачу, чтобы отдельные задачи выполнялись в рабочем процессе:

>>> add.chunks(zip(range(100), range(100)), 10).apply_async()

Фрагменты также можно преобразовать в группу:

>>> group = add.chunks(zip(range(100), range(100)), 10).group()

а затем с помощью skew увеличивать countdown каждой задачи на единицу:

>>> group.skew(start=1, stop=10)()

Это означает, что для первой задачи countdown будет равен одной секунде, для второй — двум секундам и так далее.

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

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

Например,

>>> sig1 = add.si(2, 2)
>>> sig1_res = sig1.freeze()
>>> g = group(sig1, add.si(3, 3))
>>> g.stamp(stamp='your_custom_stamp')
>>> res = g.apply_async()
>>> res.get(timeout=TIMEOUT)
[4, 6]
>>> sig1_res._get_task_meta()['stamp']
['your_custom_stamp']

инициализирует группу g и помечает её компоненты меткой your_custom_stamp.

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

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

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

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

class InGroupVisitor(StampingVisitor):
    def __init__(self):
        self.in_group = False

    def on_group_start(self, group, **headers) -> dict:
        self.in_group = True
        return {"in_group": [self.in_group], "stamped_headers": ["in_group"]}

    def on_group_end(self, group, **headers) -> None:
        self.in_group = False

    def on_chain_start(self, chain, **headers) -> dict:
        return {"in_group": [self.in_group], "stamped_headers": ["in_group"]}

    def on_signature(self, sig, **headers) -> dict:
        return {"in_group": [self.in_group], "stamped_headers": ["in_group"]}

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

class MonitoringIdStampingVisitor(StampingVisitor):
    def on_signature(self, sig, **headers) -> dict:
        return {'monitoring_id': uuid4().hex}

Важно

Ключ stamped_headers в словаре, возвращаемом on_signature() (или любым другим методом посетителя), является необязательным:

# Approach 1: Without stamped_headers - ALL keys are treated as stamps
def on_signature(self, sig, **headers) -> dict:
    return {'monitoring_id': uuid4().hex}  # monitoring_id becomes a stamp

# Approach 2: With stamped_headers - ONLY listed keys are stamps
def on_signature(self, sig, **headers) -> dict:
    return {
        'monitoring_id': uuid4().hex,      # This will be a stamp
        'other_data': 'value',             # This will NOT be a stamp
        'stamped_headers': ['monitoring_id']  # Only monitoring_id is stamped
    }

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

Далее рассмотрим, как использовать пример посетителя маркировки MonitoringIdStampingVisitor.

sig_example = signature('t1')
sig_example.stamp(visitor=MonitoringIdStampingVisitor())

group_example = group([signature('t1'), signature('t2')])
group_example.stamp(visitor=MonitoringIdStampingVisitor())

chord_example = chord([signature('t1'), signature('t2')], signature('t3'))
chord_example.stamp(visitor=MonitoringIdStampingVisitor())

chain_example = chain(signature('t1'), group(signature('t2'), signature('t3')), signature('t4'))
chain_example.stamp(visitor=MonitoringIdStampingVisitor())

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

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

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

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

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

class CustomStampingVisitor(StampingVisitor):
    def on_signature(self, sig, **headers) -> dict:
        # 'header' will automatically be treated as a stamped header
        # without needing to specify 'stamped_headers': ['header']
        return {'header': 'value'}

    def on_callback(self, callback, **header) -> dict:
        # 'on_callback' will automatically be treated as a stamped header
        return {'on_callback': True}

    def on_errback(self, errback, **header) -> dict:
        # 'on_errback' will automatically be treated as a stamped header
        return {'on_errback': True}

Этот пользовательский посетитель маркировки пометит сигнатуру, обратные вызовы и обработчики ошибок значением {'header': 'value'}, а обратные вызовы и обработчики ошибок — значениями {'on_callback': True} и {'on_errback': True} соответственно, как показано ниже.

c = chord([add.s(1, 1), add.s(2, 2)], xsum.s())
callback = signature('sig_link')
errback = signature('sig_link_error')
c.link(callback)
c.link_error(errback)
c.stamp(visitor=CustomStampingVisitor())

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

>>> c.options
{'header': 'value', 'stamped_headers': ['header']}
>>> c.tasks.tasks[0].options
{'header': 'value', 'stamped_headers': ['header']}
>>> c.tasks.tasks[1].options
{'header': 'value', 'stamped_headers': ['header']}
>>> c.body.options
{'header': 'value', 'stamped_headers': ['header']}
>>> c.body.options['link'][0].options
{'header': 'value', 'on_callback': True, 'stamped_headers': ['header', 'on_callback']}
>>> c.body.options['link_error'][0].options
{'header': 'value', 'on_errback': True, 'stamped_headers': ['header', 'on_errback']}

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

Spec-Zone.ru

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