Пошук уроків, статей та іншого контенту
Створюйте Writable Streams і записуйте частини даних із керуванням завершенням та помилками.
Writable Stream — це потік, у який можна послідовно записувати частини даних.
Замість того щоб накопичувати весь результат у пам’яті, програма записує його частинами:
у файл;
у мережеве з’єднання;
у базу даних;
у власний обробник даних;
у стандартний потік виведення.
Основний метод для запису — write():
writable.write(chunk);Коли запис завершено, потрібно викликати:
writable.end();Виклик end() означає: нових даних більше не буде.
Власний Writable Stream можна створити за допомогою класу Writable:
const { Writable } = require('node:stream');
const writable = new Writable({
write(chunk, encoding, callback) {
console.log('Отримано:', chunk.toString());
// Повідомляємо потік, що поточний фрагмент оброблено
callback();
}
});
writable.write('Перший фрагмент');
writable.write('Другий фрагмент');
writable.end('Останній фрагмент');Метод _write або параметр write отримує три аргументи:
chunk — фрагмент даних;
encoding — кодування, якщо фрагмент передано як рядок;
callback — функція, яку потрібно викликати після завершення обробки.
Поки callback не викликано, Writable Stream вважає, що поточний фрагмент ще обробляється.
Для складнішої логіки зручно успадкуватися від Writable і реалізувати метод _write:
const { Writable } = require('node:stream');
class ConsoleWritable extends Writable {
_write(chunk, encoding, callback) {
const text = chunk.toString();
console.log(`Запис: ${text}`);
callback();
}
}
const output = new ConsoleWritable();
output.write('Один');
output.write('Два');
output.end('Три');Метод _write викликається окремо для кожного фрагмента. Node.js не передає наступний фрагмент у _write, доки попередній виклик не завершиться через callback.
Це дає змогу працювати з асинхронними операціями:
const { Writable } = require('node:stream');
class DelayedWritable extends Writable {
_write(chunk, encoding, callback) {
setTimeout(() => {
console.log(chunk.toString());
// Завершуємо асинхронну обробку
callback();
}, 100);
}
}
const output = new DelayedWritable();
output.write('Перший');
output.write('Другий');
output.end('Третій');Фрагменти будуть оброблені послідовно, навіть якщо обробка кожного з них асинхронна.
Writable Stream має кілька важливих подій:
finish — усі дані записано, викликано _final, якщо він існує;
close — потік і пов’язані ресурси закрито;
error — під час запису сталася помилка;
drain — внутрішній буфер знову готовий приймати дані.
Подію finish можна використати для виконання дій після завершення запису:
const { Writable } = require('node:stream');
const output = new Writable({
write(chunk, encoding, callback) {
console.log(chunk.toString());
callback();
}
});
output.on('finish', () => {
console.log('Усі дані записано');
});
output.write('Рядок 1');
output.write('Рядок 2');
output.end();Метод end() також може отримати останній фрагмент:
output.end('Останній рядок');Це еквівалентно:
output.write('Останній рядок');
output.end();_finalМетод _final використовується для фінальних асинхронних операцій перед подією finish.
const { Writable } = require('node:stream');
class FinalWritable extends Writable {
_write(chunk, encoding, callback) {
console.log('Дані:', chunk.toString());
callback();
}
_final(callback) {
console.log('Виконуємо фінальні операції');
setTimeout(() => {
console.log('Фінальні операції завершено');
callback();
}, 100);
}
}
const output = new FinalWritable();
output.on('finish', () => {
console.log('Потік завершено');
});
output.end('Дані');callback у _final потрібно викликати обов’язково. Якщо цього не зробити, подія finish не настане.
Якщо під час обробки фрагмента сталася помилка, потрібно передати її в callback:
const { Writable } = require('node:stream');
const output = new Writable({
write(chunk, encoding, callback) {
const text = chunk.toString();
if (text.length === 0) {
callback(new Error('Порожній фрагмент'));
return;
}
console.log(text);
callback();
}
});
output.on('error', (error) => {
console.error('Помилка потоку:', error.message);
});
output.write('Коректні дані');
output.write('');Виклик:
callback(error);має такі наслідки:
потік переходить у стан помилки;
подія error повідомляє про проблему;
подальший запис може бути припинено;
код, який використовує потік, повинен обробити цю помилку.
Не слід ігнорувати подію error. Якщо для потоку немає обробника error, помилка може завершити процес Node.js.
Метод write() повертає логічне значення:
true — потік може приймати нові дані;
false — внутрішній буфер переповнюється, потрібно призупинити запис.
Коли дані з буфера будуть оброблені, потік згенерує подію drain.
const canContinue = writable.write(chunk);
if (!canContinue) {
writable.once('drain', () => {
// Тепер можна продовжити запис
});
}Це називається backpressure — механізм, який не дозволяє швидкому джерелу даних перевантажити повільний Writable Stream.
Якщо записувати великі обсяги даних у циклі та ігнорувати результат write(), програма може використати надмірно багато пам’яті.
highWaterMarkПараметр highWaterMark визначає приблизний обсяг даних, який може перебувати у внутрішньому буфері потоку:
const { Writable } = require('node:stream');
const output = new Writable({
highWaterMark: 1024,
write(chunk, encoding, callback) {
setTimeout(callback, 50);
}
});Для звичайного потоку цей розмір вимірюється в байтах. Для потоку з objectMode: true — у кількості об’єктів.
highWaterMark не обмежує максимальний розмір одного фрагмента. Він визначає момент, коли write() починає повертати false.
У прикладі нижче власний Writable Stream записує рядки у файл. Запис кожного фрагмента виконується асинхронно. Для об’єднання джерела та приймача використовується pipeline, який автоматично передає помилки між потоками та завершує їх.
const { Readable, Writable } = require('node:stream');
const { pipeline } = require('node:stream/promises');
const { appendFile, rm } = require('node:fs/promises');
const fileName = './audit.log';
class AuditLogWritable extends Writable {
constructor(filePath) {
super({
decodeStrings: false,
highWaterMark: 16
});
this.filePath = filePath;
}
_write(chunk, encoding, callback) {
const text = Buffer.isBuffer(chunk)
? chunk.toString('utf8')
: String(chunk);
if (text.trim() === '') {
callback(new Error('Не можна записати порожній рядок'));
return;
}
appendFile(this.filePath, `${text}\n`, 'utf8')
.then(() => callback())
.catch((error) => callback(error));
}
_final(callback) {
appendFile(this.filePath, '-- КІНЕЦЬ ЖУРНАЛУ --\n', 'utf8')
.then(() => callback())
.catch((error) => callback(error));
}
}
async function main() {
await rm(fileName, { force: true });
const records = [
'Користувач увійшов у систему',
'Користувач відкрив профіль',
'Користувач оновив налаштування'
];
const source = Readable.from(records);
const destination = new AuditLogWritable(fileName);
destination.on('finish', () => {
console.log('Запис у журнал завершено');
});
try {
await pipeline(source, destination);
console.log(`Дані записано у файл ${fileName}`);
} catch (error) {
console.error('Не вдалося записати журнал:', error.message);
}
}
main().catch((error) => {
console.error('Непередбачена помилка:', error);
process.exitCode = 1;
});У цьому прикладі:
Readable.from(records) створює джерело даних;
AuditLogWritable приймає рядки;
_write додає кожен рядок до файлу;
_final записує фінальний маркер;
pipeline очікує повного завершення;
помилка з _write або _final потрапляє до catch;
подія finish виникає після успішного завершення всіх записів.
end() і callback завершенняДля ручного запису даних можна передати callback у end():
const { Writable } = require('node:stream');
const output = new Writable({
write(chunk, encoding, callback) {
console.log(chunk.toString());
callback();
}
});
output.on('error', (error) => {
console.error(error);
});
output.write('Перший фрагмент');
output.end('Останній фрагмент', () => {
console.log('Запис завершено');
});Callback end() викликається після завершення запису, але для обробки помилок все одно потрібно слухати подію error.
end() використовується, коли всі дані успішно передано.
destroy() використовується для примусового припинення роботи потоку:
writable.destroy(new Error('Операцію скасовано'));Після цього потік генерує подію error, а потім зазвичай close.
Не слід викликати destroy() замість end(), якщо потрібно коректно завершити запис. destroy() перериває роботу і може залишити незаписані дані.
callbackНеправильно:
_write(chunk, encoding, callback) {
console.log(chunk.toString());
}У такому випадку потік вважає, що обробка ще не завершилася.
Правильно:
_write(chunk, encoding, callback) {
console.log(chunk.toString());
callback();
}callback кілька разівcallback потрібно викликати рівно один раз. Подвійний виклик може призвести до непередбачуваної поведінки або помилки.
Неправильно:
_write(chunk, encoding, callback) {
saveData(chunk).then(() => {
callback();
});
}Якщо saveData відхилить проміс, помилка не буде передана потоку.
Правильно:
_write(chunk, encoding, callback) {
saveData(chunk)
.then(() => callback())
.catch((error) => callback(error));
}end()Після виклику end() не можна використовувати write():
writable.end();
writable.write('Нові дані');Це призводить до помилки, оскільки потік уже прийняв сигнал про завершення.
write()Якщо write() повернув false, не слід продовжувати безконтрольно передавати дані. Потрібно дочекатися drain або використовувати pipeline, який керує backpressure автоматично.
Writable призначений для приймання та обробки частин даних.
У _write потрібно викликати callback після завершення обробки.
Помилки передаються через callback(error) і подію error.
end() повідомляє, що нових даних більше не буде.
_final призначений для фінальних операцій перед завершенням.
Подія finish означає успішне завершення запису.
write() може повернути false, якщо внутрішній буфер переповнений.
Подія drain повідомляє, що запис можна продовжити.
pipeline спрощує керування завершенням, помилками та backpressure.