Пошук уроків, статей та іншого контенту
Поєднуйте читання й запис у Duplex Streams та змінюйте дані за допомогою Transform Streams.
Duplex Stream — це потік, який одночасно має дві незалежні сторони:
Writable — приймає дані через write();
Readable — віддає дані через read(), подію data або pipe().
На відміну від звичайного Readable або Writable, Duplex може читати й записувати дані.
Прикладом Duplex у Node.js є мережеве TCP-з’єднання: ми можемо записувати дані в сокет і паралельно читати відповідь із нього.
Важливо: дві сторони Duplex не обов’язково безпосередньо пов’язані між собою. Записані дані можуть:
одразу перетворюватися на результат читання;
передаватися в зовнішню систему;
зберігатися в буфері;
взагалі не з’являтися на readable-стороні.
Створюючи власний Duplex, зазвичай реалізують два методи:
_write(chunk, encoding, callback) — викликається для кожної порції вхідних даних;
_read(size) — викликається, коли споживач готовий отримати нові дані.
_write()callback()callback(error);Метод _read() може додавати дані у readable-сторону за допомогою:
this.push(data);Коли дані завершилися, потрібно передати null:
this.push(null);Спрощена структура Duplex:
const { Duplex } = require('node:stream');
class ExampleDuplex extends Duplex {
_read(size) {
// Тут потік може додати нові дані для читання
}
_write(chunk, encoding, callback) {
// Тут потік обробляє дані, які в нього записали
callback();
}
}Реалізація _read() і _write() залежить від задачі. Node.js не знає, як саме ваш потік має отримувати, зберігати або генерувати дані.
Transform — це спеціальний різновид Duplex.
Він одночасно:
приймає дані як Writable;
змінює їх;
віддає результат як Readable.
Схема роботи:
вхідні дані → Transform Stream → змінені даніДля Transform не потрібно окремо реалізовувати _read(). Замість цього реалізують метод _transform():
const { Transform } = require('node:stream');
class MyTransform extends Transform {
_transform(chunk, encoding, callback) {
const result = /* перетворення chunk */;
callback(null, result);
}
}Перший аргумент callback — помилка, другий — результат перетворення:
callback(null, transformedData);Якщо результату для цієї порції немає, можна передати лише успішний callback:
callback();Створимо Transform Stream, який перетворює весь вхідний текст на верхній регістр.
Файл uppercase.js:
const { Transform, pipeline } = require('node:stream');
class UppercaseTransform extends Transform {
_transform(chunk, encoding, callback) {
try {
const text = chunk.toString('utf8');
const result = text.toUpperCase();
callback(null, result);
} catch (error) {
callback(error);
}
}
}
const uppercase = new UppercaseTransform();
pipeline(
process.stdin,
uppercase,
process.stdout,
(error) => {
if (error) {
console.error('Помилка обробки потоку:', error.message);
process.exitCode = 1;
}
}
);Запустити програму можна так:
node uppercase.jsПісля цього введіть текст у терміналі та натисніть Enter. Щоб завершити введення в Unix-подібних системах, натисніть Ctrl+D, а у Windows — Ctrl+Z, потім Enter.
Також можна передати файл через стандартний ввід:
node uppercase.js < input.txtРезультат буде виведено у стандартний вивід. Наприклад, якщо input.txt містить:
Node.js streams are usefulпрограма виведе:
NODE.JS STREAMS ARE USEFULStream не гарантує, що один виклик _transform() отримає цілий рядок, повідомлення або файл. Дані можуть надходити частинами:
"Node"
".js "
"streams"або навіть посередині слова:
"No"
"de.j"
"s"Тому Transform Stream має коректно працювати з кожним окремим chunk.
Для перетворення незалежних частин, наприклад зміни регістру, цього достатньо. Але для операцій над цілими рядками або повідомленнями може знадобитися власний буфер:
const { Transform } = require('node:stream');
class LineTransform extends Transform {
constructor() {
super();
this.buffer = '';
}
_transform(chunk, encoding, callback) {
this.buffer += chunk.toString('utf8');
const lines = this.buffer.split('\n');
// Останній елемент може бути неповним рядком
this.buffer = lines.pop();
for (const line of lines) {
this.push(`${line.toUpperCase()}\n`);
}
callback();
}
_flush(callback) {
if (this.buffer.length > 0) {
this.push(this.buffer.toUpperCase());
}
callback();
}
}Метод _flush() викликається перед завершенням Transform Stream. Він дає змогу обробити дані, які залишилися у внутрішньому буфері.
pipe()Метод pipe() з’єднує readable-сторону одного потоку з writable-стороною іншого:
readable.pipe(transform).pipe(writable);У прикладі з верхнім регістром:
process.stdin
.pipe(new UppercaseTransform())
.pipe(process.stdout);Тут:
process.stdin — Readable Stream;
UppercaseTransform — одночасно Writable і Readable;
process.stdout — Writable Stream.
Під час використання pipe() Node.js автоматично передає дані частинами та враховує backpressure — ситуацію, коли наступний потік не встигає обробляти дані.
pipeline()Ланцюжок pipe() зручний, але помилки в довгих ланцюжках можуть бути складнішими для обробки. Функція pipeline() допомагає обробити завершення всього ланцюжка:
const { pipeline } = require('node:stream');
pipeline(
source,
transform,
destination,
(error) => {
if (error) {
console.error(error);
} else {
console.log('Обробку завершено');
}
}
);Якщо один із потоків завершується з помилкою, pipeline() повідомляє про це у callback і коректно завершує пов’язаний ланцюжок.
Transform Stream може отримувати дані швидше, ніж встигає їх обробляти. Node.js у такому випадку використовує backpressure:
Writable-сторона приймає дані;
Transform обробляє їх;
Readable-сторона передає результат далі;
якщо наступний потік переповнений, передавання тимчасово сповільнюється.
У власній реалізації Transform не потрібно вручну керувати всією цією логікою. Важливо:
не викликати callback() багато разів;
не викликати його до завершення асинхронної операції;
не накопичувати необмежену кількість даних у власних масивах або буферах;
використовувати this.push() для передачі результату на readable-сторону.
Приклад асинхронного перетворення:
const { Transform } = require('node:stream');
class DelayedTransform extends Transform {
_transform(chunk, encoding, callback) {
setTimeout(() => {
const result = chunk.toString('utf8').toUpperCase();
// callback викликається після завершення асинхронної операції
callback(null, result);
}, 100);
}
}Поки callback не викликано, потік вважає поточну порцію даних необробленою.
Один chunk не обов’язково відповідає одному повідомленню або рядку. Якщо формат даних має межі, наприклад \n, їх потрібно шукати з урахуванням даних із попередніх chunk.
Якщо _transform() не викликає callback, потік зупиниться на поточній порції:
_transform(chunk, encoding, callback) {
const result = chunk.toString().toUpperCase();
// Без callback наступні дані можуть не оброблятися
callback(null, result);
}Не можна викликати callback і в try, і в catch після того, як він уже був викликаний. Кожен виклик _transform() має завершитися рівно одним викликом callback.
Помилки можуть виникати під час читання, перетворення або запису. Для з’єднання кількох потоків краще використовувати pipeline() і перевіряти його callback.
Transform Stream призначений для потокової обробки. Не варто без потреби збирати весь файл у пам’яті перед перетворенням.
Duplex Stream має незалежні readable- і writable-сторони.
Transform — це спеціальний Duplex, який перетворює вхідні дані на вихідні.
Для власного Transform Stream реалізують _transform(chunk, encoding, callback).
Результат передають через callback(null, result) або this.push(result).
Дані надходять частинами, тому один chunk не слід сприймати як повне повідомлення.
Метод _flush() використовується для обробки даних, що залишилися у внутрішньому буфері.
pipeline() спрощує з’єднання потоків і обробку помилок.