Spec-Zone.ru › Celery

celery.worker.consumer

Потребитель рабочего процесса.

classcelery.worker.consumer.Agent(c, **kwargs)

Агент запускает акторы https://pypi.org/project/cell/.

conditional=True
create(c)

Создать шаг.

name='celery.worker.consumer.agent.Agent'
requires=(step:celery.worker.consumer.connection.Connection{()},)
classcelery.worker.consumer.Connection(c, **kwargs)

Служба управления подключением потребителя к брокеру.

info(c)
name='celery.worker.consumer.connection.Connection'
shutdown(c)
start(c)
classcelery.worker.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.Control(c, **kwargs)

Служба команд удалённого управления.

include_if(c)

Возвращает true, если загрузочный шаг следует включить.

Можно определить это как необязательное условие, которое решает, следует ли создавать этот шаг.

name='celery.worker.consumer.control.Control'
requires=(step:celery.worker.consumer.tasks.Tasks{(step:celery.worker.consumer.mingle.Mingle{(step:celery.worker.consumer.events.Events{(step:celery.worker.consumer.connection.Connection{()},)},)},)},)
classcelery.worker.consumer.Events(c, task_events=True, without_heartbeat=False, without_gossip=False, **kwargs)

Служба отправки событий мониторинга.

name='celery.worker.consumer.events.Events'
requires=(step:celery.worker.consumer.connection.Connection{()},)
shutdown(c)
start(c)
stop(c)
classcelery.worker.consumer.Gossip(c, without_gossip=False, interval=5.0, heartbeat_interval=2.0, **kwargs)

Шаг запуска, получающий события от других воркеров.

Он поддерживает актуальное значение логических часов.

call_task(task)
compatible_transport(app)
compatible_transports={'amqp', 'redis'}
election(id, topic, action=None)
get_consumers(channel)
label='Gossip'
name='celery.worker.consumer.gossip.Gossip'
on_elect(event)
on_elect_ack(event)
on_message(prepare, message)
on_node_join(worker)
on_node_leave(worker)
on_node_lost(worker)
periodic()
register_timer()
requires=(step:celery.worker.consumer.mingle.Mingle{(step:celery.worker.consumer.events.Events{(step:celery.worker.consumer.connection.Connection{()},)},)},)
start(c)
classcelery.worker.consumer.Heart(c, without_heartbeat=False, heartbeat_interval=None, **kwargs)

Шаг запуска, отправляющий сигналы активности.

Эта служба отправляет сообщение worker-heartbeat каждые n секунд.

Примечание

Не путать с сигналами активности на уровне протокола AMQP.

name='celery.worker.consumer.heart.Heart'
requires=(step:celery.worker.consumer.events.Events{(step:celery.worker.consumer.connection.Connection{()},)},)
shutdown(c)
start(c)
stop(c)
classcelery.worker.consumer.Mingle(c, without_mingle=False, **kwargs)

Шаг запуска, синхронизирующий состояние с соседними воркерами.

При запуске или перезапуске потребителя выполняются следующие действия:

  • Синхронизация логических часов.

  • Синхронизация отозванных задач.

compatible_transport(app)
compatible_transports={'amqp', 'gcpubsub', 'redis'}
label='Mingle'
name='celery.worker.consumer.mingle.Mingle'
on_clock_event(c, clock)
on_node_reply(c, nodename, reply)
on_revoked_received(c, revoked)
requires=(step:celery.worker.consumer.events.Events{(step:celery.worker.consumer.connection.Connection{()},)},)
send_hello(c)
start(c)
sync(c)
sync_with_node(c, clock=None, revoked=None, **kwargs)
classcelery.worker.consumer.Tasks(c, **kwargs)

Шаг запуска, запускающий потребитель сообщений задач.

info(c)

Возвращает сведения о потребителе задач.

name='celery.worker.consumer.tasks.Tasks'
qos_global(c) → bool

Определяет, следует ли применять глобальный QoS.

Дополнительная информация:

https://www.rabbitmq.com/docs/consumer-prefetch https://www.rabbitmq.com/docs/quorum-queues#global-qos

requires=(step:celery.worker.consumer.mingle.Mingle{(step:celery.worker.consumer.events.Events{(step:celery.worker.consumer.connection.Connection{()},)},)},)
shutdown(c)

Останавливает потребителя задач.

start(c)

Запускает потребителя задач.

stop(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/reference/celery.worker.consumer.html

Spec-Zone.ru

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