Spec-Zone.ru › Celery

celery.worker.consumer.consumer

Схема потребителя рабочего процесса.

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

classcelery.worker.consumer.consumer.Consumer(on_task_request, init_callback=<function noop>, hostname=None, pool=None, app=None, timer=None, controller=None, hub=None, amqheartbeat=None, worker_options=None, disable_rate_limits=False, initial_prefetch_count=2, prefetch_multiplier=1, **kwargs)

Схема потребителя.

classBlueprint(steps=None, name=None, on_start=None, on_close=None, on_stopped=None)

Схема потребителя.

default_steps=['celery.worker.consumer.connection:Connection', 'celery.worker.consumer.mingle:Mingle', 'celery.worker.consumer.events:Events', 'celery.worker.consumer.gossip:Gossip', 'celery.worker.consumer.heart:Heart', 'celery.worker.consumer.control:Control', 'celery.worker.consumer.tasks:Tasks', 'celery.worker.consumer.delayed_delivery:DelayedDelivery', 'celery.worker.consumer.consumer:Evloop', 'celery.worker.consumer.agent:Agent']
name='Consumer'
shutdown(parent)
Strategies

псевдоним для dict

add_task_queue(queue, exchange=None, exchange_type=None, routing_key=None, **options)
apply_eta_task(task)

Метод, вызываемый таймером для выполнения задачи с ETA/обратным отсчётом.

broker_connection_retry_attempt=0

Счётчик попыток подключения к брокеру. Сбрасывается в 0 после успешного подключения.

bucket_for_task(type)
call_soon(p, *args, **kwargs)
cancel_active_requests()

Отменяет активные запросы при завершении работы.

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

Не отменяет успешно выполненные задачи, даже если они ещё не были подтверждены.

cancel_task_queue(queue)
connect()

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

Повторяет попытки установить соединение, если включён параметр broker_connection_retry.

connection_for_read(heartbeat=None)
connection_for_write(url=None, heartbeat=None)
create_task_handler(promise=<class 'vine.promises.promise'>)
ensure_connected(conn)
first_connection_attempt=True

Этот флаг отключается после первой неудачной попытки подключения.

init_callback=None

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

loop_args()
propertymax_prefetch_count
on_close()
on_connection_error_after_connected(exc)
on_connection_error_before_connected(exc)
on_decode_error(message, exc)

Функция обратного вызова, вызываемая при ошибке декодирования сообщения.

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

Параметры:
  • message (kombu.Message) – Полученное сообщение.

  • exc (Exception) – Обрабатываемое исключение.

on_invalid_task(body, message, exc)
on_ready()
on_send_event_buffered()
on_unknown_message(body, message)
on_unknown_task(body, message, exc)
perform_pending_operations()
pool=None

Текущий экземпляр пула рабочего процесса.

register_with_event_loop(hub)
reset_rate_limits()
restart_count=-1
shutdown()
start()
stop()
timer=None

Таймер для внутренних задач с высоким приоритетом, например отправки сигналов сердцебиения.

update_strategies()
classcelery.worker.consumer.consumer.Evloop(parent, **kwargs)

Служба цикла событий.

Примечание

Этот компонент всегда запускается последним.

label='event loop'
last=True
name='celery.worker.consumer.consumer.Evloop'
patch_all(c)
start(c)
celery.worker.consumer.consumer.dump_body(m, body)

Форматирует тело сообщения для отладки.

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.worker.consumer.consumer.html

Spec-Zone.ru

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