Spec-Zone.ru › Celery

Задачи маршрутизации

Примечание

Альтернативные концепции маршрутизации, такие как topic и fanout, доступны не для всех транспортов; см. таблицу сравнения транспортов.

Самый простой способ настроить маршрутизацию — использовать параметр task_create_missing_queues (включён по умолчанию).

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

Предположим, у вас есть два сервера, x и y, которые обрабатывают обычные задачи, и один сервер z, который обрабатывает только задачи, связанные с лентами. Можно использовать такую конфигурацию:

task_routes = {'feed.tasks.import_feed': {'queue': 'feeds'}}

При включённом маршруте задачи импорта лент будут направляться в очередь “feeds”, а все остальные задачи — в очередь по умолчанию (по историческим причинам называемую “celery”).

Также можно использовать сопоставление с шаблонами glob или даже регулярные выражения, чтобы сопоставлять все задачи в пространстве имён feed.tasks:

app.conf.task_routes = {'feed.tasks.*': {'queue': 'feeds'}}

Если важен порядок сопоставления шаблонов, следует указать маршрутизатор в формате items:

task_routes = ([
    ('feed.tasks.*', {'queue': 'feeds'}),
    ('web.tasks.*', {'queue': 'web'}),
    (re.compile(r'(video|image)\.tasks\..*'), {'queue': 'media'}),
],)

Примечание

Параметр task_routes может быть словарём или списком объектов маршрутизатора, поэтому в этом случае нужно указать параметр в виде кортежа, содержащего список.

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

user@z:/$ celery-Aprojworker-Qfeeds

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

user@z:/$ celery-Aprojworker-Qfeeds,celery

Изменить имя очереди по умолчанию можно с помощью следующей конфигурации:

app.conf.task_default_queue = 'default'

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

Очередь с именем “video” будет создана с такими параметрами:

{'exchange':'video',
'exchange_type':'direct',
'routing_key':'video'}

Бэкенды, не использующие AMQP, например Redis или SQS, не поддерживают обменники, поэтому требуют, чтобы обменник имел то же имя, что и очередь. Благодаря такому решению конфигурация будет работать и с ними.

Предположим, у вас есть два сервера, x и y, которые обрабатывают обычные задачи, и один сервер z, который обрабатывает только задачи, связанные с лентами. Можно использовать такую конфигурацию:

from kombu import Queue

app.conf.task_default_queue = 'default'
app.conf.task_queues = (
    Queue('default',    routing_key='task.#'),
    Queue('feed_tasks', routing_key='feed.#'),
)
app.conf.task_default_exchange = 'tasks'
app.conf.task_default_exchange_type = 'topic'
app.conf.task_default_routing_key = 'task.default'

task_queues — это список экземпляров Queue. Если для ключа не заданы значения обменника или типа обменника, они будут взяты из параметров task_default_exchange и task_default_exchange_type.

Чтобы направить задачу в очередь feed_tasks, можно добавить запись в параметр task_routes:

task_routes = {
        'feeds.tasks.import_feed': {
            'queue': 'feed_tasks',
            'routing_key': 'feed.import',
        },
}

Это также можно переопределить с помощью аргумента routing_key для Task.apply_async() или send_task():

>>> from feeds.tasks import import_feed
>>> import_feed.apply_async(args=['http://cnn.com/rss'],
...                         queue='feed_tasks',
...                         routing_key='feed.import')

Чтобы сервер z получал задачи только из очереди лент, запустите его с параметром celery worker -Q:

user@z:/$ celery-Aprojworker-Qfeed_tasks--hostname=z@%h

Серверы x и y должны быть настроены на получение задач из очереди по умолчанию:

user@x:/$ celery-Aprojworker-Qdefault--hostname=x@%h
user@y:/$ celery-Aprojworker-Qdefault--hostname=y@%h

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

user@z:/$ celery-Aprojworker-Qfeed_tasks,default--hostname=z@%h

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

from kombu import Exchange, Queue

