Пошук уроків, статей та іншого контенту
Виносьте тривалі операції у фон і проєктуйте надійне виконання задач поза HTTP-запитом.
HTTP-запит має обмежений життєвий цикл:
клієнт надсилає запит;
сервер виконує обробник;
сервер повертає відповідь;
з'єднання може бути закрито або перервано.
Якщо всередині обробника виконувати тривалу операцію — наприклад, створення звіту, обробку великого файлу чи надсилання великої кількості повідомлень, — користувач змушений чекати завершення всієї роботи.
Це створює проблеми:
HTTP-запит може завершитися за тайм-аутом;
клієнт може закрити з'єднання;
операція може бути перервана під час перезапуску сервера;
одночасні довгі операції займуть усі доступні ресурси;
користувач не отримує швидкої відповіді.
Замість цього HTTP-обробник має лише поставити задачу в чергу й одразу повідомити клієнту, що роботу прийнято.
Клієнт → HTTP API → черга задач → worker → результатЗазвичай API повертає статус 202 Accepted. Це означає: сервер прийняв запит, але робота ще не завершена.
Черга задач зберігає задачі, які потрібно виконати.
Worker — окремий процес або група процесів, які забирають задачі з черги та виконують їх.
Таке розділення дає змогу:
швидко відповідати на HTTP-запити;
обмежувати кількість паралельних задач;
повторювати невдалі задачі;
запускати кілька worker-процесів;
переживати перезапуск HTTP-сервера;
відокремити помилки фонової роботи від помилок API.
setTimeout — не повноцінна чергаТакий код не є надійним способом запуску фонових задач:
server.post("/reports", (request, response) => {
setTimeout(() => {
generateReport();
}, 60_000);
response.send("Прийнято");
});setTimeout лише планує виконання в пам'яті поточного процесу. Якщо процес Node.js завершиться або перезапуститься, задача зникне.
Також setTimeout:
не зберігає стан задачі;
не виконує повторні спроби;
не координує кілька процесів;
не гарантує виконання після аварії.
Для надійного виконання потрібне зовнішнє сховище черги, наприклад Redis або база даних.
BullMQ — бібліотека для черг задач у Node.js, яка використовує Redis як зовнішнє сховище.
У прикладі HTTP-сервер додає задачі до черги, а окремий worker їх обробляє.
Потрібні Node.js, Redis і пакет bullmq.
npm init -y
npm install bullmqЗапустити Redis локально можна, наприклад, у Docker:
docker run --rm --name lesson-redis -p 6379:6379 redis:7Створіть два файли: api.mjs і worker.mjs.
// api.mjs
import { createServer } from "node:http";
import { Queue } from "bullmq";
import { randomUUID } from "node:crypto";
const connection = {
host: "127.0.0.1",
port: 6379
};
const reportsQueue = new Queue("reports", { connection });
function readJson(request) {
return new Promise((resolve, reject) => {
let body = "";
request.setEncoding("utf8");
request.on("data", (chunk) => {
body += chunk;
if (body.length > 1_000_000) {
reject(new Error("Тіло запиту завелике"));
request.destroy();
}
});
request.on("end", () => {
try {
resolve(body ? JSON.parse(body) : {});
} catch {
reject(new Error("Некоректний JSON"));
}
});
request.on("error", reject);
});
}
const server = createServer(async (request, response) => {
if (request.method !== "POST" || request.url !== "/reports") {
response.writeHead(404, { "Content-Type": "application/json" });
response.end(JSON.stringify({ error: "Маршрут не знайдено" }));
return;
}
try {
const body = await readJson(request);
if (!body.userId) {
response.writeHead(400, { "Content-Type": "application/json" });
response.end(JSON.stringify({ error: "Поле userId є обов'язковим" }));
return;
}
const reportId = randomUUID();
const job = await reportsQueue.add(
"generate-report",
{
reportId,
userId: body.userId
},
{
attempts: 5,
backoff: {
type: "exponential",
delay: 1_000
},
removeOnComplete: 100,
removeOnFail: 1_000
}
);
response.writeHead(202, { "Content-Type": "application/json" });
response.end(
JSON.stringify({
reportId,
jobId: job.id,
status: "queued"
})
);
} catch (error) {
console.error("Не вдалося поставити задачу в чергу:", error);
response.writeHead(500, { "Content-Type": "application/json" });
response.end(JSON.stringify({ error: "Внутрішня помилка сервера" }));
}
});
server.listen(3000, () => {
console.log("HTTP API працює на http://localhost:3000");
});
async function shutdown() {
console.log("Завершення роботи API...");
server.close(async () => {
await reportsQueue.close();
process.exit(0);
});
}
process.on("SIGINT", shutdown);
process.on("SIGTERM", shutdown);Запит до API:
curl -X POST http://localhost:3000/reports \
-H "Content-Type: application/json" \
-d '{"userId":"user-42"}'Приклад відповіді:
{
"reportId": "6c0d4e8c-2d7f-4c76-ae1c-1c8f2995cc10",
"jobId": "1",
"status": "queued"
}API не чекає на створення звіту. Воно лише додає задачу до Redis і повертає відповідь.
// worker.mjs
import { Worker } from "bullmq";
const connection = {
host: "127.0.0.1",
port: 6379
};
function wait(milliseconds) {
return new Promise((resolve) => {
setTimeout(resolve, milliseconds);
});
}
async function generateReport(data) {
console.log(`Початок створення звіту ${data.reportId}`);
// Імітація тривалої асинхронної операції
await wait(5_000);
console.log(`Звіт ${data.reportId} створено для ${data.userId}`);
}
const worker = new Worker(
"reports",
async (job) => {
console.log(`Отримано задачу ${job.id}: ${job.name}`);
await generateReport(job.data);
// Повернене значення можна використати як результат задачі
return {
reportId: job.data.reportId,
completedAt: new Date().toISOString()
};
},
{
connection,
concurrency: 2
}
);
worker.on("completed", (job, result) => {
console.log(`Задача ${job.id} завершена`, result);
});
worker.on("failed", (job, error) => {
if (job) {
console.error(
`Задача ${job.id} не виконана. Спроба ${job.attemptsMade}:`,
error.message
);
} else {
console.error("Неідентифікована задача не виконана:", error.message);
}
});
worker.on("error", (error) => {
console.error("Помилка worker:", error);
});
async function shutdown() {
console.log("Завершення роботи worker...");
// Worker перестає брати нові задачі й чекає завершення поточної роботи
await worker.close();
process.exit(0);
}
process.on("SIGINT", shutdown);
process.on("SIGTERM", shutdown);
console.log("Worker запущено");Запустіть у двох терміналах:
node worker.mjsnode api.mjsПісля надсилання запиту API одразу поверне відповідь, а worker окремо виконає задачу.
Параметр concurrency: 2 означає, що один worker може обробляти не більше двох задач одночасно. Це не створює два потоки JavaScript, але дозволяє виконувати кілька асинхронних операцій паралельно в межах процесу.
Фонові задачі можуть завершитися помилкою через:
тимчасову недоступність зовнішнього сервісу;
мережеву помилку;
короткочасне перевантаження бази даних;
тимчасову помилку файлової системи.
У такому разі немає сенсу одразу втрачати задачу. У прикладі використано:
{
attempts: 5,
backoff: {
type: "exponential",
delay: 1_000
}
}Це дозволяє виконати до п'яти спроб. Інтервал між спробами збільшується:
приблизно 1 секунда;
приблизно 2 секунди;
приблизно 4 секунди;
і так далі.
Worker має сигналізувати про помилку через throw або відхилення Promise:
const worker = new Worker("reports", async (job) => {
const result = await callExternalService(job.data);
if (!result.ok) {
throw new Error("Зовнішній сервіс повернув помилку");
}
return result;
}, { connection });Якщо worker не повідомить про помилку, черга вважатиме задачу успішною і не запустить повторну спробу.
Через повторні спроби одна задача може бути виконана більше одного разу. Наприклад:
worker створив звіт;
процес завершився до підтвердження успішного завершення;
черга запустила задачу повторно;
звіт створюється ще раз.
Тому фонова задача має бути ідемпотентною — повторне виконання не повинно створювати некоректний результат або дублікати.
Поширені підходи:
перед створенням результату перевіряти, чи він уже існує;
використовувати унікальний reportId;
зберігати статус задачі в базі даних;
використовувати унікальний ключ для операцій над зовнішнім сервісом;
робити повторний виклик безпечним для отримувача.
Наприклад, замість безумовного створення звіту можна виконувати операцію за схемою:
async function generateReport(data) {
const existingReport = await findReportById(data.reportId);
if (existingReport) {
return existingReport;
}
return createReport({
id: data.reportId,
userId: data.userId
});
}Це лише приклад логіки. Перевірка та створення мають бути захищені на рівні бази даних унікальним обмеженням, якщо кілька worker-процесів можуть виконати одну задачу одночасно.
Worker не повинен одразу завершувати процес під час перезапуску. Спочатку потрібно:
припинити отримання нових задач;
дочекатися завершення поточних задач;
закрити з'єднання;
завершити процес.
Саме тому в прикладі використано:
await worker.close();Без коректного завершення процес може бути примусово зупинений посеред операції. Черга зазвичай зможе виявити незавершену задачу та повернути її в обробку, але сама операція все одно має бути ідемпотентною.
Типовий життєвий цикл задачі:
waiting — задача очікує worker;
active — задача виконується;
completed — задача успішно завершена;
failed — задача завершилася помилкою;
повторне виконання після помилки — відповідно до налаштувань повторних спроб.
Клієнту не потрібно чекати завершення задачі в початковому HTTP-запиті. Замість цього API може:
повернути jobId або reportId;
надати окремий endpoint для перевірки статусу;
повідомити клієнта через WebSocket або інший механізм, коли робота завершиться.
Важливо відділяти факт прийняття задачі від факту її завершення:
202 Accepted — задачу прийнято;
200 OK або готовий ресурс — результат доступний;
помилка задачі — фонова операція не завершилася успішно.
Хороша фонова задача:
виконує одну логічну операцію;
має всі необхідні дані для роботи;
не залежить від локального стану HTTP-запиту;
може бути повторена;
має зрозумілий результат або стан помилки.
Не варто передавати в чергу:
об'єкт request;
об'єкт response;
відкритий файловий дескриптор;
з'єднання з базою даних;
великі об'єкти, які можна замінити ідентифікатором.
Краще передавати компактні дані:
await reportsQueue.add("generate-report", {
reportId: "report-123",
userId: "user-42"
});Worker самостійно завантажує необхідні дані за ідентифікаторами.
setTimeoutresponse.end("Прийнято");
setTimeout(() => {
doImportantWork();
}, 10_000);Це працює лише доти, доки живий конкретний процес. Для важливих задач потрібне зовнішнє сховище.
const result = await generateLargeReport();
response.end(JSON.stringify(result));Так HTTP-запит знову стає довгим. Потрібно поставити задачу в чергу та повернути 202.
Тимчасова помилка мережі може призвести до остаточної втрати задачі. Для операцій, які можна повторити, налаштовуйте обмежену кількість спроб і затримку між ними.
Якщо повторний запуск надсилає лист, списує кошти або створює запис без перевірки, можна отримати дублікати. Проєктуйте такі задачі ідемпотентними.
Велика кількість одночасних задач може перевантажити базу даних або зовнішній API. Обмежуйте concurrency і враховуйте ліміти залежностей.
Для кожної задачі корисно логувати:
ідентифікатор задачі;
тип операції;
початок і завершення;
номер спроби;
тривалість;
текст помилки.
Без цього складно зрозуміти, чому задача зависла або повторюється.
Тривалі операції не слід виконувати всередині HTTP-запиту.
API має швидко додати задачу в зовнішню чергу та повернути 202 Accepted.
Worker окремо отримує задачі й виконує їх.
Надійна черга має зберігати задачі поза пам'яттю HTTP-процесу.
Для тимчасових помилок потрібні повторні спроби з затримкою.
Фонові задачі мають бути ідемпотентними, оскільки повторне виконання можливе.
Під час завершення процесу worker має перестати приймати нові задачі та коректно закритися.
setTimeout у процесі Node.js підходить для простого планування, але не замінює надійну чергу задач.