Практика эксплуатации
Очередь задач AI-агента: как пережить всплеск нагрузки
Если запускать каждое входящее задание сразу, одновременные вызовы модели, инструменты и история диалога начинают конкурировать за память, процессор и сетевые соединения. Решение — поставить задания в очередь задач, ограничить число активных обработчиков и отклонять либо задерживать новые запросы, когда внутренний буфер заполнен.
Что именно нужно контролировать
Ограничение параллельности отвечает на вопрос «сколько заданий выполняется сейчас». Backpressure — «что происходит с новым заданием, когда система уже не успевает». Это разные механизмы: один без другого лишь переносит проблему.
Например, лимит в четыре обработчика защищает исполнитель, но не память: если принять сто тысяч заданий в массив, процесс всё равно может исчерпать доступный объём памяти. Поэтому очередь должна иметь конечную ёмкость.
Целевая схема: входящий запрос → проверка ёмкости → ограниченная очередь → фиксированное число обработчиков → результат.
Для минимальной рабочей конфигурации понадобятся четыре параметра:
CONCURRENCY— максимальное число одновременно выполняемых заданий;MAX_QUEUE— максимальное число ожидающих заданий;TASK_TIMEOUT_MS— предельное время одного задания;RETRY_AFTER_SECONDS— подсказка перегруженному клиенту, когда повторить запрос.
Шаг 1. Создайте ограниченную очередь
Ниже — самостоятельный пример на Node.js без сторонних пакетов. Это учебная реализация, а не утверждение о конкретной производственной системе. Она показывает механику: очередь не принимает больше заданного числа ожидающих работ, а число активных работ никогда не превышает лимит.
class QueueFullError extends Error {
constructor(message = "Task queue is full") {
super(message);
this.name = "QueueFullError";
}
}
class AgentQueue {
constructor({ concurrency, maxQueue, taskTimeoutMs }) {
if (!Number.isInteger(concurrency) || concurrency < 1) {
throw new TypeError("concurrency must be a positive integer");
}
if (!Number.isInteger(maxQueue) || maxQueue < 0) {
throw new TypeError("maxQueue must be a non-negative integer");
}
if (!Number.isInteger(taskTimeoutMs) || taskTimeoutMs < 1) {
throw new TypeError("taskTimeoutMs must be a positive integer");
}
this.concurrency = concurrency;
this.maxQueue = maxQueue;
this.taskTimeoutMs = taskTimeoutMs;
this.active = 0;
this.pending = [];
this.accepted = 0;
this.rejected = 0;
this.completed = 0;
this.failed = 0;
}
submit(task) {
if (typeof task !== "function") {
return Promise.reject(new TypeError("task must be a function"));
}
if (this.active >= this.concurrency &&
this.pending.length >= this.maxQueue) {
this.rejected += 1;
return Promise.reject(new QueueFullError());
}
this.accepted += 1;
return new Promise((resolve, reject) => {
this.pending.push({ task, resolve, reject });
this.drain();
});
}
drain() {
while (
this.active < this.concurrency &&
this.pending.length > 0
) {
const job = this.pending.shift();
this.active += 1;
void this.run(job);
}
}
async run({ task, resolve, reject }) {
let timer;
try {
const timeout = new Promise((_, rejectTimeout) => {
timer = setTimeout(() => {
rejectTimeout(new Error("Task timed out"));
}, this.taskTimeoutMs);
});
const result = await Promise.race([
Promise.resolve().then(task),
timeout
]);
this.completed += 1;
resolve(result);
} catch (error) {
this.failed += 1;
reject(error);
} finally {
clearTimeout(timer);
this.active -= 1;
this.drain();
}
}
stats() {
return {
active: this.active,
queued: this.pending.length,
accepted: this.accepted,
rejected: this.rejected,
completed: this.completed,
failed: this.failed
};
}
}
Проверка ёмкости учитывает только ожидающие задания. Если есть свободный обработчик, новое задание можно принять даже при MAX_QUEUE=0. Это позволяет использовать очередь как «нулевой буфер»: выполнять только то, для чего уже есть свободная ёмкость.
Шаг 2. Подключите обработчик агента
Не помещайте секреты, токены или полные пользовательские запросы в пример и диагностический вывод. Функция runAgent ниже намеренно условна: замените её вызовом собственного агента и передавайте секреты через штатное защищённое хранилище среды.
const queue = new AgentQueue({
concurrency: readInteger("AGENT_CONCURRENCY", 4, 1, 64),
maxQueue: readInteger("AGENT_MAX_QUEUE", 100, 0, 10000),
taskTimeoutMs: readInteger(
"AGENT_TASK_TIMEOUT_MS",
120000,
1000,
900000
)
});
function readInteger(name, fallback, min, max) {
const raw = process.env[name];
if (raw === undefined) return fallback;
const value = Number(raw);
if (!Number.isInteger(value) || value < min || value > max) {
throw new Error(`${name} must be an integer from ${min} to ${max}`);
}
return value;
}
async function runAgent(input, { signal }) {
// Пример: здесь вызывается ваш исполнитель агента.
// Он должен учитывать AbortSignal и освобождать свои ресурсы.
if (signal.aborted) throw signal.reason;
return { taskId: input.taskId, status: "done" };
}
function makeTask(input) {
return async () => {
const controller = new AbortController();
return runAgent(input, { signal: controller.signal });
};
}
Важное ограничение: Promise.race прекращает ожидание, но сам по себе не останавливает фоновую операцию. Реальный исполнитель должен поддерживать отмену. Практичнее передавать созданный очередью AbortSignal и вызывать abort() при тайм-ауте; конкретная интеграция зависит от используемого HTTP-клиента, SDK и инструментов агента.
Шаг 3. Верните явный сигнал перегрузки
Для синхронного HTTP API переполнение очереди удобно обозначать кодом 503 Service Unavailable и заголовком Retry-After. Код 429 также применим, если перегрузка оформлена как лимит запросов конкретного потребителя. Главное — выбрать одну семантику и документировать её.
async function handleTaskRequest(request, response) {
const input = await readValidatedInput(request);
try {
const result = await queue.submit(makeTask(input));
response.writeHead(200, {
"Content-Type": "application/json; charset=utf-8"
});
response.end(JSON.stringify(result));
} catch (error) {
if (error instanceof QueueFullError) {
response.writeHead(503, {
"Content-Type": "application/json; charset=utf-8",
"Retry-After": "5",
"Cache-Control": "no-store"
});
response.end(JSON.stringify({
error: "agent_overloaded",
retryable: true
}));
return;
}
response.writeHead(500, {
"Content-Type": "application/json; charset=utf-8",
"Cache-Control": "no-store"
});
response.end(JSON.stringify({
error: "task_failed",
retryable: false
}));
}
}
Функция readValidatedInput здесь является заглушкой. В рабочем сервисе она должна ограничивать размер тела запроса, проверять типы и удалять поля, которые агенту не нужны. Не записывайте необработанный ввод в логи.
Шаг 4. Задайте безопасные начальные пределы
Следующие значения — пример отправной точки, а не универсальная рекомендация. Они не содержат секретов и безопасны для локального эксперимента:
AGENT_CONCURRENCY=4 \
AGENT_MAX_QUEUE=100 \
AGENT_TASK_TIMEOUT_MS=120000 \
node server.js
Не увеличивайте CONCURRENCY, пока не измерили пиковое потребление памяти одним заданием. Приближённая верхняя граница:
доступная память для заданий
────────────────────────────── = теоретический предел параллельности
пиковая память одного задания
Затем уменьшите результат, оставив запас для процесса, библиотек, сетевых буферов и сборщика мусора. Если задания неоднородны, ориентируйтесь не на среднее, а на высокий наблюдаемый перцентиль и проверяйте систему нагрузочным экспериментом в своей среде.
Шаг 5. Проверьте результат воспроизводимым сценарием
Для проверки механизма не нужен реальный вызов модели. Временная функция ниже имитирует работу длительностью 200 мс и фиксирует максимальное число одновременных запусков:
let running = 0;
let observedMax = 0;
function simulatedTask(id) {
return async () => {
running += 1;
observedMax = Math.max(observedMax, running);
try {
await new Promise((resolve) => setTimeout(resolve, 200));
return id;
} finally {
running -= 1;
}
};
}
const checkQueue = new AgentQueue({
concurrency: 2,
maxQueue: 3,
taskTimeoutMs: 2000
});
const attempts = Array.from({ length: 8 }, (_, id) =>
checkQueue.submit(simulatedTask(id))
.then(() => ({ outcome: "completed" }))
.catch((error) => ({ outcome: error.name }))
);
const outcomes = await Promise.all(attempts);
console.log({
observedMax,
outcomes,
stats: checkQueue.stats()
});
Сохраните очередь и проверочный фрагмент в локальном файле, например queue-check.mjs, затем выполните:
node queue-check.mjs
Проверка считается успешной, если одновременно выполняются следующие инварианты:
observedMaxне превышает2;stats.activeиstats.queuedпосле завершения равны нулю;- часть одновременных попыток завершается с
QueueFullError, потому что ёмкость исполнителей и буфера конечна; accepted + rejectedравно общему числу попыток;- процесс завершается самостоятельно и не оставляет активных таймеров.
Не закрепляйте в проверке порядок завершения разных заданий, если он не является контрактом системы. При параллельном исполнении такой порядок может меняться.
Наблюдаемость и настройка
Одного счётчика длины очереди недостаточно. Минимально полезный набор метрик:
- текущие
activeиqueued; - число принятых, завершённых, ошибочных и отклонённых заданий;
- время ожидания до начала исполнения отдельно от времени обработки;
- доля тайм-аутов и отмен;
- потребление памяти процесса и частота длительных пауз сборщика мусора.
Стабильно заполненная очередь означает, что входящий поток превышает пропускную способность. Увеличение буфера лишь откладывает отказ и повышает задержку. Сначала сокращайте стоимость задания, масштабируйте исполнителей либо ограничивайте источник нагрузки; размер очереди меняйте только после измерений.
Типовые ошибки
Безразмерный массив ожидания
Лимит обработчиков есть, но все входящие задания сохраняются в памяти. Под всплеском растут задержка, память и вероятность аварийного завершения. Исправление: конечный MAX_QUEUE и явный отказ при переполнении.
Повторы без случайной задержки
Если все клиенты повторяют запрос ровно через пять секунд, возникает новый синхронный всплеск. Клиентам следует использовать экспоненциальную задержку со случайным разбросом и ограничивать число попыток.
Тайм-аут без отмены
HTTP-ответ уже вернул ошибку, а модель или инструмент продолжает работу и занимает слот внешнего ресурса. Используйте отменяемые операции и проверяйте, что соединения и файловые дескрипторы освобождаются.
Смешение коротких и тяжёлых заданий
Несколько долгих заданий могут заблокировать лёгкие. Возможные меры: отдельные очереди по классу стоимости, взвешенные лимиты или планировщик с приоритетами. Приоритеты требуют защиты от вечного голодания низкоприоритетных работ.
Повтор неидемпотентной операции
После тайм-аута клиент не всегда знает, выполнилось ли действие. Для операций с побочными эффектами используйте идентификатор задания или ключ идемпотентности и храните состояние выполнения вне памяти одного процесса.
Логи с контекстом агента
История диалога, результаты инструментов и системные инструкции могут содержать чувствительные данные. Для метрик достаточно идентификатора, статуса, длительности и размеров; содержимое задания логируйте только по явно определённой политике.
Ограничения решения
Очередь в памяти подходит для одного процесса и допускает потерю ожидающих заданий при перезапуске. Она не обеспечивает распределённую координацию, долговечное хранение, повтор после сбоя или строгое выполнение ровно один раз.
При нескольких экземплярах приложения каждый экземпляр применяет собственный лимит. Общая параллельность становится примерно равна произведению числа экземпляров на локальный лимит, если нет внешнего координатора. Это важно, когда все экземпляры обращаются к одной модели, базе данных или другому общему ресурсу.
Для долговечных фоновых работ понадобится внешний брокер или хранилище состояния. При переходе к нему сохраните те же инварианты: ограниченную выдачу работ, контроль числа активных заданий, видимый сигнал перегрузки, отмену, дедлайн и идемпотентную обработку. Выбор конкретной технологии должен следовать из требований к доставке и восстановлению, а не только из максимальной пропускной способности.
Критерий готовности
Система переживает всплеск предсказуемо, если число активных заданий ограничено, очередь имеет конечный размер, переполнение быстро возвращает документированную ошибку, тайм-ауты освобождают реальные ресурсы, а метрики позволяют отличить ожидание от исполнения. После этого рост нагрузки превращается из неконтролируемого расхода памяти в управляемую деградацию сервиса.