Что такое Streams в Node.js?
Что такое Streams в Node.js
Streams (потоки) — это интерфейс для чтения и записи данных в Node.js, построенный на событийной модели. Вместо того чтобы загружать весь файл или ответ целиком в память, поток обрабатывает данные небольшими кусками (chunks) по мере их поступления.
Это фундаментальный механизм, на котором построены HTTP-запросы/ответы, чтение файлов, работа с сокетами и многое другое.
Четыре типа потоков
Readable (читаемый)
Источник данных. Примеры: fs.createReadStream(), http.IncomingMessage, process.stdin.
Writable (записываемый)
Приёмник данных. Примеры: fs.createWriteStream(), http.ServerResponse, process.stdout.
Duplex (двунаправленный)
Одновременно читаемый и записываемый. Пример: net.Socket (TCP-сокет).
Transform (преобразующий)
Разновидность Duplex, которая преобразует данные на лету. Примеры: zlib.createGzip(), crypto.createCipher().
Режимы работы Readable-потоков
Readable-поток работает в двух режимах:
- flowing — данные читаются автоматически и передаются через события
data - paused — данные читаются явно методом
.read()
Переход между режимами происходит при подписке/отписке на событие data или при вызове .pipe() / .resume() / .pause().
Метод pipe()
pipe() — самый удобный способ соединить потоки: он автоматически управляет backpressure, передавая данные из readable в writable с нужной скоростью.
Backpressure
Backpressure — механизм контроля скорости: если получатель не успевает обрабатывать данные, поток приостанавливает чтение, предотвращая переполнение памяти. pipe() реализует это автоматически. При ручном управлении нужно проверять возвращаемое значение writable.write() и слушать событие drain.
Практические преимущества
- Экономия памяти: 1 ГБ файл обрабатывается без загрузки в память целиком
- Меньшая задержка: клиент получает первые данные раньше, не дожидаясь полной обработки
- Компонуемость: потоки соединяются в pipeline, каждый отвечает за одну задачу
stream.pipeline()
С Node.js 10+ рекомендуется использовать stream.pipeline() вместо цепочки .pipe(), так как он корректно обрабатывает ошибки и освобождает ресурсы.
import { pipeline } from 'stream/promises';
import { createReadStream, createWriteStream } from 'fs';
import { createGzip } from 'zlib';
await pipeline(
createReadStream('input.txt'),
createGzip(),
createWriteStream('output.txt.gz')
);
Что хочет услышать интервьюер
Понимание, что потоки решают проблему памяти при работе с большими данными
Знание четырёх типов потоков: Readable, Writable, Duplex, Transform
Объяснение механизма backpressure и зачем он нужен
Умение рассказать о pipe() и предпочтительном stream.pipeline()
Примеры реального применения: сжатие файлов, HTTP-ответы, ETL-пайплайны
Пример: Чтение большого файла через Readable stream
import { createReadStream } from 'fs';
const stream = createReadStream('./big-file.log', {
encoding: 'utf8',
highWaterMark: 64 * 1024 // размер чанка 64 КБ
});
stream.on('data', (chunk: string) => {
// обрабатываем каждый кусок отдельно, не храня весь файл
console.log(`Получен чанк размером ${chunk.length} символов`);
});
stream.on('end', () => {
console.log('Файл полностью прочитан');
});
stream.on('error', (err: Error) => {
console.error('Ошибка чтения:', err.message);
});
Пример: Transform-поток для преобразования данных
import { Transform, TransformCallback } from 'stream';
// Поток, переводящий текст в верхний регистр
class UpperCaseTransform extends Transform {
_transform(
chunk: Buffer,
encoding: BufferEncoding,
callback: TransformCallback
): void {
// передаём преобразованный чанк дальше по pipeline
this.push(chunk.toString().toUpperCase());
callback();
}
}
// Использование в pipeline
import { pipeline } from 'stream/promises';
import { createReadStream, createWriteStream } from 'fs';
await pipeline(
createReadStream('input.txt'),
new UpperCaseTransform(),
createWriteStream('output.txt')
);
Пример: Backpressure при ручной записи
import { createWriteStream, WriteStream } from 'fs';
async function writeWithBackpressure(
stream: WriteStream,
data: string[]
): Promise<void> {
for (const chunk of data) {
// write() возвращает false, если буфер переполнен
const canContinue = stream.write(chunk);
if (!canContinue) {
// ждём, пока буфер освободится
await new Promise<void>((resolve) => stream.once('drain', resolve));
}
}
// явно закрываем поток после записи
await new Promise<void>((resolve, reject) => {
stream.end((err?: Error | null) => (err ? reject(err) : resolve()));
});
}
const ws = createWriteStream('output.txt');
const lines = Array.from({ length: 100000 }, (_, i) => `Строка ${i}\n`);
await writeWithBackpressure(ws, lines);
Типичные ошибки
Путают потоки с буферами — не понимают разницы между Buffer и Stream
Игнорируют backpressure при ручном использовании writable.write(), что приводит к утечкам памяти
Используют цепочку .pipe() без обработки ошибок вместо stream.pipeline()
Не знают разницы между режимами flowing и paused у Readable-потоков
Считают Transform-поток отдельной концепцией, а не частью семейства Duplex


