Пошук уроків, статей та іншого контенту
Передавайте дані між потоками через pipe та керуйте ланцюжками потоків за допомогою pipeline.
pipe та pipelineУ Node.js потоки часто потрібно з’єднувати між собою:
читати дані з файлу;
обробляти їх;
записувати результат в інший файл;
передавати дані через кілька послідовних перетворень.
Для цього використовують:
stream.pipe() — з’єднання двох потоків;
stream.pipeline() — керування повним ланцюжком потоків із централізованою обробкою помилок і завершенням.
Типовий ланцюжок має такий вигляд:
Readable → Transform → WritableНаприклад:
файл → стиснення gzip → інший файлpipeМетод pipe() передає дані з потоку для читання (Readable) у потік для запису (Writable або Transform).
readable.pipe(writable);Після виклику:
Node.js слухає події читального потоку;
отримані дані передаються в потік для запису;
коли читання завершується, потік для запису за замовчуванням також завершується;
механізм backpressure не дає джерелу надсилати дані швидше, ніж їх можна обробити.
const fs = require('node:fs');
const source = fs.createReadStream('input.txt');
const destination = fs.createWriteStream('output.txt');
source.pipe(destination);
destination.on('finish', () => {
console.log('Копіювання завершено');
});У цьому прикладі:
createReadStream() створює потік для читання;
createWriteStream() створює потік для запису;
pipe() передає вміст input.txt у output.txt;
подія finish означає, що потік запису завершив роботу.
Файл output.txt буде створено автоматично, якщо його ще немає. Якщо він існує, його вміст буде перезаписано.
pipeМетод pipe() повертає потік призначення. Завдяки цьому виклики можна об’єднувати в ланцюжок.
Наприклад, текст можна прочитати, перетворити на великі літери та записати в інший файл:
const fs = require('node:fs');
const { Transform } = require('node:stream');
const uppercase = new Transform({
transform(chunk, encoding, callback) {
const result = chunk.toString().toUpperCase();
// Передаємо перетворені дані далі
callback(null, result);
}
});
fs.createReadStream('input.txt')
.pipe(uppercase)
.pipe(fs.createWriteStream('output.txt'));Transform одночасно є:
потоком для читання результату;
потоком для запису вхідних даних.
Метод transform() викликається для кожної порції даних. У прикладі вона перетворюється на рядок у верхньому регістрі.
endЗа замовчуванням після завершення джерела pipe() завершує потік призначення. Це можна змінити параметром end:
source.pipe(destination, { end: false });У такому разі завершення source не закриє destination. Це корисно, коли потрібно послідовно передати кілька джерел в один потік запису.
const fs = require('node:fs');
const first = fs.createReadStream('first.txt');
const second = fs.createReadStream('second.txt');
const output = fs.createWriteStream('combined.txt');
first.pipe(output, { end: false });
first.on('end', () => {
second.pipe(output);
});У цьому прикладі обидва файли записуються в combined.txt послідовно.
Потоки можуть працювати з різною швидкістю:
диск може швидко читати великі обсяги даних;
перетворення може обробляти їх повільніше;
мережеве з’єднання може тимчасово приймати лише невеликі порції.
Якщо безконтрольно передавати всі дані одразу, вони накопичуватимуться в пам’яті. Це називається проблемою backpressure.
pipe() автоматично враховує backpressure:
якщо потік призначення перевантажений, читання тимчасово призупиняється;
коли потік призначення знову готовий приймати дані, читання продовжується.
Саме тому для передавання великих обсягів даних краще використовувати потоки, а не повністю читати файл через fs.readFile().
pipepipe() зручний для простих ланцюжків, але його потрібно використовувати обережно під час обробки помилок.
Наприклад:
source
.pipe(transform)
.pipe(destination);Якщо один із потоків завершується з помилкою, потрібно правильно обробити цю помилку. Окреме додавання обробників до всіх потоків може бути громіздким і легко призвести до неповного очищення ресурсів.
Для складніших ланцюжків варто використовувати pipeline().
pipelinepipeline() з’єднує кілька потоків і централізовано керує їхнім життєвим циклом.
Він:
передає дані між усіма потоками;
враховує backpressure;
викликає callback після успішного завершення;
передає помилку в callback;
завершує або знищує пов’язані потоки, якщо один із них завершується з помилкою.
Імпортувати pipeline можна з модуля node:stream.
const { pipeline } = require('node:stream');const fs = require('node:fs');
const { Transform, pipeline } = require('node:stream');
const uppercase = new Transform({
transform(chunk, encoding, callback) {
const result = chunk.toString().toUpperCase();
// Передаємо результат наступному потоку
callback(null, result);
}
});
pipeline(
fs.createReadStream('input.txt'),
uppercase,
fs.createWriteStream('output.txt'),
(error) => {
if (error) {
console.error('Помилка під час обробки:', error.message);
return;
}
console.log('Обробку завершено');
}
);На відміну від ланцюжка з pipe(), у pipeline() є одна централізована точка обробки помилки.
Модуль node:zlib надає потік createGzip(), який стискає дані у форматі gzip.
const fs = require('node:fs');
const { createGzip } = require('node:zlib');
const { pipeline } = require('node:stream');
pipeline(
fs.createReadStream('input.txt'),
createGzip(),
fs.createWriteStream('input.txt.gz'),
(error) => {
if (error) {
console.error('Не вдалося стиснути файл:', error.message);
process.exitCode = 1;
return;
}
console.log('Файл стиснено в input.txt.gz');
}
);Порядок потоків тут такий:
input.txt → gzip → input.txt.gzЦей приклад можна зберегти у файл compress.js і запустити командою:
node compress.jsПоруч зі скриптом має бути файл input.txt.
pipeline з промісомУ сучасному Node.js можна використовувати проміс-версію pipeline. Вона доступна в node:stream/promises.
const fs = require('node:fs');
const { createGzip } = require('node:zlib');
const { pipeline } = require('node:stream/promises');
async function compressFile() {
await pipeline(
fs.createReadStream('input.txt'),
createGzip(),
fs.createWriteStream('input.txt.gz')
);
console.log('Файл стиснено');
}
compressFile().catch((error) => {
console.error('Помилка:', error.message);
process.exitCode = 1;
});Якщо будь-який потік завершується з помилкою, pipeline() відхиляє проміс. Помилку можна обробити через try...catch:
const fs = require('node:fs');
const { createGzip } = require('node:zlib');
const { pipeline } = require('node:stream/promises');
async function compressFile() {
try {
await pipeline(
fs.createReadStream('input.txt'),
createGzip(),
fs.createWriteStream('input.txt.gz')
);
console.log('Файл стиснено');
} catch (error) {
console.error('Помилка під час стиснення:', error.message);
process.exitCode = 1;
}
}
compressFile();Проміс-версія зручна в асинхронних функціях, де інші операції також виконуються через await.
pipe та pipelinepipeПідходить, коли:
ланцюжок простий;
життєвий цикл потоків добре контролюється;
потрібен короткий і зрозумілий запис.
source.pipe(transform).pipe(destination);pipelineПідходить, коли:
у ланцюжку кілька потоків;
важливо централізовано обробляти помилки;
потрібно коректно завершувати всі потоки;
код використовує async/await.
pipeline(source, transform, destination, callback);У прикладному коді для надійних ланцюжків зазвичай варто надавати перевагу pipeline().
Погано:
const fs = require('node:fs');
fs.createReadStream('missing.txt')
.pipe(fs.createWriteStream('output.txt'));Якщо файл missing.txt не існує, читальний потік видасть помилку. Надійніше використати pipeline():
const fs = require('node:fs');
const { pipeline } = require('node:stream');
pipeline(
fs.createReadStream('missing.txt'),
fs.createWriteStream('output.txt'),
(error) => {
if (error) {
console.error('Не вдалося скопіювати файл:', error.message);
}
}
);readFileДля великих файлів не варто спочатку завантажувати весь вміст у пам’ять:
const fs = require('node:fs/promises');
const content = await fs.readFile('large-file.txt');Потокове копіювання або обробка використовує пам’ять ефективніше:
const fs = require('node:fs');
const { pipeline } = require('node:stream');
pipeline(
fs.createReadStream('large-file.txt'),
fs.createWriteStream('copy.txt'),
(error) => {
if (error) {
console.error(error.message);
}
}
);pipe та pipeline без потребиМожна використовувати pipe() всередині ланцюжка, але це ускладнює керування завершенням і помилками. Якщо весь ланцюжок відомий заздалегідь, краще передати всі потоки безпосередньо в pipeline():
pipeline(
source,
firstTransform,
secondTransform,
destination,
callback
);pipe() з’єднує потік для читання з потоком призначення.
Виклики pipe() можна об’єднувати в ланцюжок.
pipe() автоматично підтримує backpressure.
Параметр { end: false } забороняє автоматично завершувати потік призначення.
pipeline() з’єднує потоки та централізовано обробляє помилки.
pipeline() також завершує пов’язані потоки, якщо в ланцюжку виникає помилка.
Для callback-стилю використовується node:stream.
Для async/await доступна проміс-версія з node:stream/promises.
Для надійних ланцюжків із кількох потоків зазвичай краще використовувати pipeline().