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(enumWebhookType) и тело 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).