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