Skip to content

Celery и задачи

Конфигурация

Файл: src/utils/celery.py

app = Celery("tpm", broker=redis_url)
app.conf.result_backend = redis_url
app.conf.update(result_extended=True)

Celery задачи автообнаруживаются (autodiscover_tasks) из модулей: - src.moysklad.entities - src.services.expenseitems - src.services.set_prefix - src.utils.mail

Декоратор async_task

В celery.py определён декоратор async_task, который оборачивает async функции в Celery-задачи через asgiref.sync.AsyncToSync. Используется для задач set_prefix.tasks и др.

Центральная задача: save_webhook

@app.task(autoretry_for=(Exception,), retry_backoff=30)
def save_webhook(wh_type, body):
    ...
  • Принимает wh_type (enum WebhookType) и тело webhook.
  • URL: берётся из переменной окружения INTERNAL_BACKEND_URL (default: http://backend0:8081). В production указывает на внутренний endpoint Traefik (http://rep2-traefik:8080), который балансирует backend0 и backend1.
  • Transport retry: при ошибках 502/503/504 retry на уровне requests (3 попытки, backoff 1s) — защита от кратковременных сбоев.
  • Celery retry: при общем сбое — экспоненциальный backoff (base 30s).
  • Timeout: 60s (было 600s — блокировало воркер на 10 минут при отказе backend).

Эта задача — связующее звено между webhook-processor (быстрый ответ) и internal-webhook (бизнес-логика).

Периодические задачи (beat_schedule)

В настоящий момент beat_schedule пуст (все периодические задачи replace-buyer удалены вместе с приложением).

Worker init

@worker_init.connect() — инициализирует logfire для каждого воркера:

logfire.configure(service_name="worker", distributed_tracing=True)
logfire.instrument_celery()
logfire.instrument_requests(excluded_urls='moysklad.ru')

Flower

Мониторинг задач доступен по поддомену FLOWER_SUBDOMAIN_NAME (контейнер flower, порт 5555).