Пошук уроків, статей та іншого контенту
Пояснить роботу черг, буферизацію навантаження, розподіл задач і відокремлення відправників від отримувачів.
Черга повідомлень — це посередник між компонентами системи, який зберігає повідомлення до моменту, коли їх обробить отримувач.
Замість прямої взаємодії:
Сервіс замовлень → Сервіс emailсервіси взаємодіють через чергу:
Сервіс замовлень → Черга → Сервіс emailВідправник називається producer, а отримувач — consumer.
Повідомлення може містити:
назву події або тип задачі;
ідентифікатор сутності;
дані для обробки;
час створення;
ідентифікатор повідомлення.
Наприклад:
{
"id": "message-123",
"type": "send-welcome-email",
"userId": "user-456",
"email": "user@example.com"
}Без черги сервіс замовлень повинен знати:
адресу сервісу email;
формат його API;
час очікування відповіді;
спосіб обробки помилок;
що робити, якщо сервіс email тимчасово недоступний.
Черга зменшує зв’язаність між компонентами. Producer лише додає повідомлення, а consumer самостійно забирає та обробляє його.
Це дозволяє змінювати або масштабувати сервіси незалежно один від одного.
Уявімо, що протягом секунди надходить 10 000 задач, але система може обробити лише 1 000 задач за секунду.
Без черги частина запитів може завершитися помилкою або перевантажити сервіс.
Черга тимчасово зберігає надлишкові задачі:
Вхідний потік: 10 000 задач/с
↓
Черга
↓
Обробка: 1 000 задач/сПісля зменшення навантаження consumer поступово обробить накопичені повідомлення.
Черга не усуває обмеження продуктивності, а розподіляє навантаження в часі.
Не всі операції потрібно виконувати до відповіді користувачу.
Наприклад, після створення замовлення можна:
зберегти замовлення;
додати повідомлення в чергу;
одразу відповісти клієнту;
окремо відправити email або сформувати документ.
Користувач не повинен чекати завершення повільної другорядної операції.
Producer створює та надсилає повідомлення в чергу.
Наприклад, сервіс замовлень може додати задачу:
{
"type": "generate-invoice",
"orderId": "order-42"
}Producer зазвичай не знає:
який саме worker обробить повідомлення;
де він розташований;
скільки worker-ів працює;
коли саме завершиться обробка.
Queue зберігає повідомлення та видає їх consumer-ам.
Типова черга має такі характеристики:
порядок повідомлень;
максимальний розмір;
час життя повідомлень;
політику повторних спроб;
спосіб підтвердження обробки;
поведінку під час переповнення.
Consumer читає повідомлення з черги та виконує задачу.
Після успішної обробки він повинен повідомити чергу, що повідомлення можна видалити. Це називається acknowledgement, або скорочено ack.
Якщо обробка не вдалася, повідомлення можна:
повернути в чергу;
повторити пізніше;
перенести в окрему чергу помилок;
позначити як остаточно невдале.
Кілька worker-ів можуть читати одну чергу. У такому разі кожне повідомлення дістається одному з них.
┌──────────┐
Producer ────────►│ Queue │
└────┬─────┘
│
┌────────────┼────────────┐
▼ ▼ ▼
Worker 1 Worker 2 Worker 3Це називається конкуруючими споживачами.
Переваги:
можна обробляти повідомлення паралельно;
легко збільшити пропускну здатність;
вихід з ладу одного worker-а не зупиняє всю обробку.
Якщо один worker обробляє 10 задач за секунду, а три worker-и працюють незалежно, теоретична пропускна здатність може зрости до 30 задач за секунду. На практиці результат залежить від бази даних, зовнішніх сервісів та інших обмежень.
Важливо відрізняти отримання повідомлення від завершення його обробки.
Небезпечний порядок:
consumer отримав повідомлення;
queue одразу видалила його;
consumer завершився через помилку;
повідомлення втрачено.
Безпечніший порядок:
consumer отримав повідомлення;
виконав задачу;
переконався, що задача завершилася успішно;
надіслав ack;
queue видалила повідомлення.
Queue ── повідомлення ──► Consumer
│
├─ помилка → повідомлення залишається
│
└─ успіх ──► ACK → повідомлення видаляєтьсяТакий підхід часто забезпечує семантику at-least-once delivery — повідомлення буде доставлено щонайменше один раз.
Однак у разі збою після виконання задачі, але до надсилання ack, повідомлення може бути доставлено повторно. Тому consumer повинен бути ідемпотентним.
Ідемпотентна операція дає той самий кінцевий результат, навіть якщо її виконати кілька разів.
Наприклад, небезпечно просто створювати платіж під час кожної доставки повідомлення:
Повторна доставка → Повторне списання коштівКраще використовувати унікальний ідентифікатор операції та перевіряти, чи вона вже виконувалася:
messageId = payment-123
Якщо payment-123 вже оброблено:
не виконувати операцію повторно
Інакше:
виконати операцію
зберегти payment-123 як обробленийІдемпотентність особливо важлива для:
платежів;
надсилання повідомлень;
зміни стану замовлення;
створення документів;
взаємодії із зовнішніми API.
Тимчасові помилки не завжди означають, що повідомлення неправильне.
Наприклад:
база даних тимчасово недоступна;
зовнішній API повернув помилку;
виникло короткочасне мережеве переривання.
Для таких помилок використовують повторні спроби. Бажано збільшувати інтервал між ними:
Спроба 1: одразу
Спроба 2: через 1 секунду
Спроба 3: через 5 секунд
Спроба 4: через 30 секундЦей підхід називається експоненційною затримкою, або exponential backoff.
Не всі помилки потрібно повторювати. Якщо повідомлення має неправильний формат, нова спроба не допоможе. Після обмеженої кількості невдалих спроб таке повідомлення переносять у dead-letter queue, або DLQ.
DLQ потрібна для:
аналізу помилок;
ручного виправлення даних;
повторного запуску після виправлення проблеми;
запобігання нескінченному циклу повторних спроб.
Backpressure — це механізм, за допомогою якого система не дозволяє producer-ам безмежно збільшувати навантаження на consumer-ів.
Якщо consumer-и працюють повільніше за producer-а, довжина черги зростає. Це сигнал, що система не встигає обробляти вхідний потік.
Можливі стратегії:
обмежити кількість повідомлень у черзі;
тимчасово відхиляти нові задачі;
зменшити швидкість producer-а;
збільшити кількість consumer-ів;
обмежити кількість повідомлень, які один consumer обробляє одночасно;
застосувати пріоритети для важливих задач.
Без backpressure черга може спожити всю доступну пам’ять або дисковий простір.
Нижче наведено навчальний приклад черги в пам’яті. Він демонструє:
додавання задач;
розподіл задач між worker-ами;
підтвердження успішної обробки;
повторну спробу після тимчасової помилки.
Такий варіант не підходить для production, оскільки всі повідомлення зникнуть після завершення процесу. Для production використовують надійний брокер повідомлень або інше постійне сховище.
class MessageQueue {
constructor() {
this.messages = [];
this.waitingConsumers = [];
}
publish(message) {
const consumer = this.waitingConsumers.shift();
if (consumer) {
consumer(message);
return;
}
this.messages.push(message);
}
consume() {
return new Promise((resolve) => {
const message = this.messages.shift();
if (message) {
resolve(message);
return;
}
this.waitingConsumers.push(resolve);
});
}
acknowledge(messageId) {
console.log(`Повідомлення ${messageId} успішно підтверджено`);
}
}
const queue = new MessageQueue();
async function processMessage(workerName, message) {
console.log(`${workerName} обробляє ${message.id}`);
// Імітуємо тимчасову помилку для першої спроби
if (message.failOnce && message.attempts === 0) {
message.attempts += 1;
throw new Error("Тимчасова помилка зовнішнього сервісу");
}
await new Promise((resolve) => setTimeout(resolve, 300));
queue.acknowledge(message.id);
console.log(`${workerName} завершив ${message.id}`);
}
async function startWorker(workerName) {
while (true) {
const message = await queue.consume();
try {
await processMessage(workerName, message);
} catch (error) {
console.log(`${workerName}: ${error.message}`);
if (message.attempts < 3) {
console.log(`Повторна спроба для ${message.id}`);
queue.publish(message);
} else {
console.log(`Повідомлення ${message.id} перенесено в DLQ`);
}
}
}
}
startWorker("Worker 1");
startWorker("Worker 2");
for (let index = 1; index <= 5; index += 1) {
queue.publish({
id: `task-${index}`,
attempts: 0,
failOnce: index === 3
});
}У прикладі два worker-и конкурують за повідомлення. Задача task-3 навмисно завершується помилкою під час першої спроби, після чого повертається в чергу.
У реальній системі потрібно також продумати:
збереження повідомлень після перезапуску;
атомарність між виконанням задачі та ack;
обмеження кількості повторних спроб;
моніторинг довжини черги;
окрему обробку нерозпізнаних повідомлень.
У простій черзі повідомлення можуть оброблятися в порядку надходження. Але після додавання кількох worker-ів порядок завершення вже не гарантований:
Надходження: A → B → C
Завершення: B → A → CЦе важливо, якщо повідомлення змінюють одну й ту саму сутність.
Наприклад:
1. Змінити статус замовлення на "оплачено"
2. Змінити статус замовлення на "відправлено"Якщо друга задача завершиться раніше за першу, стан може стати некоректним.
Коли порядок критично важливий, потрібно використовувати відповідне розподілення повідомлень, наприклад за ідентифікатором замовлення, або реалізувати перевірку версії чи послідовності подій на стороні consumer-а.
Збільшення кількості worker-ів часто покращує пропускну здатність, але може зруйнувати глобальний порядок обробки.
Черга є частиною критичного шляху системи, тому її потрібно спостерігати.
Корисні метрики:
поточна кількість повідомлень;
швидкість надходження;
швидкість обробки;
середній час очікування;
кількість невдалих спроб;
кількість повідомлень у DLQ;
вік найстарішого повідомлення;
кількість активних worker-ів.
Особливо важливий вік найстарішого повідомлення. Навіть якщо довжина черги невелика, старі повідомлення можуть свідчити про те, що consumer-и зависли або працюють занадто повільно.
Producer створює повідомлення
↓
Queue зберігає повідомлення
↓
Consumer отримує повідомлення
↓
Consumer обробляє задачу
↓
Успіх ─────────────► ACK і видалення
│
└─ Помилка ─► повторна спроба
│
└─ ліміт спроб перевищено → DLQЯкщо повідомлення видаляється одразу після отримання, збій consumer-а призведе до втрати задачі.
ack потрібно надсилати після успішної обробки.
Черга може доставити те саме повідомлення повторно. Якщо операція не є ідемпотентною, повторна доставка може створити дубль або виконати небезпечну дію двічі.
Повідомлення з постійною помилкою може нескінченно повертатися в чергу та блокувати систему.
Потрібно встановлювати максимальну кількість спроб і використовувати DLQ.
Черга не повинна бути сховищем великих файлів або повних об’єктів домену. Зазвичай повідомлення містить ідентифікатор, тип задачі та мінімально необхідні дані.
Великі дані краще зберігати окремо, а в повідомленні передавати посилання або ідентифікатор.
Якщо producer генерує повідомлення швидше, ніж consumer їх обробляє, черга зростатиме без обмежень.
Потрібні ліміти, моніторинг і стратегія backpressure.
Наявність черги не гарантує, що повідомлення завершаться в тому самому порядку, в якому вони надійшли. Особливо це стосується кількох worker-ів і повторних спроб.
Черга повідомлень відокремлює producer-ів від consumer-ів.
Вона дає змогу виконувати задачі асинхронно та буферизувати нерівномірне навантаження.
Кілька worker-ів можуть розподіляти задачі між собою та збільшувати пропускну здатність.
Повідомлення потрібно підтверджувати після успішної обробки.
Через повторні доставки consumer-и мають бути ідемпотентними.
Тимчасові помилки можна обробляти повторними спробами з паузами.
Повідомлення, які не вдалося обробити після встановленого ліміту, варто переносити в DLQ.
Потрібно контролювати довжину черги, час очікування та кількість помилок.
Додавання worker-ів може порушити порядок обробки повідомлень.