LLM-вызов — худший гость в HTTP request-response цикле: медленный, внешний, иногда падающий и дорогой при повторении. Всё, что длиннее UX-бюджета в несколько секунд, должно жить в очереди задач. Для Python-стека это чаще всего Celery + Redis, и на этом стеке отлично видны все грабли.
Ситуация из реальной команды
Команда выносит обработку документов в Celery: одна задача process_document, внутри — шесть шагов пайплайна из главы 13. Первый месяц всё хорошо. Потом:
- на большом документе воркер падает по памяти на шаге 5 — Celery ретраит задачу, шаги 1–4 выполняются заново, счёт за токены удваивается;
- деплой перезапускает воркеры — двадцать «зависших» документов рестартуют с нуля, provider выдаёт rate limit, ретраи наслаиваются;
- один клиент грузит 500 документов, и его batch на полтора часа блокирует очередь для всех остальных;
- в Redis копятся результаты по 2 МБ (кто-то возвращал полный текст из задачи), и однажды Redis выселяет ключи брокера.
Celery здесь ни при чём. Проблемы появились потому, что LLM-пайплайн оформили как одну большую задачу без сохранённого состояния.
Принцип: задача = шаг, состояние = PostgreSQL
Правильная декомпозиция:
- одна Celery-задача = один шаг пайплайна, у которого есть вход, выход и право упасть в одиночестве;
- состояние пайплайна живёт в PostgreSQL, а не в аргументах задач и не в Redis;
- Redis — это транспорт, брокер сообщений, а не хранилище результатов.
create table pipeline_run (
id uuid primary key,
pipeline text not null, -- document_extract:v2
entity_id text not null, -- document_id
status text not null, -- pending / running / done / failed
current_step text,
idempotency_key text unique not null,
created_at timestamptz not null default now(),
updated_at timestamptz not null default now()
);
create table pipeline_step_result (
run_id uuid references pipeline_run(id),
step text not null,
status text not null,
attempt int not null default 0,
output jsonb,
input_tokens int,
output_tokens int,
cost_usd numeric(12, 6),
error text,
finished_at timestamptz,
primary key (run_id, step)
);Задача-шаг при старте смотрит в pipeline_step_result: если шаг уже done — отдаёт сохранённый результат и не трогает модель. Это превращает retry из «сжечь всё ещё раз» в «продолжить с места падения».
@celery_app.task(bind=True, max_retries=2, acks_late=True)
def run_step(self, run_id: str, step: str) -> None:
if step_already_done(run_id, step):
enqueue_next_step(run_id, step)
return
try:
output = execute_step(run_id, step) # внутри — LLM-вызов со своим timeout
except TransientLLMError as exc:
raise self.retry(exc=exc, countdown=backoff(self.request.retries))
except PermanentError as exc:
mark_run_failed(run_id, step, str(exc)) # в ручную очередь, без ретраев
return
save_step_result(run_id, step, output)
enqueue_next_step(run_id, step)Ошибки бывают двух сортов, и это главное
Retry имеет смысл только для транзиентных ошибок: таймаут, 429, 5xx провайдера, обрыв сети. Ретраить постоянные ошибки — невалидный вход, документ не парсится, контент отклонён политикой модели — значит жечь деньги на предсказуемый результат. Маппинг «исключение → сорт ошибки» должен жить в клиенте из главы 4, а не размазываться по задачам.
Для транзиентных — экспоненциальный backoff с jitter. Без jitter сто упавших задач ретраятся синхронно и устраивают провайдеру второй rate limit.
Celery-специфика, которая кусается именно с LLM
acks_late=True+ идемпотентность шагов. Иначе убитый воркер теряет задачу молча. Сacks_lateзадача выполнится повторно — что безопасно, только если шаг умеет не повторять сделанное.task_time_limitбольше, чем timeout LLM-вызова. Если Celery убивает задачу раньше, чем httpx отдаст таймаут, вы не увидите ни ошибки, ни usage. Правило:llm_timeout < task_soft_time_limit < task_time_limit.- Отдельные очереди. Минимум:
llm_interactive(пользователь ждёт статуса) иllm_batch(массовые прогоны). Иначе batch одного клиента откладывает всех. Worker-ы с разной concurrency: LLM-задачи — I/O-bound, им можно высокий prefetch, но помните про rate limit провайдера. worker_prefetch_multiplier=1для длинных задач. Стандартный prefetch заставляет воркер «прихватить» задачи, которые он будет держать часами, пока другие воркеры простаивают.- Не возвращайте большие результаты из задач. Результат — в PostgreSQL, из задачи — только
run_id. Redis-backend с мегабайтными payload-ами — бомба замедленного действия. - Оркестрация цепочки — через
enqueue_next_step, а не жёсткийchain. Celery chains/chords работают, но чинить упавший chord в проде — отдельная профессия. Явная передача «шаг закончился → ставим следующий» проще дебажится и позволяет перезапустить пайплайн с любого шага.
Статус для пользователя
Асинхронная обработка задаёт UX-контракт: клиенту нужен способ узнать, что происходит.
POST /documents/{id}/extract → 202 { "run_id": ... }
GET /runs/{run_id} → { "status": "running",
"current_step": "risk_detect",
"steps_done": 3, "steps_total": 6,
"cost_so_far_usd": 0.011 }Стоимость в статусе — не украшение: именно здесь владелец фичи впервые видит, что документ обошёлся в доллар, и приходит к вам с правильными вопросами до счёта в конце месяца.
Правило главы
Асинхронный LLM-пайплайн — это конечный автомат в PostgreSQL, у которого Celery всего лишь исполняет переходы. Если состояние живёт в аргументах задач, retry становится роскошью, деплой — лотереей, а счёт за токены — сюрпризом.
Чеклист главы
- Один шаг пайплайна = одна задача; состояние и результаты шагов — в PostgreSQL.
- Выполненный шаг не выполняется повторно при retry — проверено убийством воркера.
- Ошибки разделены на транзиентные и постоянные; ретраятся только первые, с backoff и jitter.
-
acks_late,task_time_limitи timeout LLM-вызова согласованы между собой. - Интерактивные и batch-задачи разведены по очередям.
- Из задач не возвращаются большие payload-ы; Redis — брокер, не хранилище.
- У клиента есть endpoint статуса с прогрессом и накопленной стоимостью.
Практикум
Задача 15.1 — пайплайн, переживающий убийство воркера. ⭑⭑⭑
Реализуйте пайплайн обработки документа из трёх шагов (извлечение секций → извлечение полей по секциям → сводка) на Celery + Redis + PostgreSQL по схеме этой главы. Документы — 5–10 сгенерированных «договоров» по 3–5 страниц.
Критерии приёмки:
-
docker killворкера посреди шага 2 → после рестарта пайплайн продолжается с шага 2; шаг 1 не перевыполняется — доказано поllm_call_log(количество вызовов) и поattemptв step_result. - Повторный
POSTтого же документа с тем же idempotency key не создаёт второй run. - Транзиентная ошибка (мок 429 на втором шаге) ретраится с backoff; постоянная (документ-пустышка) сразу уводит run в
failedбез ретраев. -
GET /runs/{id}показывает прогресс и стоимость в реальном времени.
Подсказки:
- Сначала напишите конечный автомат и таблицы, потом Celery. Если начать с Celery, состояние неизбежно расползётся по аргументам задач.
- Тест на идемпотентность шага пишется без Celery вообще: вызовите
run_stepдважды синхронно и посчитайте LLM-вызовы.
Задача 15.2 — batch, который не мешает людям. ⭑⭑
Добавьте к пайплайну из 15.1 массовую загрузку: endpoint принимает 100 документов разом. Разведите интерактивную и batch-обработку по очередям и воркерам.
Критерии приёмки:
- Пока batch из 100 документов обрабатывается, одиночный интерактивный документ проходит за обычное время — измерено.
- Batch уважает rate limit провайдера: параллелизм ограничен, при 429 скорость снижается, а не нарастает лавина ретраев.
- Прогресс batch-а виден одним запросом: сколько документов в каком статусе, сколько потрачено.
- Отмена batch-а останавливает невыполненные шаги и не оставляет «зомби-задач» в очереди.
Подсказки:
- Ограничение параллелизма к провайдеру удобно держать не в Celery concurrency, а семафором вокруг LLM-клиента — тогда лимит общий для всех очередей.