Сигналы
Сигналы позволяют слабо связанные приложениям получать уведомления о действиях, происходящих в других частях приложения.
Celery включает множество сигналов, к которым приложение может подключаться для изменения поведения определённых действий.
Сигналы вызываются событиями нескольких типов. Вы можете подключаться к этим сигналам, чтобы выполнять действия при их вызове.
Пример подключения к сигналу after_task_publish:
from celery.signals import after_task_publish
@after_task_publish.connect
def task_sent_handler(sender=None, headers=None, body=None, **kwargs):
# information about task are located in headers for task messages
# using the task protocol version 2.
info = headers if 'task' in headers else body
print('after_task_publish for task id {info[id]}'.format(
info=info,
))
У некоторых сигналов также есть отправитель, по которому можно фильтровать. Например, сигнал after_task_publish использует имя задачи в качестве отправителя, поэтому, передав аргумент sender в connect, можно подключить обработчик, который будет вызываться каждый раз при публикации задачи с именем «proj.tasks.add»:
@after_task_publish.connect(sender='proj.tasks.add')
def task_sent_handler(sender=None, headers=None, body=None, **kwargs):
# information about task are located in headers for task messages
# using the task protocol version 2.
info = headers if 'task' in headers else body
print('after_task_publish for task id {info[id]}'.format(
info=info,
))
В сигналах используется та же реализация, что и в django.core.dispatch. Поэтому по умолчанию все обработчики сигналов получают и другие именованные параметры (например, signal).
Рекомендуется, чтобы обработчики сигналов принимали произвольные именованные аргументы (то есть **kwargs). Благодаря этому новые версии Celery смогут добавлять дополнительные аргументы, не нарушая работу пользовательского кода.
Добавлено в версии 3.1.
Вызывается перед публикацией задачи. Обратите внимание, что сигнал выполняется в процессе, отправляющем задачу.
Отправитель — имя отправляемой задачи.
Передаваемые аргументы:
-
body -
exchangeИмя биржи, в которую нужно отправить сообщение, или объект
Exchange. -
routing_keyКлюч маршрутизации, используемый при отправке сообщения.
-
headersОтображение заголовков приложения (можно изменить).
-
propertiesСвойства сообщения (можно изменить).
-
declare -
retry_policyОтображение параметров повтора. Может содержать любые аргументы
kombu.Connection.ensure()и может быть изменено.
Вызывается после отправки задачи брокеру. Обратите внимание, что сигнал выполняется в процессе, отправившем задачу.
Отправитель — имя отправляемой задачи.
Передаваемые аргументы:
-
headers -
body -
exchangeИмя биржи или использованный объект
Exchange. -
routing_keyИспользованный ключ маршрутизации.
Вызывается перед выполнением задачи.
Отправитель — выполняемый объект задачи.
Передаваемые аргументы:
-
task_idИдентификатор выполняемой задачи.
-
taskВыполняемая задача.
-
argsПозиционные аргументы задачи.
-
kwargsИменованные аргументы задачи.
Вызывается после выполнения задачи.
Отправитель — выполненный объект задачи.
Передаваемые аргументы:
-
task_idИдентификатор выполняемой задачи.
-
taskВыполняемая задача.
-
argsПозиционные аргументы задачи.
-
kwargsИменованные аргументы задачи.
-
retvalВозвращаемое задачей значение.
-
stateИмя результирующего состояния.
Вызывается, когда задача будет запущена повторно.
Отправитель — объект задачи.
Передаваемые аргументы:
-
requestТекущий запрос задачи.
-
reasonПричина повтора (обычно экземпляр исключения, но всегда может быть преобразована в
str). -
einfoПодробная информация об исключении, включая трассировку стека (объект
billiard.einfo.ExceptionInfo).
Примечание
Гарантируется, что во всех случаях будет передан только аргумент request. Аргументы reason и einfo могут быть None или не передаваться в некоторых ситуациях, например, когда задача отменяется и запускается повторно. Обработчики сигналов не должны считать, что эти аргументы присутствуют всегда.
Вызывается при успешном выполнении задачи.
Отправитель — выполненный объект задачи.
Передаваемые аргументы
-
result-
Возвращаемое задачей значение.
Вызывается при сбое задачи.
Отправитель — выполненный объект задачи.
Передаваемые аргументы:
-
task_idИдентификатор задачи.
-
exceptionВозникший экземпляр исключения.
-
argsПозиционные аргументы, с которыми была вызвана задача.
-
kwargsИменованные аргументы, с которыми была вызвана задача.
-
tracebackОбъект трассировки стека.
-
einfoЭкземпляр
billiard.einfo.ExceptionInfo.
Вызывается при возникновении внутренней ошибки Celery во время выполнения задачи.
Отправитель — выполненный объект задачи.
Передаваемые аргументы:
-
task_idИдентификатор задачи.
-
argsПозиционные аргументы, с которыми была вызвана задача.
-
kwargsИменованные аргументы, с которыми была вызвана задача.
-
requestИсходный словарь запроса. Он передаётся, поскольку к моменту возникновения исключения объект
task.requestможет быть ещё не готов. -
exceptionВозникший экземпляр исключения.
-
tracebackОбъект трассировки стека.
-
einfoЭкземпляр
billiard.einfo.ExceptionInfo.
Вызывается, когда задача получена от брокера и готова к выполнению.
Отправитель — объект потребителя.
Передаваемые аргументы:
-
requestЭто экземпляр
Request, а неtask.request. При использовании пула prefork этот сигнал вызывается в родительском процессе, поэтомуtask.requestнедоступен и не должен использоваться. Используйте вместо него этот объект: у них много общих полей.
Вызывается, когда задача отзывается или завершается рабочим процессом.
Отправитель — отозванный или завершённый объект задачи.
Передаваемые аргументы:
-
requestЭто экземпляр
Context, а неtask.request. При использовании пула prefork этот сигнал вызывается в родительском процессе, поэтомуtask.requestнедоступен и не должен использоваться. Используйте вместо него этот объект: у них много общих полей. -
terminatedУстановлено в значение
True, если задача была завершена. -
signumНомер сигнала, использованного для завершения задачи. Если это значение равно
None, а terminated равноTrue, следует считать, что использовано значениеTERM. -
expiredУстановлено в значение
True, если срок действия задачи истёк.
Вызывается, когда рабочий процесс получает сообщение для задачи, не зарегистрированной в системе.
Отправитель — рабочий процесс Consumer.
Передаваемые аргументы:
-
nameИмя задачи, не найденной в реестре.
-
idИдентификатор задачи, указанный в сообщении.
-
messageОбъект исходного сообщения.
-
excВозникшая ошибка.
Вызывается, когда рабочий процесс получает сообщение неизвестного типа в одну из своих очередей задач.
Отправитель — рабочий процесс Consumer.
Передаваемые аргументы:
-
messageОбъект исходного сообщения.
-
excВозникшая ошибка (если есть).
Этот сигнал отправляется, когда программа (рабочий процесс, beat, оболочка и т. д.) запрашивает импорт модулей, указанных в настройках include и imports.
Отправитель — экземпляр приложения.
Этот сигнал отправляется после настройки экземпляра рабочего процесса, но до вызова им метода run. Это означает, что уже включены все очереди, указанные параметром celery worker -Q, настроено ведение журналов и т. д.
Сигнал можно использовать для добавления пользовательских очередей, из которых нужно получать сообщения всегда, независимо от параметра celery worker -Q. В следующем примере для каждого рабочего процесса создаётся отдельная прямая очередь. Затем эти очереди можно использовать для маршрутизации задачи к конкретному рабочему процессу:
from celery.signals import celeryd_after_setup
@celeryd_after_setup.connect
def setup_direct_queue(sender, instance, **kwargs):
queue_name = '{0}.dq'.format(sender) # sender is the nodename of the worker
instance.app.amqp.queues.select_add(queue_name)
Передаваемые аргументы:
-
senderИмя узла рабочего процесса.
-
instanceИнициализируемый экземпляр
celery.apps.worker.Worker. Обратите внимание, что на данный момент заданы только атрибутыappиhostname(nodename), а выполнение остальной части__init__ещё не началось. -
confКонфигурация текущего приложения.
Это первый сигнал, отправляемый при запуске celery worker. Значение sender — имя хоста рабочего процесса, поэтому этот сигнал можно использовать для настройки отдельных рабочих процессов:
from celery.signals import celeryd_init
@celeryd_init.connect(sender='worker12@example.com')
def configure_worker12(conf=None, **kwargs):
conf.task_default_rate_limit = '10/m'
Чтобы настроить несколько рабочих процессов, можно не указывать отправителя при подключении:
from celery.signals import celeryd_init
@celeryd_init.connect
def configure_workers(sender=None, conf=None, **kwargs):
if sender in ('worker1@example.com', 'worker2@example.com'):
conf.task_default_rate_limit = '10/m'
if sender == 'worker3@example.com':
conf.worker_prefetch_multiplier = 0
Передаваемые аргументы:
-
senderИмя узла рабочего процесса.
-
instanceИнициализируемый экземпляр
celery.apps.worker.Worker. Обратите внимание, что на данный момент заданы только атрибутыappиhostname(nodename), а выполнение остальной части__init__ещё не началось. -
confКонфигурация текущего приложения.
-
optionsПараметры, переданные рабочему процессу через аргументы командной строки (включая значения по умолчанию).
Вызывается перед запуском рабочего процесса.
Вызывается в родительском процессе непосредственно перед созданием нового дочернего процесса в пуле prefork. Сигнал можно использовать для очистки объектов, которые некорректно работают при ветвлении процесса.
@signals.worker_before_create_process.connect
def clean_channels(**kwargs):
grpc_singleton.clean_channel()
Вызывается, когда рабочий процесс готов принимать задания.
Вызывается при отправке Celery сигнала проверки активности рабочего процесса.
Отправитель — экземпляр celery.worker.heartbeat.Heart.
Вызывается при начале завершения работы рабочего процесса.
Передаваемые аргументы:
-
sigПолученный сигнал POSIX.
-
howМетод завершения работы: штатный или принудительный.
-
exitcodeКод завершения, используемый при выходе основного процесса.
Вызывается при запуске всех дочерних процессов пула.
Обратите внимание: обработчики этого сигнала не должны блокировать выполнение более чем на 4 секунды, иначе процесс будет завершён, поскольку система сочтёт, что его запуск завершился с ошибкой.
Вызывается непосредственно перед завершением всех дочерних процессов пула.
Примечание. Нет гарантии, что этот сигнал будет вызван. Как и в случае с блоками finally, невозможно гарантировать вызов обработчиков при завершении работы; кроме того, их выполнение может быть прервано.
Передаваемые аргументы:
-
pidИдентификатор процесса завершающегося дочернего процесса.
-
exitcodeКод завершения, используемый при выходе дочернего процесса.
Вызывается непосредственно перед завершением работы рабочего процесса.
Вызывается при запуске celery beat (как отдельно, так и во встроенном режиме).
Отправитель — экземпляр celery.beat.Service.
Вызывается дополнительно к сигналу beat_init, когда celery beat запускается как встроенный процесс.
Отправитель — экземпляр celery.beat.Service.
Отправляется после запуска пула eventlet.
Отправитель — экземпляр celery.concurrency.eventlet.TaskPool.
Отправляется при завершении работы рабочего процесса, непосредственно перед тем, как пул eventlet запрашивают дождаться завершения оставшихся рабочих процессов.
Отправитель — экземпляр celery.concurrency.eventlet.TaskPool.
Отправляется после объединения пула, когда рабочий процесс готов к завершению.
Отправитель — экземпляр celery.concurrency.eventlet.TaskPool.
Отправляется каждый раз при передаче задачи пулу.
Отправитель — экземпляр celery.concurrency.eventlet.TaskPool.
Передаваемые аргументы:
-
targetЦелевая функция.
-
argsПозиционные аргументы.
-
kwargsИменованные аргументы.
Celery не будет настраивать регистраторы, если подключён этот сигнал. Поэтому его можно использовать для полной замены конфигурации журналирования на собственную.
Чтобы дополнить конфигурацию журналирования, настроенную Celery, используйте сигналы after_setup_logger и after_setup_task_logger.
Передаваемые аргументы:
-
loglevelУровень объекта журналирования.
-
logfileИмя файла журнала.
-
formatСтрока формата журнала.
-
colorizeУказывает, должны ли сообщения журнала быть цветными.
Отправляется после настройки каждого глобального регистратора (но не регистраторов задач). Используется для дополнения конфигурации журналирования.
Передаваемые аргументы:
-
loggerОбъект регистратора.
-
loglevelУровень объекта журналирования.
-
logfileИмя файла журнала.
-
formatСтрока формата журнала.
-
colorizeУказывает, должны ли сообщения журнала быть цветными.
Отправляется после настройки каждого регистратора задач. Используется для дополнения конфигурации журналирования.
Передаваемые аргументы:
-
loggerОбъект регистратора.
-
loglevelУровень объекта журналирования.
-
logfileИмя файла журнала.
-
formatСтрока формата журнала.
-
colorizeУказывает, должны ли сообщения журнала быть цветными.
Этот сигнал отправляется после того, как любая из программ командной строки Celery завершит разбор пользовательских параметров предварительной загрузки.
Его можно использовать для добавления дополнительных аргументов командной строки к общей команде celery:
from celery import Celery, signals
from click import Option
app = Celery()
# Celery 5.0+ uses click for its command-line interface.
# Use click.option to add new command-line arguments.
app.user_options['preload'].add(Option(
('--monitoring',), is_flag=True,
help='Enable our external monitoring utility, blahblah',
))
@signals.user_preload_options.connect
def handle_preload_options(options, **kwargs):
if options.get('monitoring'):
enable_monitoring()
Отправитель — экземпляр Command; его значение зависит от вызванной программы (например, для общей команды это будет объект CeleryCommand).
Передаваемые аргументы:
-
appЭкземпляр приложения.
-
optionsОтображение разобранных пользовательских параметров предварительной загрузки (со значениями по умолчанию).
Этот сигнал устарел; используйте вместо него after_task_publish.
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/signals.html