Spec-Zone.ru › Celery

Оптимизация

Введение

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

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

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

Обеспечение работоспособности

В книге «Programming Pearls» Джон Бентли знакомит читателя с концепцией приблизительных расчётов «на салфетке», задавая вопрос:

❝ Сколько воды вытекает из реки Миссисипи за день? ❞

Смысл этого упражнения [*] — показать, что существует предел объёма данных, который система может обработать за приемлемое время. Приблизительные расчёты «на салфетке» можно использовать, чтобы заранее спланировать работу с учётом этого предела.

В Celery: если выполнение задачи занимает 10 минут, а каждую минуту поступает 10 новых задач, очередь никогда не опустеет. Поэтому очень важно следить за длиной очередей!

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

Общие настройки

Пулы подключений к брокеру

Пул подключений к брокеру включён по умолчанию начиная с версии 2.5.

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

Использование временных очередей

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

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

from kombu import Exchange, Queue

task_queues = (
    Queue('celery', routing_key='celery'),
    Queue('transient', Exchange('transient', delivery_mode=1),
          routing_key='transient', durable=False),
)

или воспользоваться настройкой task_routes:

task_routes = {
    'proj.tasks.add': {'queue': 'celery', 'delivery_mode': 'transient'}
}

Настройка delivery_mode изменяет способ доставки сообщений в эту очередь. Значение 1 означает, что сообщение не будет записано на диск, а значение 2 (по умолчанию) — что сообщение может быть записано на диск.

Чтобы направить задачу в новую временную очередь, можно указать аргумент queue (или использовать настройку task_routes):

task.apply_async(args, queue='transient')

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

Настройки исполнителей

Ограничения предварительной выборки

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

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

Значение предварительной выборки исполнителя по умолчанию равно настройке worker_prefetch_multiplier, умноженной на количество слотов параллельного выполнения [‡] (процессов/потоков/зелёных потоков).

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

Однако если у вас много коротких задач и для вас важны пропускная способность и задержка обмена данными, это значение должно быть большим. Исполнитель может обрабатывать больше задач в секунду, если сообщения уже предварительно выбраны и находятся в памяти. Возможно, придётся поэкспериментировать, чтобы найти подходящее значение. В таких случаях могут подойти значения 50 или 150 — например, 64 или 128.

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

Резервирование одной задачи за раз

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

Обратите внимание: исключение в Celery считается штатной ситуацией и задача будет подтверждена. Подтверждения нужны для защиты от сбоев, которые нельзя обработать обычным механизмом исключений Python (например, отключения питания, повреждения памяти, аппаратного сбоя, фатального сигнала и т. д.). Для обработки обычных исключений используйте task.retry(), чтобы повторить задачу.

См. также

Примечания в разделе Использовать retry или acks_late?.

При использовании раннего подтверждения, включённого по умолчанию, множитель предварительной выборки, равный одному, означает, что исполнитель зарезервирует не более одной дополнительной задачи для каждого процесса исполнителя. Иными словами, если исполнитель запущен с параметром -c 10, он может в любой момент зарезервировать не более 20 задач (10 подтверждённых выполняемых задач и 10 неподтверждённых зарезервированных задач).

Пользователи часто спрашивают, можно ли отключить «предварительную выборку задач». Это возможно, но с некоторыми оговорками. Можно настроить исполнителя так, чтобы он резервировал только столько задач, сколько у него процессов, при условии, что подтверждение задач выполняется поздно (10 неподтверждённых выполняемых задач для -c 10).

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

Включить такое поведение можно с помощью следующих параметров конфигурации:

task_acks_late = True
worker_prefetch_multiplier = 1

Если для ваших задач невозможно использовать позднее подтверждение, можно отключить предварительную выборку брокером, включив настройку worker_disable_prefetch. При такой настройке исполнитель получает новую задачу, только когда освобождается слот для выполнения, благодаря чему задачи не ждут в очереди за долгими задачами на загруженных исполнителях. Это также можно настроить из командной строки с помощью параметра --disable-prefetch. В настоящее время эта возможность поддерживается только при использовании Redis в качестве брокера.

Использование памяти

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

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

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

В Celery есть два основных способа снизить использование памяти из-за «высокого водяного знака» и/или утечек памяти в дочерних процессах: настройки worker_max_tasks_per_child и worker_max_memory_per_child.

Не задавайте для этих параметров слишком низкие значения, иначе исполнители будут тратить большую часть времени на перезапуск дочерних процессов вместо обработки задач. Например, если задать для worker_max_tasks_per_child значение 1, а запуск дочернего процесса занимает 1 секунду, этот процесс сможет обработать не более 60 задач в минуту (если считать, что задача выполняется мгновенно). Аналогичная проблема может возникнуть, если ваши задачи всегда превышают значение worker_max_memory_per_child.

Сноски

[*]

Эту главу можно бесплатно прочитать здесь: Приблизительные расчёты «на салфетке». Эта книга — классический труд. Настоятельно рекомендуем её прочитать.

[†]

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

[‡]

Это настройка параллельного выполнения: worker_concurrency или параметр celery worker -c.

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

Spec-Zone.ru

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