Расширения и этапы запуска
Возможно, вы захотите встроить пользовательские потребители Kombu для ручной обработки сообщений.
Для этого существует специальный класс этапа запуска ConsumerStep, в котором нужно определить только метод get_consumers, возвращающий список объектов kombu.Consumer, которые следует запускать при каждом установлении соединения:
from celery import Celery
from celery import bootsteps
from kombu import Consumer, Exchange, Queue
my_queue = Queue('custom', Exchange('custom'), 'routing_key')
app = Celery(broker='amqp://')
class MyConsumerStep(bootsteps.ConsumerStep):
def get_consumers(self, channel):
return [Consumer(channel,
queues=[my_queue],
callbacks=[self.handle_message],
accept=['json'])]
def handle_message(self, body, message):
print('Received message: {0!r}'.format(body))
message.ack()
app.steps['consumer'].add(MyConsumerStep)
def send_me_a_message(who, producer=None):
with app.producer_or_acquire(producer) as producer:
producer.publish(
{'hello': who},
serializer='json',
exchange=my_queue.exchange,
routing_key='routing_key',
declare=[my_queue],
retry=True,
)
if __name__ == '__main__':
send_me_a_message('world!')
Примечание
У потребителей Kombu есть два разных механизма вызова обработчиков сообщений. Первый — аргумент callbacks, принимающий список обработчиков с сигнатурой (body, message); второй — аргумент on_message, принимающий один обработчик с сигнатурой (message,). Последний автоматически не декодирует и не десериализует полезную нагрузку.
def get_consumers(self, channel):
return [Consumer(channel, queues=[my_queue],
on_message=self.on_message)]
def on_message(self, message):
payload = message.decode()
print(
'Received message: {0!r} {props!r} rawlen={s}'.format(
payload, props=message.properties, s=len(message.body),
))
message.ack()
Этапы запуска — это способ добавления функциональности рабочим процессам. Этап запуска — пользовательский класс, определяющий хуки для выполнения пользовательских действий на разных этапах работы рабочего процесса. Каждый этап запуска относится к схеме, а в рабочем процессе сейчас определены две схемы: Worker и Consumer.
- Рисунок A: Этапы запуска в схемах Worker и Consumer. Запуск
-
снизу вверх: первым этапом в схеме рабочего процесса является Timer, а последним — запуск схемы Consumer, которая затем устанавливает соединение с брокером и начинает получать сообщения.
Первой запускается схема Worker; вместе с ней запускаются основные компоненты, такие как цикл событий, пул обработки и таймер для задач ETA и других событий, запланированных на определённое время.
После полного запуска рабочего процесса запускается схема Consumer, которая настраивает выполнение задач, подключается к брокеру и запускает потребителей сообщений.
WorkController — основная реализация рабочего процесса; она содержит несколько методов и атрибутов, которые можно использовать в своём этапе запуска.
- app
-
Текущий экземпляр приложения.
- hostname
-
Имя узла рабочего процесса (например, worker1@example.com).
- blueprint
-
Это
Blueprintрабочего процесса.
- hub
-
Объект цикла событий (
Hub). С его помощью можно зарегистрировать обработчики в цикле событий.Поддерживается только транспортами с включённым асинхронным вводом-выводом (amqp, redis); в этом случае необходимо установить атрибут worker.use_eventloop.
Чтобы использовать этот объект, этап запуска рабочего процесса должен требовать этап Hub:
class WorkerStep(bootsteps.StartStopStep): requires = {'celery.worker.components:Hub'}
- pool
-
Текущий пул процессов/eventlet/gevent/потоков. См.
celery.concurrency.base.BasePool.Чтобы использовать этот объект, этап запуска рабочего процесса должен требовать этап Pool:
class WorkerStep(bootsteps.StartStopStep): requires = {'celery.worker.components:Pool'}
- timer
-
Timer, используемый для планирования функций.Чтобы использовать этот объект, этап запуска рабочего процесса должен требовать этап Timer:
class WorkerStep(bootsteps.StartStopStep): requires = {'celery.worker.components:Timer'}
- statedb
-
Database <celery.worker.state.Persistent>`для сохранения состояния между перезапусками рабочего процесса.Определяется только при включённом аргументе
statedb.Чтобы использовать этот объект, этап запуска рабочего процесса должен требовать этап
Statedb:class WorkerStep(bootsteps.StartStopStep): requires = {'celery.worker.components:Statedb'}
- autoscaler
-
Autoscalerдля автоматического увеличения и уменьшения числа процессов в пуле.Определяется только при включённом аргументе
autoscale.Чтобы использовать этот объект, этап запуска рабочего процесса должен требовать этап Autoscaler:
class WorkerStep(bootsteps.StartStopStep): requires = ('celery.worker.autoscaler:Autoscaler',)
Пример этапа запуска Worker:
from celery import bootsteps
class ExampleWorkerStep(bootsteps.StartStopStep):
requires = {'celery.worker.components:Pool'}
def __init__(self, worker, **kwargs):
print('Called when the WorkController instance is constructed')
print('Arguments to WorkController: {0!r}'.format(kwargs))
def create(self, worker):
# this method can be used to delegate the action methods
# to another object that implements ``start`` and ``stop``.
return self
def start(self, worker):
print('Called when the worker is started.')
def stop(self, worker):
print('Called when the worker shuts down.')
def terminate(self, worker):
print('Called when the worker terminates')
Каждый метод получает текущий экземпляр WorkController в качестве первого аргумента.
В другом примере таймер используется для регулярного пробуждения:
from celery import bootsteps
class DeadlockDetection(bootsteps.StartStopStep):
requires = {'celery.worker.components:Timer'}
def __init__(self, worker, deadlock_timeout=3600):
self.timeout = deadlock_timeout
self.requests = []
self.tref = None
def start(self, worker):
# run every 30 seconds.
self.tref = worker.timer.call_repeatedly(
30.0, self.detect, (worker,), priority=10,
)
def stop(self, worker):
if self.tref:
self.tref.cancel()
self.tref = None
def detect(self, worker):
# update active requests
for req in worker.active_requests:
if req.time_start and time() - req.time_start > self.timeout:
raise SystemExit()
Рабочий процесс Celery отправляет сообщения в подсистему журналирования Python для различных событий на протяжении всего жизненного цикла задачи. Эти сообщения можно настроить, переопределив форматные строки LOG_<TYPE>, определённые в celery/app/trace.py. Например:
import celery.app.trace celery.app.trace.LOG_SUCCESS = "This is a custom message"
Для форматирования % всем форматным строкам передаются имя и идентификатор задачи; некоторые из них также получают дополнительные поля, например возвращаемое значение или исключение, из-за которого задача завершилась с ошибкой. Эти поля можно использовать в пользовательских форматных строках, например так:
import celery.app.trace celery.app.trace.LOG_REJECTED = "%(name)r is cursed and I won't run it: %(exc)s"
Схема Consumer устанавливает соединение с брокером и запускается заново при каждой потере этого соединения. Этапы запуска Consumer включают отправку сигналов проверки активности рабочего процесса, потребитель команд удалённого управления и, что особенно важно, потребитель задач.
Создавая этапы запуска потребителя, учитывайте, что должна быть возможность перезапустить вашу схему. Для этапов запуска потребителя определён дополнительный метод «shutdown», который вызывается при завершении работы рабочего процесса.
- app
-
Текущий экземпляр приложения.
- controller
-
Родительский объект
WorkController, создавший этот потребитель.
- hostname
-
Имя узла рабочего процесса (например, worker1@example.com).
- blueprint
-
Это
Blueprintрабочего процесса.
- hub
-
Объект цикла событий (
Hub). С его помощью можно зарегистрировать обработчики в цикле событий.Поддерживается только транспортами с включённым асинхронным вводом-выводом (amqp, redis); в этом случае необходимо установить атрибут worker.use_eventloop.
Чтобы использовать этот объект, этап запуска рабочего процесса должен требовать этап Hub:
class WorkerStep(bootsteps.StartStopStep): requires = {'celery.worker.components:Hub'}
- connection
-
Текущее соединение с брокером (
kombu.Connection).Чтобы использовать этот объект, этап запуска потребителя должен требовать этап «Connection»:
class Step(bootsteps.StartStopStep): requires = {'celery.worker.consumer.connection:Connection'}
- event_dispatcher
-
Объект
app.events.Dispatcher, который можно использовать для отправки событий.Чтобы использовать этот объект, этап запуска потребителя должен требовать этап Events.
class Step(bootsteps.StartStopStep): requires = {'celery.worker.consumer.events:Events'}
- gossip
-
Широковещательная связь между рабочими процессами (
Gossip).Чтобы использовать этот объект, этап запуска потребителя должен требовать этап Gossip.
class RatelimitStep(bootsteps.StartStopStep): """Rate limit tasks based on the number of workers in the cluster.""" requires = {'celery.worker.consumer.gossip:Gossip'} def start(self, c): self.c = c self.c.gossip.on.node_join.add(self.on_cluster_size_change) self.c.gossip.on.node_leave.add(self.on_cluster_size_change) self.c.gossip.on.node_lost.add(self.on_node_lost) self.tasks = [ self.app.tasks['proj.tasks.add'] self.app.tasks['proj.tasks.mul'] ] self.last_size = None def on_cluster_size_change(self, worker): cluster_size = len(list(self.c.gossip.state.alive_workers())) if cluster_size != self.last_size: for task in self.tasks: task.rate_limit = 1.0 / cluster_size self.c.reset_rate_limits() self.last_size = cluster_size def on_node_lost(self, worker): # may have processed heartbeat too late, so wake up soon # in order to see if the worker recovered. self.c.timer.call_after(10.0, self.on_cluster_size_change)Обработчики
-
<set> gossip.on.node_joinВызывается при каждом присоединении нового узла к кластеру и получает экземпляр
Worker. -
<set> gossip.on.node_leaveВызывается при каждом выходе нового узла из кластера (завершении работы) и получает экземпляр
Worker. -
<set> gossip.on.node_lostВызывается, если для рабочего процесса в кластере не получен сигнал проверки активности (сигнал не получен или не обработан вовремя), и получает экземпляр
Worker.Это не обязательно означает, что рабочий процесс действительно отключён, поэтому используйте механизм тайм-аута, если стандартного тайм-аута сигнала проверки активности недостаточно.
-
- pool
-
Текущий пул процессов/eventlet/gevent/потоков. См.
celery.concurrency.base.BasePool.
- timer
-
Timer <celery.utils.timer2.Schedule, используемый для планирования функций.
- heart
-
Отвечает за отправку сигналов проверки активности рабочего процесса (
Heart).Чтобы использовать этот объект, этап запуска потребителя должен требовать этап Heart:
class Step(bootsteps.StartStopStep): requires = {'celery.worker.consumer.heart:Heart'}
- task_consumer
-
Объект
kombu.Consumer, используемый для получения сообщений задач.Чтобы использовать этот объект, этап запуска потребителя должен требовать этап Tasks:
class Step(bootsteps.StartStopStep): requires = {'celery.worker.consumer.tasks:Tasks'}
- strategies
-
Для каждого зарегистрированного типа задачи в этом словаре есть запись, значение которой используется для выполнения входящего сообщения этого типа задачи (стратегия выполнения задачи). Этот словарь создаётся этапом Tasks при запуске потребителя:
for name, task in app.tasks.items(): strategies[name] = task.start_strategy(app, consumer) task.__trace__ = celery.app.trace.build_tracer( name, task, loader, hostname )Чтобы использовать этот словарь, этап запуска потребителя должен требовать этап Tasks:
class Step(bootsteps.StartStopStep): requires = {'celery.worker.consumer.tasks:Tasks'}
- task_buckets
-
defaultdict, используемый для поиска ограничения частоты выполнения задачи по её типу. Значениями в этом словаре могут быть None (если ограничение отсутствует) или экземплярTokenBucket, реализующийconsume(tokens)иexpected_time(tokens).TokenBucket реализует алгоритм корзины токенов, но можно использовать любой алгоритм, если он соответствует тому же интерфейсу и определяет два указанных выше метода.
- qos
-
Объект
QoSможно использовать для изменения текущего значения prefetch_count канала задач:# increment at next cycle consumer.qos.increment_eventually(1) # decrement at next cycle consumer.qos.decrement_eventually(1) consumer.qos.set(10)
- consumer.reset_rate_limits()
-
Обновляет словарь
task_bucketsдля всех зарегистрированных типов задач.
- consumer.bucket_for_task(type, Bucket=TokenBucket)
-
Создаёт корзину ограничения частоты для задачи, используя её атрибут
task.rate_limit.
- consumer.add_task_queue(name, exchange=None, exchange_type=None,
- routing_key=None, \*\*options):
-
Добавляет новую очередь для получения сообщений. Изменение сохраняется при перезапуске соединения.
- consumer.cancel_task_queue(name)
-
Прекращает получение сообщений из очереди с указанным именем. Изменение сохраняется при перезапуске соединения.
- apply_eta_task(request)
-
Планирует выполнение задачи ETA на основе атрибута
request.eta. (Request)
app.steps['worker'] и app.steps['consumer'] можно изменить, чтобы добавить новые этапы запуска:
>>> app = Celery()
>>> app.steps['worker'].add(MyWorkerStep) # < add class, don't instantiate
>>> app.steps['consumer'].add(MyConsumerStep)
>>> app.steps['consumer'].update([StepA, StepB])
>>> app.steps['consumer']
{step:proj.StepB{()}, step:proj.MyConsumerStep{()}, step:proj.StepA{()}
Порядок этапов здесь не имеет значения, поскольку он определяется результирующим графом зависимостей (Step.requires).
Чтобы показать, как устанавливать этапы запуска и как они работают, рассмотрим пример этапа, который выводит бесполезную отладочную информацию. Его можно добавить как этап запуска как рабочего процесса, так и потребителя:
from celery import Celery
from celery import bootsteps
class InfoStep(bootsteps.Step):
def __init__(self, parent, **kwargs):
# here we can prepare the Worker/Consumer object
# in any way we want, set attribute defaults, and so on.
print('{0!r} is in init'.format(parent))
def start(self, parent):
# our step is started together with all other Worker/Consumer
# bootsteps.
print('{0!r} is starting'.format(parent))
def stop(self, parent):
# the Consumer calls stop every time the consumer is
# restarted (i.e., connection is lost) and also at shutdown.
# The Worker will call stop at shutdown only.
print('{0!r} is stopping'.format(parent))
def shutdown(self, parent):
# shutdown is called by the Consumer at shutdown, it's not
# called by Worker.
print('{0!r} is shutting down'.format(parent))
app = Celery(broker='amqp://')
app.steps['worker'].add(InfoStep)
app.steps['consumer'].add(InfoStep)
При запуске рабочего процесса с установленным этим этапом появятся следующие записи журнала:
<Worker: w@example.com (initializing)> is in init
<Consumer: w@example.com (initializing)> is in init
[2013-05-29 16:18:20,544: WARNING/MainProcess]
<Worker: w@example.com (running)> is starting
[2013-05-29 16:18:21,577: WARNING/MainProcess]
<Consumer: w@example.com (running)> is starting
<Consumer: w@example.com (closing)> is stopping
<Worker: w@example.com (closing)> is stopping
<Consumer: w@example.com (terminating)> is shutting down
После инициализации рабочего процесса операторы print перенаправляются в подсистему журналирования, поэтому строки «is starting» снабжаются отметками времени. Можно заметить, что при завершении работы этого больше не происходит: методы stop и shutdown вызываются внутри обработчика сигнала, а использовать журналирование внутри такого обработчика небезопасно. Журналирование с помощью модуля logging в Python не является реентерабельным: это означает, что нельзя прервать функцию, а затем вызвать её снова. Важно, чтобы написанные вами методы stop и shutdown также были реентерабельными.
Запуск рабочего процесса с параметром --loglevel=debug позволит увидеть больше информации о процессе запуска:
[2013-05-29 16:18:20,509: DEBUG/MainProcess] | Worker: Preparing bootsteps.
[2013-05-29 16:18:20,511: DEBUG/MainProcess] | Worker: Building graph...
<celery.apps.worker.Worker object at 0x101ad8410> is in init
[2013-05-29 16:18:20,511: DEBUG/MainProcess] | Worker: New boot order:
{Hub, Pool, Timer, StateDB, Autoscaler, InfoStep, Beat, Consumer}
[2013-05-29 16:18:20,514: DEBUG/MainProcess] | Consumer: Preparing bootsteps.
[2013-05-29 16:18:20,514: DEBUG/MainProcess] | Consumer: Building graph...
<celery.worker.consumer.Consumer object at 0x101c2d8d0> is in init
[2013-05-29 16:18:20,515: DEBUG/MainProcess] | Consumer: New boot order:
{Connection, Mingle, Events, Gossip, InfoStep, Agent,
Heart, Control, Tasks, event loop}
[2013-05-29 16:18:20,522: DEBUG/MainProcess] | Worker: Starting Hub
[2013-05-29 16:18:20,522: DEBUG/MainProcess] ^-- substep ok
[2013-05-29 16:18:20,522: DEBUG/MainProcess] | Worker: Starting Pool
[2013-05-29 16:18:20,542: DEBUG/MainProcess] ^-- substep ok
[2013-05-29 16:18:20,543: DEBUG/MainProcess] | Worker: Starting InfoStep
[2013-05-29 16:18:20,544: WARNING/MainProcess]
<celery.apps.worker.Worker object at 0x101ad8410> is starting
[2013-05-29 16:18:20,544: DEBUG/MainProcess] ^-- substep ok
[2013-05-29 16:18:20,544: DEBUG/MainProcess] | Worker: Starting Consumer
[2013-05-29 16:18:20,544: DEBUG/MainProcess] | Consumer: Starting Connection
[2013-05-29 16:18:20,559: INFO/MainProcess] Connected to amqp://guest@127.0.0.1:5672//
[2013-05-29 16:18:20,560: DEBUG/MainProcess] ^-- substep ok
[2013-05-29 16:18:20,560: DEBUG/MainProcess] | Consumer: Starting Mingle
[2013-05-29 16:18:20,560: INFO/MainProcess] mingle: searching for neighbors
[2013-05-29 16:18:21,570: INFO/MainProcess] mingle: no one here
[2013-05-29 16:18:21,570: DEBUG/MainProcess] ^-- substep ok
[2013-05-29 16:18:21,571: DEBUG/MainProcess] | Consumer: Starting Events
[2013-05-29 16:18:21,572: DEBUG/MainProcess] ^-- substep ok
[2013-05-29 16:18:21,572: DEBUG/MainProcess] | Consumer: Starting Gossip
[2013-05-29 16:18:21,577: DEBUG/MainProcess] ^-- substep ok
[2013-05-29 16:18:21,577: DEBUG/MainProcess] | Consumer: Starting InfoStep
[2013-05-29 16:18:21,577: WARNING/MainProcess]
<celery.worker.consumer.Consumer object at 0x101c2d8d0> is starting
[2013-05-29 16:18:21,578: DEBUG/MainProcess] ^-- substep ok
[2013-05-29 16:18:21,578: DEBUG/MainProcess] | Consumer: Starting Heart
[2013-05-29 16:18:21,579: DEBUG/MainProcess] ^-- substep ok
[2013-05-29 16:18:21,579: DEBUG/MainProcess] | Consumer: Starting Control
[2013-05-29 16:18:21,583: DEBUG/MainProcess] ^-- substep ok
[2013-05-29 16:18:21,583: DEBUG/MainProcess] | Consumer: Starting Tasks
[2013-05-29 16:18:21,606: DEBUG/MainProcess] basic.qos: prefetch_count->80
[2013-05-29 16:18:21,606: DEBUG/MainProcess] ^-- substep ok
[2013-05-29 16:18:21,606: DEBUG/MainProcess] | Consumer: Starting event loop
[2013-05-29 16:18:21,608: WARNING/MainProcess] celery@example.com ready.
Параметры отдельных команд
Дополнительные параметры командной строки можно добавить к командам worker, beat и events, изменив атрибут user_options экземпляра приложения.
Команды Celery используют модуль click для разбора аргументов командной строки. Поэтому для добавления пользовательских аргументов нужно добавить экземпляры click.Option в соответствующий набор.
Пример добавления пользовательского параметра к команде celery worker:
from celery import Celery
from click import Option
app = Celery(broker='amqp://')
app.user_options['worker'].add(Option(('--enable-my-option',),
is_flag=True,
help='Enable custom option.'))
Теперь все этапы запуска будут получать этот аргумент как именованный аргумент метода Bootstep.__init__:
from celery import bootsteps
class MyBootstep(bootsteps.Step):
def __init__(self, parent, enable_my_option=False, **options):
super().__init__(parent, **options)
if enable_my_option:
party()
app.steps['worker'].add(MyBootstep)
Параметры предварительной загрузки
Команда-оболочка celery поддерживает понятие «параметров предварительной загрузки». Это специальные параметры, передаваемые всем подкомандам.
Можно добавить новые параметры предварительной загрузки, например для указания шаблона конфигурации:
from celery import Celery
from celery import signals
from click import Option
app = Celery()
app.user_options['preload'].add(Option(('-Z', '--template'),
default='default',
help='Configuration template to use.'))
@signals.user_preload_options.connect
def on_preload_parsed(options, **kwargs):
use_template(options['template'])
В команду-оболочку celery можно добавлять новые команды с помощью точек входа setuptools.
Точки входа — это специальные метаданные, которые можно добавить в пакеты для регистрации программы setup.py, а после установки считать из системы с помощью модуля importlib.
Для установки дополнительных подкоманд Celery распознаёт точки входа celery.commands; значение точки входа должно указывать на допустимую команду Click.
Так расширение для мониторинга https://pypi.org/project/Flower/ может добавить команду celery flower: для этого в setup.py добавляется точка входа:
setup(
name='flower',
entry_points={
'celery.commands': [
'flower = flower.command:flower',
],
}
)
Определение команды состоит из двух частей, разделённых знаком равенства: первая часть — имя подкоманды (flower), вторая — полный путь к символу функции, реализующей команду:
flower.command:flower
Как показано выше, путь к модулю и имя атрибута следует разделять двоеточием.
В модуле flower/command.py функцию команды можно определить следующим образом:
import click
@click.command()
@click.option('--port', default=8888, type=int, help='Webserver port')
@click.option('--debug', is_flag=True)
def flower(port, debug):
print('Running our command')
Hub — асинхронный цикл событий рабочего процесса
- поддерживаемые транспорты:
-
amqp, redis
Добавлено в версии 3.0.
При использовании транспортов брокера amqp или redis рабочий процесс применяет асинхронный ввод-вывод. В конечном счёте планируется, что все транспорты будут использовать цикл событий, но на это потребуется время, поэтому в остальных транспортах пока используется решение на основе потоков.
- hub.add(fd, callback, flags)
- hub.add_reader(fd, callback, \*args)
-
Добавляет обработчик, который будет вызван, когда
fdстанет доступен для чтения.Обработчик останется зарегистрированным, пока его явно не удалят с помощью
hub.remove(fd)или пока файловый дескриптор не будет автоматически отброшен из-за того, что он больше недействителен.Обратите внимание: для одного файлового дескриптора одновременно можно зарегистрировать только один обработчик. Поэтому повторный вызов
addудалит обработчик, ранее зарегистрированный для этого дескриптора.Файловый дескриптор — это любой файловый объект, поддерживающий метод
fileno, либо номер файлового дескриптора (int).
- hub.add_writer(fd, callback, \*args)
-
Добавляет обработчик, который будет вызван, когда
fdстанет доступен для записи. См. также примечания выше дляhub.add_reader().
- hub.remove(fd)
-
Удаляет из цикла все обработчики для файлового дескриптора
fd.
- timer.call_after(secs, callback, args=(), kwargs=(),
- priority=0)
- timer.call_repeatedly(secs, callback, args=(), kwargs=(),
- priority=0)
- timer.call_at(eta, callback, args=(), kwargs=(),
- priority=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/extending.html