Пошук уроків, статей та іншого контенту
Застосовуйте backpressure, щоб узгодити швидкість читання й запису та уникати переповнення пам’яті.
Backpressure — це механізм узгодження швидкості між частинами потоку, коли джерело даних працює швидше за споживача.
Наприклад:
файл читається швидко;
мережеве з’єднання приймає дані повільно;
база даних обробляє записи повільніше, ніж вони надходять.
Без backpressure швидкий producer продовжує створювати дані, а повільний consumer накопичує їх у буферах. У результаті зростає використання пам’яті, а в гіршому випадку процес завершується з помилкою JavaScript heap out of memory.
У Node.js backpressure вбудований у Streams API.
WritableМетод writable.write(chunk) повертає логічне значення:
true — внутрішній буфер ще може приймати дані;
false — буфер досягнув порогу highWaterMark, тому потрібно призупинити запис.
Якщо write() повернув false, не слід одразу записувати наступні частини даних. Потрібно дочекатися події drain:
import { once } from 'node:events';
import { Writable } from 'node:stream';
const output = new Writable({
highWaterMark: 2,
write(chunk, encoding, callback) {
setTimeout(() => {
console.log(`Оброблено: ${chunk.toString()}`);
callback();
}, 300);
},
});
for (let index = 1; index <= 10; index += 1) {
const canContinue = output.write(`Запис ${index}\n`);
if (!canContinue) {
console.log('Буфер заповнений, очікуємо drain');
await once(output, 'drain');
console.log('Буфер знову готовий до запису');
}
}
output.end();
await once(output, 'finish');
console.log('Усі дані записано');У цьому прикладі Writable обробляє кожен запис із затримкою. Коли внутрішній буфер заповнюється, write() повертає false. Цикл призупиняється до появи drain.
write()Такий код може неконтрольовано збільшувати використання пам’яті:
for (const chunk of chunks) {
output.write(chunk);
}Якщо output не встигає обробляти дані, усі наступні chunks залишаються у внутрішній черзі. Кількість елементів у черзі може стати дуже великою.
Правильний алгоритм:
записати chunk;
перевірити результат write();
якщо отримано false, дочекатися drain;
продовжити запис після звільнення буфера.
highWaterMarkhighWaterMark визначає поріг, після якого потік сигналізує про backpressure.
Він не є жорстким максимальним розміром буфера. Потік може тимчасово містити більше даних, але highWaterMark визначає момент, коли producer має сповільнитися.
Значення за замовчуванням залежать від типу потоку та версії Node.js, тому для критичних ділянок їх краще задавати явно.
Для звичайних потоків байтів highWaterMark вимірюється в байтах:
import { Writable } from 'node:stream';
const byteStream = new Writable({
highWaterMark: 64 * 1024,
write(chunk, encoding, callback) {
callback();
},
});Для потоків у objectMode highWaterMark вимірюється в кількості об’єктів:
import { Writable } from 'node:stream';
const objectStream = new Writable({
objectMode: true,
highWaterMark: 16,
write(record, encoding, callback) {
callback();
},
});Це означає, що highWaterMark: 16 для object mode — приблизно 16 об’єктів, а не 16 байтів.
ReadableУ Readable аналогічний сигнал повертає метод push():
true — можна продовжити додавання даних;
false — внутрішній буфер читабельного потоку заповнений, потрібно припинити додавання.
Класичний producer має припиняти виклик push(), коли він повернув false. Node.js викличе _read() знову, коли споживач звільнить місце в буфері.
import { Readable } from 'node:stream';
class NumberSource extends Readable {
constructor(maximum) {
super({
objectMode: true,
highWaterMark: 2,
});
this.current = 1;
this.maximum = maximum;
}
_read() {
while (this.current <= this.maximum) {
const canContinue = this.push(this.current);
this.current += 1;
if (!canContinue) {
// Зупиняємо генерацію до наступного виклику _read
return;
}
}
this.push(null);
}
}
const source = new NumberSource(10);
source.on('data', (number) => {
console.log(number);
});
source.on('end', () => {
console.log('Джерело завершило роботу');
});Після push(false) producer не повинен продовжувати генерувати значення в цьому виклику _read(). Інакше backpressure буде проігноровано.
pipe() і автоматичний backpressureНайчастіше не потрібно вручну керувати write() і drain. Метод readable.pipe(writable) автоматично:
читає дані з Readable;
записує їх у Writable;
зупиняє читання, коли Writable перевантажений;
відновлює читання після drain;
завершує destination після завершення source.
import { Readable, Writable } from 'node:stream';
const source = Readable.from(
(async function* generateRecords() {
for (let index = 1; index <= 20; index += 1) {
yield { id: index };
}
})(),
{
objectMode: true,
highWaterMark: 2,
},
);
const destination = new Writable({
objectMode: true,
highWaterMark: 2,
write(record, encoding, callback) {
setTimeout(() => {
console.log(`Збережено запис ${record.id}`);
callback();
}, 100);
},
});
source.pipe(destination);
destination.on('finish', () => {
console.log('Потік завершено без переповнення буфера');
});У цьому прикладі генератор створює записи швидше, ніж destination їх обробляє. pipe() тимчасово призупиняє source, коли буфер destination досягає порогу.
pipe()Ручне керування потребує правильно обробляти:
write() === false;
подію drain;
завершення потоку;
помилки;
передчасне закриття одного з потоків.
pipe() бере на себе основну логіку backpressure. Для типового з’єднання Readable із Writable це безпечніший варіант.
pipeline() для обробки помилокpipe() автоматично узгоджує швидкість, але для складних ланцюжків зручніше використовувати pipeline().
pipeline():
з’єднує кілька потоків;
передає backpressure між ними;
коректно завершує всі потоки при помилці;
дозволяє дочекатися повного завершення через Promise.
import { Readable, Transform, Writable } from 'node:stream';
import { pipeline } from 'node:stream/promises';
const source = Readable.from(
(async function* generateNumbers() {
for (let number = 1; number <= 10; number += 1) {
yield number;
}
})(),
{ objectMode: true },
);
const double = new Transform({
objectMode: true,
transform(number, encoding, callback) {
callback(null, number * 2);
},
});
const destination = new Writable({
objectMode: true,
write(number, encoding, callback) {
setTimeout(() => {
console.log(`Отримано: ${number}`);
callback();
}, 200);
},
});
try {
await pipeline(source, double, destination);
console.log('Pipeline успішно завершено');
} catch (error) {
console.error('Pipeline завершено з помилкою:', error);
process.exitCode = 1;
}Запустити такий файл можна як модуль Node.js:
node stream-example.mjspipeline() завершує Promise лише після того, як усі дані пройшли через усі етапи. Якщо будь-який потік завершується з помилкою, Promise буде відхилено.
Іноді дані надходять не з Readable, а з асинхронного джерела. У такому разі backpressure можна реалізувати безпосередньо через write() і drain.
import { once } from 'node:events';
import { finished } from 'node:stream/promises';
import { Writable } from 'node:stream';
const destination = new Writable({
highWaterMark: 3,
write(chunk, encoding, callback) {
setTimeout(() => {
console.log(`Записано: ${chunk.toString()}`);
callback();
}, 250);
},
});
for (let index = 1; index <= 12; index += 1) {
const canContinue = destination.write(`chunk-${index}`);
if (!canContinue) {
// Подальший запис призупиняється, поки buffer не спорожніє
await once(destination, 'drain');
}
}
destination.end();
await finished(destination);
console.log('Передавання завершено');Важливо, що finished() очікує завершення потоку й також реагує на помилки. Для production-коду це надійніше, ніж очікувати лише подію finish.
TransformTransform одночасно є Readable і Writable. Через це він може мати backpressure з обох боків:
вхідні дані надходять від upstream;
результати передаються downstream;
якщо downstream повільний, transform не повинен безмежно накопичувати результати.
У власній реалізації _transform() потрібно викликати callback лише після завершення обробки chunk:
import { Transform } from 'node:stream';
const transform = new Transform({
transform(chunk, encoding, callback) {
setTimeout(() => {
const result = chunk.toString().toUpperCase();
// Передаємо результат далі лише після завершення обробки
callback(null, result);
}, 100);
},
});Якщо викликати callback() занадто рано, а потім продовжувати генерувати результати незалежно від стану downstream, можна створити неконтрольоване накопичення даних.
Для Transform не потрібно вручну викликати drain між його внутрішнім callback() і наступним потоком. Цю взаємодію забезпечують pipe() або pipeline().
Занадто малий highWaterMark:
зменшує використання пам’яті;
частіше призупиняє та відновлює потік;
може збільшити накладні витрати.
Занадто великий highWaterMark:
дозволяє обробляти більші порції даних;
може збільшити затримку;
споживає більше пам’яті під час повільного downstream.
Оптимальне значення залежить від:
розміру chunks;
швидкості джерела;
швидкості споживача;
кількості одночасних потоків;
доступної пам’яті.
highWaterMark не замінює backpressure. Навіть із правильно підібраним значенням producer повинен реагувати на сигнали потоку.
false від write()destination.write(chunk);
// Наступний chunk записується незалежно від стану буфераПотрібно перевіряти результат write() і чекати drain.
end() до завершення записуend() сигналізує, що нових даних більше не буде. Його слід викликати лише після того, як усі chunks передані в write() і виконано необхідне очікування drain.
data без контролюРучна обробка події data може ускладнити контроль швидкості. Для передавання між потоками перевагу варто надавати pipe() або pipeline().
Потік у objectMode працює з JavaScript-об’єктами, а звичайний потік — із Buffer, рядками або Uint8Array. Неправильне змішування режимів може спричинити помилки або неочікувану серіалізацію.
Очікування лише finish не гарантує успішне завершення всієї операції. Для ланцюжків потоків використовуйте pipeline() і обробляйте відхилення його Promise.
Збирання всіх даних у масив перед обробкою усуває переваги потоків. Backpressure працює, коли дані обробляються частинами, а не накопичуються повністю в пам’яті.
Під час реалізації потокового передавання:
Визначте, який компонент є producer, а який — consumer.
Для стандартного з’єднання використовуйте pipeline().
Якщо пишете вручну, завжди перевіряйте результат write().
На false призупиняйте producer і чекайте drain.
У власному Readable припиняйте виклик push() після false.
Використовуйте highWaterMark як контрольований поріг буферизації.
Перевіряйте поведінку під час навмисно повільного consumer.
Обробляйте помилки та коректно завершуйте всі потоки.
Backpressure узгоджує швидкість producer і consumer.
Writable.write() повертає false, коли його буфер перевантажений.
Після false потрібно чекати подію drain.
Readable.push() також повертає сигнал, після якого producer має призупинити генерацію.
pipe() автоматично реалізує основну логіку backpressure.
pipeline() додатково забезпечує коректне завершення та обробку помилок.
highWaterMark визначає поріг буферизації, але не замінює обробку сигналів потоку.
Ігнорування backpressure може призвести до надмірного використання пам’яті та аварійного завершення процесу.