Пошук уроків, статей та іншого контенту
Створюйте Readable Streams і читайте дані частинами з файлів та інших джерел.
Readable Stream — це об’єкт, з якого програма може послідовно читати дані частинами, або чанками.
Замість того щоб завантажувати весь файл у пам’ять, потік читає його фрагмент за фрагментом. Це особливо важливо для великих файлів, мережевих відповідей і джерел даних, обсяг яких наперед невідомий.
Основні переваги:
менше використання оперативної пам’яті;
можливість почати обробку ще до повного отримання даних;
підтримка великих файлів;
керування швидкістю читання через механізм зворотного тиску.
У Node.js Readable Stream реалізовано класом Readable з модуля node:stream.
Найпростіший спосіб створити Readable Stream для файлу — використати createReadStream з модуля node:fs.
const fs = require('node:fs');
const stream = fs.createReadStream('./data.txt', {
encoding: 'utf8',
highWaterMark: 16
});
stream.on('data', (chunk) => {
console.log('Отримано частину:', JSON.stringify(chunk));
});
stream.on('end', () => {
console.log('Читання завершено');
});
stream.on('error', (error) => {
console.error('Помилка читання:', error.message);
});У цьому прикладі:
encoding: 'utf8' перетворює отримані дані на рядки;
highWaterMark: 16 задає орієнтовний розмір внутрішнього буфера в байтах;
подія data спрацьовує для кожної отриманої частини;
подія end означає, що всі дані прочитано;
подія error повідомляє про помилки, наприклад відсутність файлу.
Якщо encoding не вказати, chunk буде об’єктом Buffer.
const fs = require('node:fs');
const stream = fs.createReadStream('./data.txt');
stream.on('data', (chunk) => {
console.log('Тип:', Buffer.isBuffer(chunk));
console.log('Розмір:', chunk.length);
});Під час роботи з бінарними даними, наприклад зображеннями або архівами, зазвичай потрібно залишати дані у форматі Buffer.
Readable Stream можна читати за допомогою for await...of. Такий синтаксис зручний, коли потрібно послідовно обробляти частини даних в асинхронній функції.
const fs = require('node:fs');
async function readFileInChunks() {
const stream = fs.createReadStream('./data.txt', {
encoding: 'utf8',
highWaterMark: 32
});
try {
for await (const chunk of stream) {
console.log('Отримано частину:', JSON.stringify(chunk));
}
console.log('Читання завершено');
} catch (error) {
console.error('Помилка читання:', error.message);
}
}
readFileInChunks();Цей підхід має кілька переваг:
код виглядає послідовним;
помилки можна обробляти через try...catch;
наступна частина читається лише після завершення обробки поточної;
не потрібно вручну підписуватися на data та end.
Наприклад, можна підрахувати кількість рядків у великому текстовому файлі:
const fs = require('node:fs');
async function countLines(filePath) {
const stream = fs.createReadStream(filePath, {
encoding: 'utf8'
});
let remainder = '';
let lineCount = 0;
for await (const chunk of stream) {
const text = remainder + chunk;
const lines = text.split('\n');
remainder = lines.pop();
lineCount += lines.length;
}
if (remainder.length > 0) {
lineCount += 1;
}
return lineCount;
}
countLines('./data.txt')
.then((count) => {
console.log('Кількість рядків:', count);
})
.catch((error) => {
console.error('Не вдалося прочитати файл:', error.message);
});Один рядок може бути розділений між двома чанками. Тому змінна remainder зберігає неповний рядок до отримання наступної частини.
Readable Stream може працювати в одному з двох основних режимів:
flowing mode — дані автоматично надходять через подію data;
paused mode — дані читаються вручну за допомогою read() або через for await...of.
Додавання обробника data переводить потік у flowing mode:
const fs = require('node:fs');
const stream = fs.createReadStream('./data.txt', {
encoding: 'utf8'
});
stream.on('data', (chunk) => {
console.log(chunk);
});У paused mode можна керувати читанням вручну:
const fs = require('node:fs');
const stream = fs.createReadStream('./data.txt', {
encoding: 'utf8'
});
stream.on('readable', () => {
let chunk;
while ((chunk = stream.read()) !== null) {
console.log('Прочитано:', JSON.stringify(chunk));
}
});
stream.on('end', () => {
console.log('Читання завершено');
});Подія readable означає, що в потоці з’явилися дані, доступні для читання. Метод read() повертає наступну частину або null, якщо доступних даних наразі немає.
На практиці для більшості задач зручніше використовувати:
for await...of;
подію data, якщо потрібна подієва модель.
Власний Readable Stream створюють через клас Readable і реалізацію методу _read().
Метод _read() викликається Node.js, коли потоку потрібні нові дані. Щоб додати дані до потоку, потрібно викликати this.push().
Коли дані закінчилися, слід викликати:
this.push(null);null — спеціальний сигнал кінця потоку.
const { Readable } = require('node:stream');
class NumberStream extends Readable {
constructor(max) {
super({
objectMode: true
});
this.current = 1;
this.max = max;
}
_read() {
if (this.current > this.max) {
// Повідомляємо, що дані закінчилися
this.push(null);
return;
}
// Додаємо до потоку одне число
this.push(this.current);
this.current += 1;
}
}
async function main() {
const stream = new NumberStream(5);
for await (const number of stream) {
console.log('Отримано число:', number);
}
console.log('Потік завершено');
}
main().catch((error) => {
console.error(error);
});Результат:
Отримано число: 1
Отримано число: 2
Отримано число: 3
Отримано число: 4
Отримано число: 5
Потік завершеноЗа замовчуванням Readable Stream працює з рядками та Buffer. Щоб передавати довільні JavaScript-значення, потрібно ввімкнути objectMode.
const { Readable } = require('node:stream');
const users = Readable.from(
[
{ id: 1, name: 'Олена' },
{ id: 2, name: 'Андрій' },
{ id: 3, name: 'Марія' }
],
{
objectMode: true
}
);
async function printUsers() {
for await (const user of users) {
console.log(`${user.id}: ${user.name}`);
}
}
printUsers().catch(console.error);У режимі об’єктів потік передає окремі JavaScript-об’єкти, а не байти. highWaterMark у такому режимі визначає кількість об’єктів, а не їхній розмір у байтах.
Метод Readable.from() дозволяє створити потік із синхронного або асинхронного ітерованого об’єкта.
const { Readable } = require('node:stream');
const stream = Readable.from(['Node.js', ' ', 'Readable', ' ', 'Stream']);
stream.on('data', (chunk) => {
process.stdout.write(chunk);
});
stream.on('end', () => {
process.stdout.write('\n');
});Асинхронний генератор зручно використовувати, коли частини даних надходять із затримкою або обчислюються поступово:
const { Readable } = require('node:stream');
async function* generateMessages() {
for (let index = 1; index <= 3; index += 1) {
await new Promise((resolve) => setTimeout(resolve, 300));
// Генеруємо чергову частину даних
yield `Повідомлення ${index}\n`;
}
}
async function main() {
const stream = Readable.from(generateMessages());
for await (const message of stream) {
process.stdout.write(message);
}
}
main().catch(console.error);Потік читає значення генератора в міру потреби, а не створює всі повідомлення заздалегідь.
Зворотний тиск виникає, коли споживач не встигає обробляти дані з тією самою швидкістю, з якою джерело їх генерує.
Readable Stream не повинен безмежно накопичувати дані в пам’яті. Node.js використовує внутрішній буфер і зупиняє запит нових частин, коли буфер заповнений.
Параметр highWaterMark визначає поріг, після якого потік має призупинити активне наповнення буфера:
const fs = require('node:fs');
const stream = fs.createReadStream('./large-file.txt', {
highWaterMark: 64 * 1024
});
stream.on('data', (chunk) => {
console.log(`Отримано ${chunk.length} байт`);
});highWaterMark не гарантує точний розмір кожного чанка. Це налаштування внутрішньої буферизації, а не вимога до кожної частини.
Під час читання через for await...of призупинення відбувається природно: наступна частина запитується після завершення обробки поточної.
Основні сигнали Readable Stream:
data — отримано частину даних;
readable — доступні дані для ручного читання;
end — усі дані прочитано;
close — ресурс потоку закрито;
error — виникла помилка.
Подія end не означає, що потік успішно закрився в усіх можливих сценаріях. Для помилок потрібно окремо передбачати обробку error або використовувати try...catch з for await...of.
Після завершення роботи потік можна зупинити методом destroy():
const fs = require('node:fs');
const stream = fs.createReadStream('./large-file.txt');
stream.on('data', (chunk) => {
console.log(`Отримано ${chunk.length} байт`);
if (chunk.length > 10000) {
// Достроково зупиняємо читання
stream.destroy(new Error('Частина даних завелика'));
}
});
stream.on('error', (error) => {
console.error('Потік завершено з помилкою:', error.message);
});
stream.on('close', () => {
console.log('Ресурс потоку закрито');
});Викликайте destroy() лише тоді, коли потрібно достроково припинити роботу потоку або повідомити про помилку.
const chunks = [];
stream.on('data', (chunk) => {
chunks.push(chunk);
});Такий підхід знову накопичує весь файл у пам’яті. Перевага потоків зникає, якщо дані потрібно лише послідовно обробити.
Потік працює з частинами байтів, а не з логічними рядками. Один рядок може бути розділений між кількома чанками, а один chunk може містити багато рядків.
Для текстової обробки потрібно зберігати неповний фрагмент між ітераціями, як у прикладі з підрахунком рядків.
push(null)У власному Readable Stream _read() має повідомити про кінець даних:
this.push(null);Якщо цього не зробити, споживач може чекати на подальші дані нескінченно.
push() після завершенняПісля this.push(null) не можна додавати нові дані. Кінець потоку має бути остаточним.
У звичайному режимі Readable Stream приймає рядки, Buffer або Uint8Array. Якщо потрібно передавати об’єкти, потрібно вказати:
objectMode: trueФайловий потік може завершитися помилкою через відсутній файл, відсутність дозволів або проблеми з файловою системою. Завжди обробляйте error або використовуйте try...catch з for await...of.
Readable Stream читає дані частинами, не завантажуючи весь ресурс у пам’ять.
fs.createReadStream() використовується для потокового читання файлів.
Дані можна отримувати через подію data, метод read() або for await...of.
Власний потік створюють через Readable і метод _read().
this.push() додає дані, а this.push(null) завершує потік.
objectMode дозволяє передавати JavaScript-об’єкти замість байтів.
highWaterMark керує внутрішньою буферизацією та допомагає обмежувати споживання пам’яті.
Потоки потрібно правильно завершувати й обробляти їхні помилки.