147 lines
10 KiB
Markdown
147 lines
10 KiB
Markdown
# Горизонтальное масштабирование очередей 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 -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 --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) — реплики полезны в первую очередь
|
||
при большой аудитории рассылок (закреплённые конференции с длинным списком
|
||
участников/приглашённых).
|