app.conf.task_queues = (
    Queue('feed_tasks',    routing_key='feed.#'),
    Queue('regular_tasks', routing_key='task.#'),
    Queue('image_tasks',   exchange=Exchange('mediatasks', type='direct'),
                           routing_key='image.compress'),
)

Если эти термины вам непонятны, ознакомьтесь с AMQP.

См. также

Помимо раздела Приоритеты сообщений Redis ниже, ознакомьтесь со статьёй Rabbits and Warrens — отличной публикацией об очередях и обменниках. Также есть руководство CloudAMQP. Пользователям RabbitMQ может пригодиться FAQ RabbitMQ.

поддерживаемые транспорты:

RabbitMQ

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

Очереди можно настроить для поддержки приоритетов, задав аргумент x-max-priority:

from kombu import Exchange, Queue

app.conf.task_queues = [
    Queue('tasks', Exchange('tasks'), routing_key='tasks',
          queue_arguments={'x-max-priority': 10}),
]

Значение по умолчанию для всех очередей можно задать с помощью параметра task_queue_max_priority:

app.conf.task_queue_max_priority = 10

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

app.conf.task_default_priority = 5
поддерживаемые транспорты:

Redis

Хотя транспорт Celery для Redis учитывает поле приоритета, сам Redis не поддерживает понятие приоритетов. Прежде чем пытаться реализовать приоритеты с помощью Redis, прочитайте это примечание: вы можете столкнуться с неожиданным поведением.

Чтобы начать планировать задачи с учётом приоритетов, необходимо настроить параметр транспорта queue_order_strategy.

app.conf.broker_transport_options = {
    'queue_order_strategy': 'priority',
}

Поддержка приоритетов реализована путём создания n списков для каждой очереди. Это означает, что, хотя существует 10 уровней приоритета (0–9), по умолчанию они объединяются в 4 уровня для экономии ресурсов. Таким образом, очередь с именем celery фактически будет разделена на 4 очереди.

Очередь с наивысшим приоритетом будет называться celery, а к именам остальных очередей будет добавлен разделитель (по умолчанию x06x16) и номер их приоритета.

['celery', 'celery\x06\x163', 'celery\x06\x166', 'celery\x06\x169']

Чтобы использовать больше уровней приоритета или другой разделитель, задайте параметры транспорта priority_steps и sep:

app.conf.broker_transport_options = {
    'priority_steps': list(range(10)),
    'sep': ':',
    'queue_order_strategy': 'priority',
}

Приведённая выше конфигурация создаст очереди со следующими именами:

['celery', 'celery:1', 'celery:2', 'celery:3', 'celery:4', 'celery:5', 'celery:6', 'celery:7', 'celery:8', 'celery:9']

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

Сообщение состоит из заголовков и тела. Celery использует заголовки для хранения типа содержимого сообщения и его кодировки. Тип содержимого обычно указывает формат сериализации, использованный для сериализации сообщения. Тело содержит имя выполняемой задачи, идентификатор задачи (UUID), аргументы, с которыми её нужно вызвать, и дополнительные метаданные — например, число повторных попыток или время ETA.

Пример сообщения задачи в виде словаря Python:

{'task':'myapp.tasks.add',
'id':'54086c5e-6193-4575-8308-dbab76798756',
'args':[4,4],
'kwargs':{}}

Клиент, отправляющий сообщения, обычно называют издателем или производителем, а получающую сообщения сущность — потребителем.

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

В материалах, посвящённых AMQP, эти термины встречаются довольно часто.

  1. Сообщения отправляются в обменники.

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

  3. Сообщение ждёт в очереди, пока его кто-нибудь не получит.

  4. Сообщение удаляется из очереди после подтверждения его получения.

Чтобы отправлять и получать сообщения, необходимо выполнить следующие действия:

  1. Создать обменник

  2. Создать очередь

  3. Привязать очередь к обменнику.

Celery автоматически создаёт сущности, необходимые для работы очередей из task_queues (если только для очереди не задан параметр auto_declare со значением False).

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

from kombu import Exchange, Queue

