Пошук уроків, статей та іншого контенту
З’ясуєте, як поширювати сигнал перевантаження між компонентами та запобігати переповненню черг і буферів.
Backpressure — це механізм, за допомогою якого повільний компонент повідомляє джерело даних, що більше не може безпечно приймати роботу з поточною швидкістю.
Без backpressure система часто працює так:
продюсер створює повідомлення;
брокер або проміжний сервіс приймає їх без обмежень;
черга зростає;
збільшується використання пам’яті та диска;
зростає затримка обробки;
компоненти починають завершуватися з помилками або ставати недоступними.
Backpressure не збільшує пропускну здатність системи. Він допомагає узгодити швидкість надходження роботи зі швидкістю її обробки.
Основне правило: якщо споживач не встигає, продюсер повинен або сповільнитися, або відмовитися від частини роботи.
У розподіленій системі вузьке місце може бути на кожному етапі:
HTTP-сервіс приймає більше запитів, ніж може обробити;
сервіс запису повільніше зберігає дані, ніж інший сервіс їх створює;
брокер повідомлень накопичує непрочитані повідомлення;
база даних вичерпує пул з’єднань;
зовнішній API обмежує кількість запитів;
один із воркерів обробляє завдання довше за інших.
Важливо поширити сигнал перевантаження не лише на безпосереднього клієнта. Якщо сервіс B не встигає обробляти запити від сервісу A, а A продовжує приймати запити від користувачів без обмежень, то проблема лише переміститься в чергу сервісу A.
У найпростішій моделі є три компоненти:
Продюсер → Черга → СпоживачСпоживач обробляє дані зі швидкістю C, а продюсер створює їх зі швидкістю P.
якщо P < C, черга зазвичай порожня або невелика;
якщо P = C, система працює на межі;
якщо P > C, черга зростає.
Черга не усуває проблему, а лише відкладає її. Тому черга повинна мати обмеження:
максимальну кількість повідомлень;
максимальний обсяг у байтах;
максимальний час очікування;
політику поведінки після досягнення ліміту.
Після заповнення черги система може:
призупинити продюсер;
відхилити нове повідомлення;
відкласти повідомлення для повторної спроби;
видалити найстаріші або найновіші повідомлення;
передати повідомлення до окремої черги помилок;
знизити якість або деталізацію роботи.
Вибір залежить від того, чи можна втрачати дані.
Продюсер не додає нове завдання, поки в черзі немає місця.
Це найпряміший варіант:
черга заповнена → продюсер чекає → споживач звільняє місце → продюсер продовжуєПеревага — дані не втрачаються. Недолік — очікування може поширитися на весь ланцюжок викликів.
Якщо HTTP-сервіс перевантажений, він може повернути:
429 Too Many Requests — клієнт перевищив допустиму швидкість;
503 Service Unavailable — сервіс тимчасово не може приймати роботу.
Відповідь може містити Retry-After, щоб клієнт не повторював запит негайно.
Клієнт повинен:
розпізнати тимчасову відмову;
дочекатися вказаного часу або застосувати backoff;
обмежити кількість повторних спроб;
припинити повтори після вичерпання дедлайну.
Без обмеження повторних спроб відмова може перетворитися на ще більше навантаження.
Споживач заздалегідь повідомляє, скільки повідомлень готовий прийняти. Продюсер може надсилати дані лише в межах отриманого кредиту.
Наприклад:
споживач: готовий прийняти 100 повідомлень
продюсер: надсилає 100 повідомлень
споживач: обробив 40 і повертає 40 кредитів
продюсер: може надіслати ще 40Цей підхід добре підходить для потокової обробки та протоколів, де передавання даних відбувається частинами.
Компонент може приймати не більше заданої кількості операцій за секунду. Якщо ліміт перевищено, нові операції очікують або відхиляються.
Обмеження швидкості може бути:
локальним — для одного процесу;
на рівні сервісу;
спільним для всіх екземплярів сервісу;
прив’язаним до клієнта, користувача або типу операції.
Локальний ліміт простіший, але не враховує загальне навантаження на кластер. Для розподіленого ліміту потрібен спільний стан або компонент, який координує дозволи.
Нижче наведено runnable-приклад на JavaScript. Продюсер створює завдання швидше, ніж воркер їх обробляє. Коли черга заповнюється, продюсер очікує сигналу про звільнення місця.
class BoundedQueue {
constructor(limit) {
this.limit = limit;
this.items = [];
this.waitingProducers = [];
this.waitingConsumers = [];
}
async push(item) {
while (this.items.length >= this.limit) {
await new Promise((resolve) => {
this.waitingProducers.push(resolve);
});
}
this.items.push(item);
const resolveConsumer = this.waitingConsumers.shift();
if (resolveConsumer) {
resolveConsumer();
}
}
async pop() {
while (this.items.length === 0) {
await new Promise((resolve) => {
this.waitingConsumers.push(resolve);
});
}
const item = this.items.shift();
const resolveProducer = this.waitingProducers.shift();
if (resolveProducer) {
resolveProducer();
}
return item;
}
size() {
return this.items.length;
}
}
const queue = new BoundedQueue(3);
function sleep(milliseconds) {
return new Promise((resolve) => {
setTimeout(resolve, milliseconds);
});
}
async function producer() {
for (let id = 1; id <= 10; id += 1) {
console.log(`Продюсер готує завдання ${id}`);
await queue.push({ id });
console.log(
`Завдання ${id} додано до черги; розмір черги: ${queue.size()}`
);
// Продюсер створює завдання швидше, ніж воркер їх обробляє.
await sleep(100);
}
}
async function worker() {
for (let i = 0; i < 10; i += 1) {
const task = await queue.pop();
console.log(
`Воркер почав завдання ${task.id}; розмір черги: ${queue.size()}`
);
// Імітація повільної обробки.
await sleep(500);
console.log(`Воркер завершив завдання ${task.id}`);
}
}
async function main() {
await Promise.all([producer(), worker()]);
console.log("Усі завдання оброблено");
}
main().catch((error) => {
console.error(error);
process.exitCode = 1;
});Ключовий момент — метод push. Він не додає елемент понад ліміт, а призупиняє продюсер до моменту, коли воркер забере завдання.
У розподіленій системі таке очікування зазвичай реалізується не спільною пам’яттю, а протоколом:
відповіддю 429 або 503;
підтвердженням прийому повідомлення;
кредитами;
відсутністю дозволу на наступне читання;
затримкою перед наступною публікацією.
Розглянемо ланцюжок:
Клієнт → API → Сервіс замовлень → Сервіс платежівЯкщо сервіс платежів перевантажений, сервіс замовлень не повинен безмежно накопичувати запити. Він може:
обмежити кількість одночасних викликів платежів;
обмежити власну внутрішню чергу;
повернути клієнту тимчасову помилку;
зберегти операцію в надійну чергу, якщо бізнес-процес допускає асинхронність;
припинити приймання нових операцій, коли досягнуто безпечного ліміту.
Сигнал повинен рухатися проти напрямку потоку даних:
сервіс платежів → сервіс замовлень → API → клієнтЯкщо сигнал зупиняється на одному рівні, попередні компоненти продовжують створювати навантаження.
Асинхронна черга корисна, коли клієнту не потрібно чекати завершення всієї операції. Наприклад, API може:
перевірити базову коректність запиту;
додати завдання до обмеженої черги;
повернути ідентифікатор операції;
обробити завдання пізніше.
Але навіть у цьому випадку потрібні обмеження:
черга не повинна бути необмеженою;
продюсер має отримувати інформацію про відмову;
потрібно розрізняти тимчасову недоступність і некоректні дані;
повторна публікація не повинна створювати дублікати;
необхідно визначити, що робити із завданнями, які довго очікують.
Якщо черга переповнена, безпечніше відмовити новому запиту, ніж прийняти його й мовчки втратити пізніше.
Підходить для операцій, які можна повторити з боку клієнта. Відмова повинна бути явною та спостережуваною.
Підходить, якщо старі повідомлення важливіші, а нові можна не приймати. Наприклад, для послідовної обробки подій це може бути неприйнятно.
Підходить для даних про поточний стан, коли новіше значення робить старе непотрібним. Для фінансових операцій або команд зміни стану такий підхід зазвичай небезпечний.
Система може тимчасово:
зменшити частоту оновлень;
відкладати необов’язкові операції;
вимкнути другорядні обчислення;
обробляти спрощене представлення даних.
Це дозволяє зберегти основний функціонал під час пікового навантаження.
Коли сервіс отримує сигнал перевантаження, клієнти часто повторюють запит пізніше. Зазвичай застосовують експоненційний backoff:
100 мс → 200 мс → 400 мс → 800 мс → ...До backoff додають випадкове відхилення — jitter. Воно не дає великій кількості клієнтів повторити запити одночасно.
Повторні спроби повинні мати:
максимальну кількість;
загальний дедлайн;
обмеження часу очікування;
коректне визначення тимчасових помилок;
захист від повторення неідемпотентних операцій.
Якщо операцію повторити небезпечно, клієнт і сервер мають використовувати ідентифікатор запиту або інший механізм дедуплікації.
Одного показника завантаження CPU недостатньо. Корисно відстежувати:
поточний і максимальний розмір черги;
час очікування повідомлення в черзі;
швидкість надходження та обробки;
кількість відхилених повідомлень;
кількість активних операцій;
кількість відповідей 429 і 503;
кількість повторних спроб;
час обробки;
частку повідомлень, які завершилися помилкою.
Особливо важливий вік найстарішого повідомлення. Черга може мати прийнятний розмір, але якщо повідомлення чекають надто довго, система вже не виконує вимоги до затримки.
Для кожного компонента визначте:
Яка максимальна кількість одночасних операцій безпечна?
Який розмір черги допустимий?
Що відбувається після заповнення черги?
Як попередній компонент дізнається про перевантаження?
Чи можна повторити відхилену операцію?
Чи можна втрачати повідомлення?
Який дедлайн очікування?
Як оператор побачить, що backpressure активувався?
Після цього перевірте весь ланцюжок. Локальний backpressure не захищає систему, якщо наступний або попередній компонент має необмежений буфер.
Черга приховує проблему, але не вирішує її. Рано чи пізно вона споживе доступну пам’ять, дисковий простір або ліміт брокера.
Клієнти, які одразу повторюють відхилені запити, створюють retry storm — хвилю повторних запитів, що посилює перевантаження.
Якщо воркер зупинив читання з черги, але API продовжує приймати запити в необмежену локальну чергу, проблема лише переміщується.
Очікування місця в черзі без тайм-ауту може залишити запит активним назавжди й вичерпати ресурси сервісу.
Для критичних команд, аналітичних подій і даних про поточний стан можуть бути потрібні різні правила переповнення. Втрата події платежу та втрата незначущого телеметричного вимірювання — різні за наслідками.
Якщо система лише сповільнюється, але не реєструє розмір черг, час очікування та кількість відмов, причину деградації буде складно знайти.
Backpressure узгоджує швидкість продюсера зі швидкістю споживача.
Кожна черга та внутрішній буфер повинні мати обмеження.
Сигнал перевантаження можна передавати очікуванням, відмовою, кредитами або обмеженням швидкості.
Відмова має поширюватися проти напрямку потоку даних до джерела навантаження.
429 і 503 повинні оброблятися з backoff, jitter, дедлайном і лімітом повторів.
Політика переповнення залежить від того, чи можна втрачати дані.
Розмір черги, вік повідомлень, час очікування та кількість відмов необхідно вимірювати.
Обмежена черга краще за необмежене накопичення, оскільки робить перевантаження контрольованим і видимим.