Spec-Zone.ru › Celery

celery.events

Получатель и диспетчер событий мониторинга.

События — это поток сообщений, отправляемых при определённых действиях, происходящих в рабочем процессе (и клиентах, если включён параметр task_send_sent_event); они используются для мониторинга.

celery.events.Event(type, _fields=None, __dict__=<class 'dict'>, __now__=<built-in function time>, **fields)

Создать событие.

Примечания

Событие — это просто словарь: единственное обязательное поле — type. Если поле timestamp не задано, ему будет присвоено текущее время.

classcelery.events.EventDispatcher(connection=None, hostname=None, enabled=True, channel=None, buffer_while_offline=True, app=None, serializer=None, groups=None, delivery_mode=1, buffer_group=None, buffer_limit=24, on_send_buffered=None)

Отправляет сообщения о событиях.

Параметры:
  • connection (kombu.Connection) – Подключение к брокеру.

  • hostname (str) – Имя узла, под которым идентифицируется объект; по умолчанию используется имя узла, возвращаемое функцией anon_nodename().

  • groups (Sequence[str]) – Список групп, для которых отправляются события. send() будет игнорировать запросы на отправку для групп, которых нет в этом списке. Если значение равно None, будут отправляться все события. Примеры групп: "task" и "worker".

  • enabled (bool) – Установите значение False, чтобы не публиковать события; в этом случае вызов send() ничего не будет делать.

  • channel (kombu.Channel) – Можно использовать вместо connection, чтобы указать конкретный канал для отправки событий.

  • buffer_while_offline (bool) – Если параметр включён, события будут буферизоваться при разрыве соединения. После восстановления соединения необходимо вызвать flush().

Примечание

После использования необходимо вызвать close().

DISABLED_TRANSPORTS={'sql'}
app=None
close()

Закрыть диспетчер событий.

disable()
enable()
extend_buffer(other)

Скопировать исходящий буфер другого экземпляра.

flush(errors=True, groups=True)

Очистить исходящий буфер.

on_disabled=None
on_enabled=None
publish(type, fields, producer, blind=False, Event=<function Event>, **kwargs)

Опубликовать событие с помощью пользовательского Producer.

Параметры:
  • type (str) – Имя типа события; группа отделяется дефисом (-). fields: словарь полей события, значения которого должны быть сериализуемы в JSON.

  • producer (kombu.Producer) – Используемый экземпляр производителя: будет вызван только метод publish.

  • retry (bool) – Повторять попытку при сбое соединения.

  • retry_policy (Mapping) – Сопоставление пользовательских параметров политики повторных попыток. См. ensure().

  • blind (bool) – Не задавать значение логических часов (а также не передавать внутренние логические часы).

  • Event (Callable) – Тип события, используемый для создания события. По умолчанию — Event().

  • utcoffset (Callable) – Функция, возвращающая текущее смещение UTC в часах.

propertypublisher
send(type, blind=False, utcoffset=<function utcoffset>, retry=False, retry_policy=None, Event=<function Event>, **fields)

Отправить событие.

Параметры:
  • type (str) – Имя типа события; группа отделяется дефисом (-).

  • retry (bool) – Повторять попытку при сбое соединения.

  • retry_policy (Mapping) – Сопоставление пользовательских параметров политики повторных попыток. См. ensure().

  • blind (bool) – Не задавать значение логических часов (а также не передавать внутренние логические часы).

  • Event (Callable) – Тип события, используемый для создания события; по умолчанию — Event().

  • utcoffset (Callable) – Функция, возвращающая текущее смещение UTC в часах.

  • **fields (Any) – Поля события; значения должны быть сериализуемы в JSON.

classcelery.events.EventReceiver(channel, handlers=None, routing_key='#', node_id=None, app=None, queue_prefix=None, accept=None, queue_ttl=None, queue_expires=None, queue_exclusive=None, queue_durable=None)

Получать события.

Параметры:
  • connection (kombu.Connection) – Подключение к брокеру.

  • handlers (Mapping[Callable]) – Обработчики событий. Это сопоставление имён типов событий с соответствующими обработчиками. Специальный обработчик “*” перехватывает все события, для которых не задан обработчик.

app=None
capture(limit=None, timeout=None, wakeup=True)

Запустить потребителя для перехвата событий.

Этот метод должен выполняться в главном процессе и никогда не остановится, если EventDispatcher.should_stop не установлено в True или выполнение не будет принудительно прервано с помощью KeyboardInterrupt или SystemExit.

propertyconnection
event_from_message(body, localize=True, now=<built-in function time>, tzfields=operator.itemgetter('utcoffset', 'timestamp'), adjust_timestamp=<function adjust_timestamp>, CLIENT_CLOCK_SKEW=-1)
get_consumers(Consumer, channel)
itercapture(limit=None, timeout=None, wakeup=True)
on_consume_ready(connection, channel, consumers, wakeup=True, **kwargs)
process(type, event)

Обработать событие, передав его настроенному обработчику.

wakeup_workers(channel=None)
celery.events.get_exchange(conn, name='celeryev')

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

Параметры:
  • conn (kombu.Connection) – Подключение, используемое для отправки и получения событий.

  • name (str) – Имя обменника. По умолчанию — celeryev.

Примечание

Тип события меняется при использовании Redis в качестве транспорта (с topic на fanout).

celery.events.group_from(type)

Получить часть имени типа события, обозначающую группу.

Пример

>>> group_from('task-sent')
'task'
>>> group_from('custom-my-event')
'custom'

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/reference/celery.events.html

Spec-Zone.ru

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