app.conf.task_queues = (
    Queue('default', Exchange('default'), routing_key='default'),
    Queue('videos',  Exchange('media'),   routing_key='media.video'),
    Queue('images',  Exchange('media'),   routing_key='media.image'),
)
app.conf.task_default_queue = 'default'
app.conf.task_default_exchange_type = 'direct'
app.conf.task_default_routing_key = 'default'

Тип обменника определяет, как сообщения маршрутизируются через него. Стандарт определяет типы обменников direct, topic, fanout и headers. Также для RabbitMQ доступны нестандартные типы обменников в виде подключаемых модулей, например модуль last-value-cache Майкла Бриджена.

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

Обменники типа topic сопоставляют ключи маршрутизации, состоящие из слов, разделённых точками, и используют подстановочные символы: * (соответствует одному слову) и # (соответствует нулю или более слов).

Для таких ключей маршрутизации, как usa.news, usa.weather, norway.news и norway.weather, привязки могут выглядеть так: *.news (все новости), usa.# (все материалы из США) или usa.weather (все материалы о погоде в США).

exchange.declare(exchange_name, type, passive,
durable, auto_delete, internal)

Объявляет обменник по имени.

См. amqp:Channel.exchange_declare.

Именованные аргументы:
  • passive — пассивный режим означает, что обменник не будет создан; этот параметр можно использовать, чтобы проверить, существует ли обменник.

  • durable — долговечные обменники сохраняются (то есть переживают перезапуск брокера).

  • auto_delete — обменник будет удалён брокером, когда его перестанут использовать очереди.

queue.declare(queue_name, passive, durable, exclusive, auto_delete)

Объявляет очередь по имени.

См. amqp:Channel.queue_declare

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

queue.bind(queue_name, exchange_name, routing_key)

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

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

См. amqp:Channel.queue_bind

queue.delete(name, if_unused=False, if_empty=False)

Удаляет очередь и её привязку.

См. amqp:Channel.queue_delete

exchange.delete(name, if_unused=False)

Удаляет обменник.

См. amqp:Channel.exchange_delete

Примечание

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

В Celery есть инструмент celery amqp для работы с API AMQP из командной строки. Он предоставляет доступ к административным задачам, например созданию и удалению очередей и обменников, очистке очередей или отправке сообщений. Его также можно использовать с брокерами, не поддерживающими AMQP, но их реализации могут не поддерживать все команды.

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

$ celery-Aprojamqp
-> connecting to amqp://guest@localhost:5672/.
-> connected.
1>

Здесь 1> — это приглашение командной строки. Число 1 — количество выполненных команд. Введите help, чтобы увидеть список доступных команд. Также поддерживается автодополнение: начните вводить команду и нажмите клавишу tab, чтобы увидеть возможные варианты.

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

$ celery-Aprojamqp
1> exchange.declare testexchange direct
ok.
2> queue.declare testqueue
ok. queue:testqueue messages:0 consumers:0.
3> queue.bind testqueue testexchange testkey
ok.

Эта команда создала обменник типа direct testexchange и очередь с именем testqueue. Очередь привязана к обменнику с помощью ключа маршрутизации testkey.

Теперь все сообщения, отправленные в обменник testexchange с ключом маршрутизации testkey, будут перемещены в эту очередь. Отправить сообщение можно командой basic.publish:

4> basic.publish 'This is a message!' testexchange testkey
ok.

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

Извлечём сообщение из очереди:

5> basic.get testqueue
{'body': 'This is a message!',
 'delivery_info': {'delivery_tag': 1,
                   'exchange': u'testexchange',
                   'message_count': 0,
                   'redelivered': False,
                   'routing_key': u'testkey'},
 'properties': {}}

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

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

Подтвердить получение сообщения можно с помощью команды basic.ack:

6> basic.ack 1
ok.

После тестирования удалите созданные сущности:

7> queue.delete testqueue
ok. 0 messages deleted.
8> exchange.delete testexchange
ok.

В Celery доступные очереди определяются параметром task_queues.

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

