Зачем агенту очередь, если API отвечает сразу
Синхронный вызов удобен до первого длинного анализа, временного отказа модели или всплеска трафика. Очередь отделяет обещание «задача принята» от фактического выполнения. Веб-слой быстро валидирует запрос, создаёт job_id и возвращает статус, а worker забирает работу тогда, когда есть квота модели, память и доступные tool-сервисы.
Это не просто ускоритель. Очередь становится границей управления: здесь задаются приоритет, дедлайн, tenant, бюджет токенов, допустимые инструменты и политика повторов.
Контракт задания
Минимальный envelope содержит job_id, tenant_id, тип операции, версию workflow, ссылку на входные данные, дедлайн и idempotency key. Большие документы не кладите в сообщение: сохраняйте их в объектном хранилище, а в очередь передавайте неизменяемую ссылку и checksum.
{
"job_id": "job_01J...",
"workflow": "research_report:v3",
"input_ref": "s3://bucket/input/sha256...",
"deadline_at": "2027-03-29T12:00:00Z",
"attempt": 1,
"trace_id": "..."
}
At-least-once означает дубликаты
Надёжные брокеры обычно предпочитают не потерять сообщение, поэтому worker может получить его повторно: подтверждение потерялось, процесс упал после внешнего действия или истёк visibility timeout. Нельзя рассчитывать, что «ровно один раз» обеспечит бизнес-результат. Перед каждым необратимым действием проверяйте idempotency key и храните итог операции в транзакционной таблице.
Visibility timeout и heartbeat
После выдачи сообщение временно скрывается от других workers. Timeout должен покрывать нормальную длительность шага с запасом, но не быть настолько большим, чтобы падение зависло на часы. Для непредсказуемо долгих задач worker отправляет heartbeat и продлевает lease. Если heartbeat прекратился, сообщение возвращается в очередь.
Разделяйте общий дедлайн задания и lease одной попытки. Продление lease не должно бесконечно продлевать пользовательский SLA.
Retries без шторма
Повторяйте только временные ошибки: 429, часть 5xx, сетевые обрывы и истечение lease. Ошибки валидации, отсутствие прав и превышение бюджета обычно неретрайбельны. Используйте exponential backoff с jitter, ограничение попыток и отдельный retry budget. Иначе массовый сбой провайдера создаст retry storm и окончательно забьёт квоту.
- 1-я попытка — сразу;
- 2-я — через 5–15 секунд;
- 3-я — через 30–90 секунд;
- дальше — медленнее, но только до дедлайна.
DLQ — не кладбище
После исчерпания попыток сообщение уходит в dead-letter queue вместе с кодом ошибки, версией worker, последним checkpoint и ссылкой на trace. Для DLQ нужны владелец, алерт по росту, runbook разбора и безопасный redrive. Перед повторным запуском исправьте причину и отфильтруйте задания, которые уже просрочены или частично выполнили внешнее действие.
Приоритеты и честность
Одна FIFO-очередь заставляет срочные короткие задачи ждать за длинным отчётом. Практичнее несколько классов обслуживания: interactive, batch и maintenance. Внутри tenant применяйте квоты, чтобы один крупный клиент не вытеснил остальных. Приоритет не отменяет дедлайн и cost ceiling.
Метрики и план запуска
Смотрите на age of oldest message, time-in-queue, execution time, success rate по попыткам, DLQ rate, saturation workers, 429 провайдера и стоимость задания. Начните с одного типа workflow и небольшого concurrency, проведите тест падения worker в середине tool-call, затем включите autoscaling по возрасту очереди, а не только по её длине.
Архитектура потока от API до результата
Входной API не должен публиковать сообщение до фиксации задания: сохраните job и outbox event одной транзакцией, затем publisher передаст событие брокеру. Worker сначала захватывает lease, проверяет дедлайн и отмену, затем выполняет шаги и атомарно фиксирует результат. Клиент получает статус через polling, webhook или server-sent events, но ни один канал доставки статуса не считается источником истины.
| Слой | Ответственность | Типичный сбой |
|---|---|---|
| API | валидация, job_id, auth | повтор запроса |
| Broker | буфер и доставка | дубликат |
| Worker | workflow и checkpoint | падение процесса |
| Result store | статус и артефакты | конфликт записи |
Backpressure и autoscaling
Масштабирование только по количеству сообщений обманчиво: десять двухчасовых отчётов тяжелее тысячи коротких классификаций. Оценивайте backlog в секундах работы, возраст p95 и доступную квоту tokens per minute. Новый worker бесполезен, если провайдер уже отвечает 429; в этом случае admission control должен замедлить приём или перевести batch в отложенный класс.
Отмена и дедлайн
Отмена — состояние, а не убийство одного процесса. Worker проверяет cancellation token между шагами, прекращает новые tool calls, отзывает временные полномочия и помечает незавершённые артефакты. Уже начавшийся внешний эффект нужно завершить, сверить или компенсировать. Результат должен различать cancelled, expired, failed и partially_completed.