# Горизонтальное масштабирование очередей Celery Описывает, как развести обработку задач по отдельным репликам/сервисам под нагрузкой, когда одного контейнера `worker` недостаточно. Здесь — инструкция для оператора, применимая к текущей и будущей конфигурации `deploy/docker-compose.yml`. ## 1. Карта очередей Маршрутизация задач по очередям задана в `workers/celery_app.py` (`app.conf.task_routes`) — приоритет между семьями задач реализован **изоляцией очередей**, а не Redis-priorities Celery (транспорт `redis` эмулирует приоритеты ненадёжно, без строгих гарантий порядка — в отличие от RabbitMQ). | Очередь | Задачи | Кто слушает по умолчанию | |----------------|-------------------------------------------------------------------------|---------------------------| | `transcription`| `run_pipeline` (faster-whisper) | `worker-transcriber` (профиль `transcribe`), `--pool=solo` — ctranslate2/faster-whisper несовместимы с prefork | | `summarize` | `summarize_session` (Qwen map-reduce) | базовый `worker` | | `notify` | `notify_session`, `send_invitations` (.ics-приглашения) | базовый `worker` | | `celery` (default) | `cleanup_conferences`, `recover_stuck_summaries`, `recover_stuck_notifications` (маршрута не имеют, обслуживание) | базовый `worker` | Базовый сервис `worker` слушает `celery,summarize,notify` — для малых пресетов поставки (1–3) это один контейнер, обрабатывающий и суммаризацию, и уведомления, и обслуживание. `worker-transcriber` — всегда отдельный процесс (независимо от пресета), т.к. пул `solo` несовместим с остальными задачами в том же процессе. Глубину каждой очереди в реальном времени видно в `GET /metrics` (`vidconf_celery_queue_depth{queue=...}`, `backend/api/metrics.py`) и в Grafana-дашборде «Пайплайны пост-обработки» (алерт `QueueGrowing`, `deploy/monitoring/alerts.yml`) — по этим показателям и принимается решение о вынесении очереди в отдельную реплику. ## 2. Общее правило: `celery beat` — только один экземпляр Периодические задачи (`app.conf.beat_schedule`) планирует `celery beat`. Базовый `worker` запускается с флагом `-B` (worker + встроенный beat в одном процессе, `deploy/docker-compose.yml`). При масштабировании **нельзя** просто поднять несколько реплик сервиса с `-B` — каждая реплика завела бы собственный планировщик, и periodic-задачи (`cleanup_conferences`, `recover_stuck_*`) ставились бы в очередь многократно на каждом тике. Правило: ровно один процесс во всей инсталляции запускается с `-B` (встроенным или отдельным `celery -A workers.celery_app beat`); все дополнительные реплики — только `worker` без `-B`, с явным `-Q` на нужные очереди. ## 3. Вынос очереди в отдельную реплику Пример: суммаризация (`summarize`) стала узким местом — очередь растёт, базовый `worker` не успевает. Решение — отдельный сервис-потребитель только этой очереди, без встроенного beat (он остаётся на базовом `worker`). Добавить в `deploy/docker-compose.override.yml` (или новый профиль по аналогии с `worker-transcriber`): ```yaml services: worker-summarize: build: context: ../backend dockerfile: Dockerfile restart: unless-stopped # Без -B: beat уже запущен на базовом worker (см. правило выше). command: ["uv", "run", "celery", "-A", "workers.celery_app", "worker", "-Q", "summarize", "--hostname=worker-summarize-%h@%h", "--loglevel=info"] env_file: - ../.env environment: DATABASE_URL: ${DATABASE_URL:-postgresql+asyncpg://vidconf:vidconf@postgres:5432/vidconf} REDIS_URL: ${REDIS_URL:-redis://redis:6379/0} PLUGINS_CONFIG_PATH: ${PLUGINS_CONFIG_PATH:-config/plugins.yaml} PYTHONPATH: /app volumes: - ../workers:/app/workers:ro - ../config:/app/config:ro - llm-models:/models/qwen:ro depends_on: redis: condition: service_healthy ``` и одновременно убрать `summarize` из списка очередей базового `worker` (его команда сужается до `-Q celery,notify`), чтобы задачи не выполнялись дважды разными процессами (Celery доставляет задачу ровно одному consumer'у одной и той же очереди — дублирования не будет, но держать лишний неиспользуемый consumer незачем). Аналогично можно выделить `notify` в `worker-notify` (та же схема, `-Q notify`) — например, если рассылка приглашений/уведомлений на большую аудиторию (SMTP-латентность) начинает задерживать саммаризацию соседних сеансов при их совместном обслуживании базовым `worker`. ## 4. Масштабирование реплик через `docker compose up --scale` Если одной выделенной очереди тоже мало (несколько ядер CPU для суммаризации/уведомлений), реплицируем сервис командой `--scale`: ```bash docker compose -f deploy/docker-compose.yml --env-file .env -f deploy/docker-compose.override.yml \ up -d --scale worker-summarize=3 ``` Требования для корректного масштабирования сервиса: 1. **Без `-B`** на масштабируемом сервисе (правило §2). 2. **Без фиксированного `--hostname`** на весь сервис — при нескольких репликах одинаковое имя узла Celery приведёт к конфликту регистрации в кластере (соединения будут путаться, `celery inspect` начнёт видеть произвольного из реплик). Использовать `%h` (имя контейнера, уникальное у каждой реплики Compose) — см. `--hostname=worker-summarize-%h@%h` в примере выше. Это же ограничение действует и для `worker-transcriber`: его текущий healthcheck (`--destination worker-transcriber@localhost`) жёстко завязан на единственную реплику; при `--scale worker-transcriber=N` healthcheck и `--hostname` в `deploy/docker-compose.yml` потребуется поменять на шаблон `%h` (и убрать `--destination`, либо адресовать каждую реплику отдельно) — вне рамок этого документа, т.к. правит основной `deploy/docker-compose.yml` (devops-часть). 3. **Без `container_name`** и без фиксированных host-портов на масштабируемом сервисе (у Celery-воркеров портов нет — ограничение не актуально для `worker*`, но актуально, если аналогичный приём применяется к `backend` за `nginx upstream`). 4. Ресурсные лимиты (`cpus`/`mem_limit`) в Compose-файле применяются к КАЖДОЙ реплике, а не разделяются между ними — планировать суммарное потребление хоста (`N × cpus`). Для `worker-transcriber` тот же приём уже применим на уровне профиля: ```bash docker compose -f deploy/docker-compose.yml --env-file .env --profile transcribe \ up -d --scale worker-transcriber=2 ``` (после снятия ограничения фиксированного `--hostname`, см. пункт 2 выше). ## 5. Когда масштабировать Ориентир — глубина очереди (`vidconf_celery_queue_depth`) растущая дольше 15 минут (алерт `QueueGrowing`, `deploy/monitoring/alerts.yml`) либо устойчиво положительная под обычной нагрузкой инстанса. Для `summarize` и `transcription` также ориентир — доля CPU/GPU (транскрибация и суммаризация — тяжёлые по вычислениям шаги, `docs/deploy/hardware-profiles.md`); `notify` почти всегда I/O-bound (SMTP) — реплики полезны в первую очередь при большой аудитории рассылок (закреплённые конференции с длинным списком участников/приглашённых).