default_exchange = Exchange('default', type='direct')
media_exchange = Exchange('media', type='direct')

app.conf.task_queues = (
    Queue('default', default_exchange, routing_key='default'),
    Queue('videos', media_exchange, routing_key='media.video'),
    Queue('images', media_exchange, routing_key='media.image')
)
app.conf.task_default_queue = 'default'
app.conf.task_default_exchange = 'default'
app.conf.task_default_routing_key = 'default'

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

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

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

from kombu import Exchange, Queue, binding

media_exchange = Exchange('media', type='direct')

CELERY_QUEUES = (
    Queue('media', [
        binding(media_exchange, routing_key='media.video'),
        binding(media_exchange, routing_key='media.image'),
    ]),
)

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

  1. Аргументы маршрутизации для Task.apply_async().

  2. Атрибуты маршрутизации, заданные непосредственно в Task.

  3. Маршрутизаторы, определённые в параметре task_routes.

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

Маршрутизатор — это функция, определяющая параметры маршрутизации задачи.

Чтобы определить новый маршрутизатор, достаточно создать функцию с сигнатурой (name, args, kwargs, options, task=None, **kw):

def route_task(name, args, kwargs, options, task=None, **kw):
        if name == 'myapp.tasks.compress_video':
            return {'exchange': 'video',
                    'exchange_type': 'topic',
                    'routing_key': 'video.compress'}

Если вернуть ключ queue, он будет дополнен параметрами соответствующей очереди, заданными в task_queues:

{'queue':'video','routing_key':'video.compress'}

преобразуется в →

{'queue':'video',
'exchange':'video',
'exchange_type':'topic',
'routing_key':'video.compress'}

Установите классы маршрутизаторов, добавив их в параметр task_routes:

task_routes = (route_task,)

Функции маршрутизаторов также можно указать по имени:

task_routes = ('myapp.routers.route_task',)

Для простых соответствий «имя задачи → маршрут», таких как в приведённом выше примере маршрутизатора, можно просто добавить словарь в task_routes, чтобы получить то же поведение:

task_routes = {
    'myapp.tasks.compress_video': {
        'queue': 'video',
        'routing_key': 'video.compress',
    },
}

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

Можно также определить несколько маршрутизаторов в последовательности:

task_routes = [
    route_task,
    {
        'myapp.tasks.compress_video': {
            'queue': 'video',
            'routing_key': 'video.compress',
    },
]

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

Если вы используете Redis или RabbitMQ, можно также указать в маршруте приоритет очереди по умолчанию.

task_routes = {
    'myapp.tasks.compress_video': {
        'queue': 'video',
        'routing_key': 'video.compress',
        'priority': 10,
    },
}

Аналогично, вызов apply_async для задачи переопределит этот приоритет по умолчанию.

task.apply_async(priority=0)

Порядок приоритетов и отзывчивость кластера

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

Обратите внимание: при использовании брокера Redis приоритеты сортируются в обратном порядке — 0 означает наивысший приоритет.

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

from kombu.common import Broadcast

app.conf.task_queues = (Broadcast('broadcast_tasks'),)
app.conf.task_routes = {
    'tasks.reload_cache': {
        'queue': 'broadcast_tasks',
        'exchange': 'broadcast_tasks'
    }
}

Теперь задача tasks.reload_cache будет отправлена каждому рабочему процессу, получающему задачи из этой очереди.

Вот ещё один пример широковещательной маршрутизации — на этот раз с расписанием celery beat:

from kombu.common import Broadcast
from celery.schedules import crontab

app.conf.task_queues = (Broadcast('broadcast_tasks'),)

app.conf.beat_schedule = {
    'test-task': {
        'task': 'tasks.reload_cache',
        'schedule': crontab(minute=0, hour='*/3'),
        'options': {'exchange': 'broadcast_tasks'}
    },
}

Широковещательная рассылка и результаты

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

В этом случае рекомендуется задать атрибут task.ignore_result.

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

Spec-Zone.ru

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