Пошук уроків, статей та іншого контенту
Будуйте черги з повторними спробами, пріоритетами й окремими job workers для обробки задач.
Черга завдань розділяє:
код, який створює завдання;
код, який зберігає завдання в очікуванні;
воркери, які виконують завдання.
Наприклад, HTTP-обробник не повинен сам генерувати великий звіт. Він може додати завдання до черги й одразу повернути відповідь клієнту. Окремий воркер сформує звіт у фоновому режимі.
Типове завдання містить:
id — унікальний ідентифікатор;
type — тип операції;
payload — дані для виконання;
priority — пріоритет;
attempt — номер поточної спроби;
maxAttempts — максимальна кількість спроб;
availableAt — час, після якого завдання можна виконувати.
Для простоти нижче буде реалізована черга в пам’яті процесу Node.js.
Така черга підходить для навчання та локальних сценаріїв. Якщо процес завершиться, усі завдання в пам’яті буде втрачено. Для production потрібне зовнішнє сховище або готова система черг із постійним зберіганням.
Черга може обирати завдання не лише за часом додавання, а й за пріоритетом.
Наприклад:
priority: 10 — термінове завдання;
priority: 5 — звичайне завдання;
priority: 1 — низький пріоритет.
Чим більше значення, тим раніше завдання має бути виконане. Якщо пріоритет однаковий, завдання виконуються в порядку додавання.
Пріоритет змінює порядок вибору, але не перериває завдання, яке вже виконується. Якщо воркер почав обробляти завдання з низьким пріоритетом, він не зможе автоматично зупинити його через появу важливішого завдання.
Тимчасові помилки не завжди означають, що завдання неможливо виконати. Наприклад:
зовнішній сервіс тимчасово недоступний;
мережевий запит завершився тайм-аутом;
база даних тимчасово перевантажена.
Після помилки воркер може повернути завдання до черги. Кількість спроб потрібно обмежувати, інакше несправне завдання може повторюватися нескінченно.
Поширена схема:
Воркер бере завдання з черги.
Збільшує номер спроби.
Виконує обробник.
Якщо обробка успішна — завдання завершується.
Якщо сталася помилка і спроби ще залишилися — завдання повертається до черги.
Перед повтором робиться затримка.
Затримку між спробами називають backoff. Один із простих варіантів — експоненційна затримка:
затримка = базова_затримка × 2^(номер_спроби - 1)Наприклад, при базовій затримці 200 мс:
після першої помилки — 200 мс;
після другої — 400 мс;
після третьої — 800 мс.
Різні типи завдань часто мають різні вимоги.
Наприклад:
відправлення електронних листів може виконуватися кількома воркерами;
генерація звітів може бути важкою операцією, тому для неї потрібен один воркер;
обробка зображень може використовувати окрему чергу.
Окремі черги та воркери дають змогу:
незалежно масштабувати різні типи роботи;
не блокувати швидкі завдання повільними;
задавати різну кількість паралельних обробників;
окремо налаштовувати повторні спроби.
Наведений приклад використовує тільки стандартні можливості Node.js. Він містить:
дві незалежні черги;
пріоритети;
повторні спроби;
експоненційний backoff;
два воркери для електронних листів;
один воркер для звітів;
коректне завершення роботи після сигналу зупинки.
Збережіть код у файл queue-workers.js і запустіть командою node queue-workers.js.
'use strict';
const { randomUUID } = require('node:crypto');
const sleep = (milliseconds) =>
new Promise((resolve) => setTimeout(resolve, milliseconds));
class JobQueue {
constructor(name) {
this.name = name;
this.jobs = [];
this.sequence = 0;
}
enqueue(input) {
const job = {
id: input.id ?? randomUUID(),
type: input.type,
payload: input.payload,
priority: input.priority ?? 0,
attempt: input.attempt ?? 0,
maxAttempts: input.maxAttempts ?? 3,
availableAt: input.availableAt ?? Date.now(),
sequence: this.sequence++,
};
this.jobs.push(job);
return job;
}
dequeue() {
const now = Date.now();
const readyJobs = this.jobs
.filter((job) => job.availableAt <= now)
.sort((a, b) => {
if (b.priority !== a.priority) {
return b.priority - a.priority;
}
return a.sequence - b.sequence;
});
const nextJob = readyJobs[0];
if (!nextJob) {
return null;
}
const index = this.jobs.indexOf(nextJob);
this.jobs.splice(index, 1);
return nextJob;
}
get size() {
return this.jobs.length;
}
}
class JobWorker {
constructor({ name, queue, handler, pollInterval = 100 }) {
this.name = name;
this.queue = queue;
this.handler = handler;
this.pollInterval = pollInterval;
}
async run(signal) {
console.log(`[${this.name}] запущено`);
while (!signal.aborted) {
const job = this.queue.dequeue();
if (!job) {
await sleep(this.pollInterval);
continue;
}
job.attempt += 1;
try {
console.log(
`[${this.name}] початок job=${job.id}, ` +
`type=${job.type}, priority=${job.priority}, ` +
`attempt=${job.attempt}`,
);
await this.handler(job);
console.log(
`[${this.name}] успішно завершено job=${job.id}`,
);
} catch (error) {
console.error(
`[${this.name}] помилка job=${job.id}: ${error.message}`,
);
if (job.attempt < job.maxAttempts) {
const baseDelay = 200;
const delay = baseDelay * 2 ** (job.attempt - 1);
job.availableAt = Date.now() + delay;
this.queue.enqueue(job);
console.log(
`[${this.name}] повторна спроба job=${job.id} ` +
`через ${delay} мс`,
);
} else {
console.error(
`[${this.name}] job=${job.id} остаточно провалено ` +
`після ${job.attempt} спроб`,
);
}
}
}
console.log(`[${this.name}] зупинено`);
}
}
async function sendEmail(job) {
await sleep(250);
console.log(
` Відправлення листа користувачу ${job.payload.userId}`,
);
}
async function generateReport(job) {
await sleep(700);
// Імітуємо тимчасову помилку під час першої спроби
if (job.id === 'report-1' && job.attempt === 1) {
throw new Error('сервіс звітів тимчасово недоступний');
}
console.log(
` Генерація звіту "${job.payload.title}"`,
);
}
async function main() {
const emailQueue = new JobQueue('emails');
const reportQueue = new JobQueue('reports');
emailQueue.enqueue({
id: 'email-low',
type: 'send-email',
payload: { userId: 101 },
priority: 1,
maxAttempts: 3,
});
emailQueue.enqueue({
id: 'email-important',
type: 'send-email',
payload: { userId: 202 },
priority: 10,
maxAttempts: 3,
});
emailQueue.enqueue({
id: 'email-normal',
type: 'send-email',
payload: { userId: 303 },
priority: 5,
maxAttempts: 3,
});
reportQueue.enqueue({
id: 'report-1',
type: 'generate-report',
payload: { title: 'Продажі за місяць' },
priority: 10,
maxAttempts: 3,
});
reportQueue.enqueue({
id: 'report-2',
type: 'generate-report',
payload: { title: 'Активність користувачів' },
priority: 1,
maxAttempts: 2,
});
const stopController = new AbortController();
const emailWorker1 = new JobWorker({
name: 'email-worker-1',
queue: emailQueue,
handler: sendEmail,
});
const emailWorker2 = new JobWorker({
name: 'email-worker-2',
queue: emailQueue,
handler: sendEmail,
});
const reportWorker = new JobWorker({
name: 'report-worker-1',
queue: reportQueue,
handler: generateReport,
});
const workers = [
emailWorker1.run(stopController.signal),
emailWorker2.run(stopController.signal),
reportWorker.run(stopController.signal),
];
// Даємо воркерам час обробити завдання
setTimeout(() => {
stopController.abort();
}, 4000);
await Promise.all(workers);
console.log(`Залишилося email-завдань: ${emailQueue.size}`);
console.log(`Залишилося report-завдань: ${reportQueue.size}`);
}
main().catch((error) => {
console.error('Критична помилка:', error);
process.exitCode = 1;
});JobQueue зберігає завдання в масиві. Метод dequeue():
знаходить лише завдання, для яких настав час виконання;
сортує їх за спаданням priority;
при однаковому пріоритеті використовує sequence;
видаляє та повертає перше завдання.
Поле availableAt потрібне для повторних спроб. Якщо завдання повернули до черги з майбутнім часом availableAt, воркер тимчасово його не вибирає.
JobWorker постійно перевіряє чергу:
while (!signal.aborted) {
const job = queue.dequeue();
if (!job) {
await sleep(pollInterval);
continue;
}
// обробка завдання
}Якщо черга порожня, воркер не завершується, а чекає перед наступною перевіркою. Це називають polling.
У прикладі є два воркери електронних листів:
const emailWorker1 = new JobWorker({
name: 'email-worker-1',
queue: emailQueue,
handler: sendEmail,
});
const emailWorker2 = new JobWorker({
name: 'email-worker-2',
queue: emailQueue,
handler: sendEmail,
});Вони конкурують за завдання в одній черзі, тому можуть обробляти два листи паралельно.
Для звітів використовується інша черга й один воркер:
const reportWorker = new JobWorker({
name: 'report-worker-1',
queue: reportQueue,
handler: generateReport,
});Тому обробка звітів не впливає на кількість паралельних відправлень листів.
Параметр maxAttempts задає загальну кількість дозволених виконань, включно з першою спробою.
Наприклад:
{
maxAttempts: 3
}означає:
перша спроба;
повторна спроба після першої помилки;
остання спроба після другої помилки.
Після третьої помилки завдання не повертається до черги:
if (job.attempt < job.maxAttempts) {
// Завдання можна повторити
} else {
// Завдання остаточно провалено
}У реальній системі остаточно провалені завдання зазвичай потрібно зберігати окремо. Це дає змогу:
переглянути причину помилки;
повторити завдання вручну;
дослідити проблемний payload;
побудувати статистику невдалих операцій.
Воркер повинен мати спосіб завершити роботу. У прикладі для цього використовується AbortController:
const controller = new AbortController();
worker.run(controller.signal);
controller.abort();Воркер перевіряє signal.aborted перед кожною новою ітерацією. Поточне завдання не переривається штучно: спочатку воно завершується, а потім воркер зупиняється.
Це важливо для цілісності операцій. Якщо примусово зупинити процес посеред запису у файл або запиту до зовнішнього сервісу, результат може бути неповним.
Навчальна черга з масивом показує основну модель, але production-рішення має враховувати додаткові властивості.
Масив у пам’яті зникає після завершення процесу. Якщо завдання не можна втрачати, його стан потрібно зберігати поза процесом воркера.
Після отримання завдання черга повинна позначити його як таке, що виконується. Якщо воркер аварійно завершиться, завдання має повернутися до обробки після тайм-ауту.
Інакше завдання може назавжди залишитися у стані «виконується».
Повторна спроба може виконати частину операції вдруге. Обробник завдання має бути ідемпотентним або перевіряти, чи операція вже була виконана.
Наприклад, перед повторною відправкою листа можна перевірити в базі даних, чи вже існує запис про успішне відправлення.
Збільшення кількості воркерів не завжди прискорює систему. Зовнішній сервіс, база даних або файлове сховище можуть мати власні обмеження.
Кількість воркерів потрібно добирати з урахуванням:
часу виконання завдання;
доступних ресурсів;
обмежень зовнішніх API;
допустимого навантаження на базу даних.
Без maxAttempts одне несправне завдання може назавжди займати воркер або створювати постійне навантаження.
Завжди задавайте максимальну кількість спроб і визначайте, що робити після її вичерпання.
Не кожну помилку потрібно повторювати. Некоректний payload або відсутній користувач зазвичай не виправляться після затримки.
Корисно розділяти:
тимчасові помилки — можна повторити;
постійні помилки — потрібно завершити без повтору.
Якщо повторювати завдання одразу, система може створити лавину запитів до вже несправного сервісу.
Використовуйте затримку між спробами, а для багатьох одночасних завдань додатково може бути потрібен випадковий компонент затримки — jitter.
Пріоритет має бути частиною алгоритму вибору завдань, а не лише полем у об’єкті. Якщо черга завжди бере перше додане завдання, поле priority нічого не змінює.
Якщо важкі звіти й швидкі листи обробляються однією чергою, звіти можуть затримувати листи.
Для різних ресурсних профілів краще використовувати окремі черги та воркери.
Якщо головний процес завершується одразу після додавання завдань, воркери можуть не встигнути їх обробити. Потрібно коректно очікувати завершення воркерів або використовувати довгоживучий процес.
Черга відокремлює створення завдань від їхнього виконання.
Воркер отримує завдання, виконує його та фіксує результат.
Пріоритети дають змогу обробляти важливіші завдання раніше.
Поля attempt і maxAttempts обмежують кількість повторів.
Backoff зменшує навантаження під час тимчасових помилок.
Окремі черги та воркери дозволяють незалежно масштабувати різні типи роботи.
Для надійної системи потрібні постійне зберігання, контроль стану завдання та ідемпотентні обробники.