Работа с большими объёмами данных в приложениях на Node.js — палка о двух концах. Возможность обрабатывать колоссальные объёмы данных чрезвычайно удобна, но может приводить к узким местам производительности и исчерпанию памяти. Традиционно разработчики решали эту задачу, читая весь набор данных в память целиком. Этот подход, хоть и интуитивен для небольших наборов данных, становится неэффективным и ресурсоёмким для больших данных (например, файлов, сетевых запросов…).
Вот тут-то и появляются потоки (streams) Node.js. Потоки предлагают принципиально иной подход, позволяя обрабатывать данные инкрементально и оптимизировать использование памяти. Обрабатывая данные управляемыми порциями (чанками), потоки дают вам возможность строить масштабируемые приложения, эффективно справляющиеся даже с самыми пугающими наборами данных. Как гласит популярная цитата, «streams are arrays over time» (потоки — это массивы во времени).
В этом руководстве мы даём обзор концепции потоков, их истории и API, а также несколько рекомендаций по их использованию и эксплуатации.
Потоки Node.js предлагают мощную абстракцию для управления потоком данных в ваших приложениях. Они великолепны при обработке больших наборов данных — например, при чтении или записи файлов и сетевых запросов — без ущерба для производительности.
Этот подход отличается от загрузки всего набора данных в память целиком. Потоки обрабатывают данные чанками, значительно снижая использование памяти. Все потоки в Node.js наследуются от класса EventEmitter, что позволяет им эмитировать события на различных стадиях обработки данных. Эти потоки могут быть читаемыми (readable), записываемыми (writable) или и теми, и другими, обеспечивая гибкость для разных сценариев работы с данными.
Node.js процветает на событийно-ориентированной архитектуре, что делает его идеальным для I/O в реальном времени. Это означает потребление ввода сразу, как только он доступен, и отправку вывода сразу, как только приложение его сгенерировало. Потоки бесшовно вписываются в этот подход, обеспечивая непрерывную обработку данных.
Достигают они этого, эмитируя события на ключевых стадиях. К этим событиям относятся сигналы о полученных данных (событие data) и о завершении потока (событие end). Разработчики могут слушать эти события и выполнять соответствующую пользовательскую логику. Эта событийно-ориентированная природа делает потоки высокоэффективными для обработки данных из внешних источников.
Потоки дают три ключевых преимущества перед другими методами работы с данными:
- Эффективность по памяти: потоки обрабатывают данные инкрементально, потребляя и обрабатывая их чанками, а не загружая весь набор данных в память. Это большое преимущество при работе с большими наборами данных, поскольку значительно снижает использование памяти и предотвращает связанные с памятью проблемы производительности.
- Улучшенное время отклика: потоки позволяют обрабатывать данные немедленно. Когда приходит чанк данных, его можно обработать, не дожидаясь получения всей полезной нагрузки или набора данных. Это снижает задержку и улучшает общую отзывчивость приложения.
- Масштабируемость для обработки в реальном времени: обрабатывая данные чанками, потоки Node.js могут эффективно справляться с большими объёмами данных при ограниченных ресурсах. Эта масштабируемость делает потоки идеальными для приложений, обрабатывающих большие объёмы данных в реальном времени.
Эти преимущества делают потоки мощным инструментом для создания высокопроизводительных масштабируемых приложений на Node.js, особенно при работе с большими наборами данных или обработке данных в реальном времени.
Если у вашего приложения все данные уже готовы в памяти, использование потоков может добавить ненужных накладных расходов, сложности и замедлить приложение.
Этот раздел — справка по истории потоков в Node.js. Если вы не работаете с кодовой базой, написанной под версию Node.js до 0.11.5 (2013), вы редко встретите старые версии API потоков, но термины всё ещё могут быть в ходу.
Первая версия потоков была выпущена одновременно с Node.js. Хотя класса Stream ещё не было, разные модули использовали эту концепцию и реализовывали функции read/write. Функция util.pump() была доступна для управления потоком данных между потоками.
С выпуском Node v0.4.0 в 2011 году были представлены класс Stream, а также метод pipe().
В 2012 году с выпуском Node v0.10.0 были представлены Streams 2. Это обновление принесло новые подклассы потоков, включая Readable, Writable, Duplex и Transform. Кроме того, было добавлено событие readable. Для сохранения обратной совместимости потоки можно было переключить в старый режим, добавив слушатель события data либо вызвав методы pause() или resume().
В 2013 году были выпущены Streams 3 с Node v0.11.5, чтобы решить проблему потока, имеющего обработчики и события data, и события readable. Это устранило необходимость выбирать между «текущим» и «старым» режимами. Streams 3 — текущая версия потоков в Node.js.
Readable — это класс, который мы используем для последовательного чтения источника данных. Типичные примеры потоков Readable в API Node.js — fs.ReadStream при чтении файлов, http.IncomingMessage при чтении HTTP-запросов и process.stdin при чтении из стандартного ввода.
Читаемый поток работает с несколькими основными методами и событиями, которые позволяют тонко управлять обработкой данных:
on('data'): это событие срабатывает всякий раз, когда из потока доступны данные. Оно очень быстрое, поскольку поток проталкивает (push) данные так быстро, как только может обработать, что делает его подходящим для сценариев с высокой пропускной способностью.on('end'): эмитируется, когда из потока больше нечего читать. Оно означает завершение доставки данных. Это событие срабатывает только тогда, когда все данные из потока потреблены.on('readable'): это событие срабатывает, когда из потока доступны данные для чтения или когда достигнут конец потока. Оно позволяет при необходимости более контролируемо читать данные.on('close'): это событие эмитируется, когда поток и его нижележащие ресурсы закрыты, и указывает, что больше событий не будет эмитировано.on('error'): это событие может быть эмитировано в любой момент, сигнализируя о том, что произошла ошибка обработки. Обработчик этого события можно использовать, чтобы избежать неперехваченных исключений.
Демонстрацию использования этих событий можно увидеть в следующих разделах.
Вот пример простой реализации читаемого потока, который генерирует данные динамически:
class MyStream extends Readable {
#count = 0;
_read(size) {
this.push(':-)');
if (++this.#count === 5) {
this.push(null);
}
}
}
const stream = new MyStream();
stream.on('data', chunk => {
console.log(chunk.toString());
});В этом коде класс MyStream расширяет Readable и переопределяет метод _read(), чтобы протолкнуть строку ":-)" во внутренний буфер. После пятикратного проталкивания строки он сигнализирует о конце потока, проталкивая null. Обработчик события on('data') логирует каждый чанк в консоль по мере получения.
Для ещё более тонкого управления потоком данных можно использовать событие readable. Это событие сложнее, но обеспечивает лучшую производительность для определённых приложений, позволяя явно контролировать, когда данные читаются из потока:
const stream = new MyStream({
highWaterMark: 1,
});
stream.on('readable', () => {
console.count('>> readable event');
let chunk;
while ((chunk = stream.read()) !== null) {
console.log(chunk.toString()); // Обрабатываем чанк
}
});
stream.on('end', () => console.log('>> end event'));Здесь событие readable используется, чтобы вручную вытягивать данные из потока по мере необходимости. Цикл внутри обработчика события readable продолжает читать данные из буфера потока, пока он не вернёт null, что указывает, что буфер временно пуст или поток завершился. Установка highWaterMark в 1 держит размер буфера маленьким, вызывая событие readable чаще и позволяя более гранулярно управлять потоком данных.
С предыдущим кодом вы получите вывод вроде
>> readable event: 1
:-):-)
:-)
:-)
:-)
>> readable event: 2
>> readable event: 3
>> readable event: 4
>> end eventДавайте это переварим. Когда мы прикрепляем событие on('readable'), оно делает первый вызов read(), потому что именно это может вызвать эмит события readable. После эмита этого события мы вызываем read на первой итерации цикла while. Вот почему мы получаем первые два смайлика в одной строке. После этого мы продолжаем вызывать read, пока не будет протолкнут null. Каждый вызов read программирует эмит нового события readable, но поскольку мы в режиме «flow» (то есть используем событие readable), эмит планируется на nextTick. Вот почему мы получаем их все в конце, когда синхронный код цикла завершён.
ПРИМЕЧАНИЕ: можно попробовать запустить код с NODE_DEBUG=stream, чтобы увидеть, что emitReadable срабатывает после каждого push.
Если мы хотим увидеть, что события readable вызываются перед каждым смайликом, можно обернуть push в setImmediate или process.nextTick вот так:
class MyStream extends Readable {
#count = 0;
_read(size) {
setImmediate(() => {
this.push(':-)');
if (++this.#count === 5) {
return this.push(null);
}
});
}
}И мы получим:
>> readable event: 1
:-)
>> readable event: 2
:-)
>> readable event: 3
:-)
>> readable event: 4
:-)
>> readable event: 5
:-)
>> readable event: 6
>> end eventПотоки Writable полезны для создания файлов, загрузки данных или любой задачи, связанной с последовательным выводом данных. Если читаемые потоки предоставляют источник данных, то записываемые потоки в Node.js выступают приёмником для ваших данных. Типичные примеры записываемых потоков в API Node.js — fs.WriteStream, process.stdout и process.stderr.
.write(): этот метод используется для записи чанка данных в поток. Он обрабатывает данные, буферизуя их до заданного предела (highWaterMark), и возвращает булево значение, указывающее, можно ли немедленно записать ещё данные..end(): этот метод сигнализирует о конце процесса записи данных. Он сообщает потоку завершить операцию записи и, возможно, выполнить необходимую очистку.
Вот пример создания записываемого потока, который переводит все входящие данные в верхний регистр перед записью в стандартный вывод:
const { once } = require('node:events');
const { Writable } = require('node:stream');
class MyStream extends Writable {
constructor() {
super({ highWaterMark: 10 /* 10 bytes */ });
}
_write(data, encode, cb) {
process.stdout.write(data.toString().toUpperCase() + '\n', cb);
}
}
async function main() {
const stream = new MyStream();
for (let i = 0; i < 10; i++) {
const waitDrain = !stream.write('hello');
if (waitDrain) {
console.log('>> wait drain');
await once(stream, 'drain');
}
}
stream.end('world');
}
// Вызываем async-функцию
main().catch(console.error);В этом коде MyStream — это пользовательский поток Writable с ёмкостью буфера (highWaterMark) в 10 байт. Он переопределяет метод _write, чтобы перевести данные в верхний регистр перед выводом.
Цикл пытается записать hello в поток десять раз. Если буфер заполняется (waitDrain становится true), он ждёт события drain, прежде чем продолжить, гарантируя, что мы не переполним буфер потока.
Вывод будет таким:
HELLO
>> wait drain
HELLO
HELLO
>> wait drain
HELLO
HELLO
>> wait drain
HELLO
HELLO
>> wait drain
HELLO
HELLO
>> wait drain
HELLO
WORLDПотоки Duplex реализуют интерфейсы и readable, и writable.
Дуплексные потоки реализуют все методы и события, описанные в разделах о читаемых и записываемых потоках.
Хороший пример дуплексного потока — класс Socket в модуле net:
const net = require('node:net');
// Создаём TCP-сервер
const server = net.createServer(socket => {
socket.write('Hello from server!\n');
socket.on('data', data => {
console.log(`Client says: ${data.toString()}`);
});
// Обрабатываем отключение клиента
socket.on('end', () => {
console.log('Client disconnected');
});
});
// Запускаем сервер на порту 8080
server.listen(8080, () => {
console.log('Server listening on port 8080');
});Предыдущий код откроет TCP-сокет на порту 8080, отправит Hello from server! любому подключающемуся клиенту и залогирует любые полученные данные.
const net = require('node:net');
// Подключаемся к серверу по localhost:8080
const client = net.createConnection({ port: 8080 }, () => {
client.write('Hello from client!\n');
});
client.on('data', data => {
console.log(`Server says: ${data.toString()}`);
});
// Обрабатываем закрытие соединения сервером
client.on('end', () => {
console.log('Disconnected from server');
});Предыдущий код подключится к TCP-сокету, отправит сообщение Hello from client и залогирует любые полученные данные.
Потоки Transform — это дуплексные потоки, где вывод вычисляется на основе ввода. Как следует из названия, их обычно используют между читаемым и записываемым потоком, чтобы преобразовывать данные по мере их прохождения.
Помимо всех методов и событий дуплексных потоков, есть:
_transform: эта функция вызывается внутренне для управления потоком данных между читаемой и записываемой частями. Она НЕ ДОЛЖНА вызываться кодом приложения.
Чтобы создать новый поток transform, мы можем передать объект options в конструктор Transform, включая функцию transform, которая обрабатывает, как выходные данные вычисляются из входных с помощью метода push.
const { Transform } = require('node:stream');
const upper = new Transform({
transform(data, enc, cb) {
this.push(data.toString().toUpperCase());
cb();
},
});Этот поток возьмёт любой ввод и выведет его в верхнем регистре.
Работая с потоками, мы обычно хотим прочитать из источника и записать в приёмник, возможно нуждаясь в некоторой трансформации данных между ними. В следующих разделах рассматриваются разные способы это сделать.
Метод .pipe() соединяет один читаемый поток с записываемым (или transform) потоком. Хотя это выглядит простым способом достичь нашей цели, он делегирует всю обработку ошибок программисту, из-за чего сделать это правильно трудно.
Следующий пример показывает pipe, пытающийся вывести текущий файл в верхнем регистре в консоль.
const fs = require('node:fs');
const { Transform } = require('node:stream');
let errorCount = 0;
const upper = new Transform({
transform(data, enc, cb) {
if (errorCount === 10) {
return cb(new Error('BOOM!'));
}
errorCount++;
this.push(data.toString().toUpperCase());
cb();
},
});
const readStream = fs.createReadStream(__filename, { highWaterMark: 1 });
const writeStream = process.stdout;
readStream.pipe(upper).pipe(writeStream);
readStream.on('close', () => {
console.log('Readable stream closed');
});
upper.on('close', () => {
console.log('Transform stream closed');
});
upper.on('error', err => {
console.error('\nError in transform stream:', err.message);
});
writeStream.on('close', () => {
console.log('Writable stream closed');
});После записи 10 символов upper вернёт ошибку в колбэке, что вызовет закрытие потока. Однако другие потоки не будут уведомлены, что приведёт к утечкам памяти. Вывод будет таким:
CONST FS =
Error in transform stream: BOOM!
Transform stream closedpipeline(): voidЧтобы избежать подводных камней и низкоуровневой сложности метода .pipe(), в большинстве случаев рекомендуется использовать метод pipeline(). Этот метод — более безопасный и надёжный способ соединять потоки в цепочку, автоматически обрабатывая ошибки и очистку.
Следующий пример демонстрирует, как использование pipeline() предотвращает подводные камни предыдущего примера:
const fs = require('node:fs');
const { Transform, pipeline } = require('node:stream');
let errorCount = 0;
const upper = new Transform({
transform(data, enc, cb) {
if (errorCount === 10) {
return cb(new Error('BOOM!'));
}
errorCount++;
this.push(data.toString().toUpperCase());
cb();
},
});
const readStream = fs.createReadStream(__filename, { highWaterMark: 1 });
const writeStream = process.stdout;
readStream.on('close', () => {
console.log('Readable stream closed');
});
upper.on('close', () => {
console.log('\nTransform stream closed');
});
writeStream.on('close', () => {
console.log('Writable stream closed');
});
pipeline(readStream, upper, writeStream, err => {
if (err) {
return console.error('Pipeline error:', err.message);
}
console.log('Pipeline succeeded');
});В этом случае все потоки будут закрыты, со следующим выводом:
CONST FS =
Transform stream closed
Writable stream closed
Pipeline error: BOOM!
Readable stream closedУ метода pipeline() также есть async-версия pipeline(), которая не принимает колбэк, а вместо этого возвращает промис, отклоняемый при сбое пайплайна.
Async-итераторы рекомендуются как стандартный способ взаимодействия с API потоков. По сравнению со всеми примитивами потоков — и в Web, и в Node.js — async-итераторы легче понять и использовать, что способствует меньшему числу багов и более поддерживаемому коду. В недавних версиях Node.js async-итераторы стали более элегантным и читаемым способом взаимодействия с потоками. Опираясь на фундамент событий, async-итераторы предоставляют абстракцию более высокого уровня, упрощающую потребление потоков.
В Node.js все читаемые потоки являются асинхронно итерируемыми (async iterables). Это означает, что вы можете использовать синтаксис for await...of, чтобы перебирать данные потока по мере их появления, обрабатывая каждый кусок данных с эффективностью и простотой асинхронного кода.
Использование async-итераторов с потоками упрощает обработку асинхронных потоков данных несколькими способами:
- Улучшенная читаемость: структура кода чище и читаемее, особенно при работе с несколькими асинхронными источниками данных.
- Обработка ошибок: async-итераторы позволяют прямолинейно обрабатывать ошибки с помощью блоков try/catch, аналогично обычным асинхронным функциям.
- Управление потоком: они по своей природе управляют backpressure, поскольку потребитель управляет потоком, ожидая (await) следующий кусок данных, что обеспечивает более эффективное использование памяти и обработку.
Async-итераторы предлагают более современный и часто более читаемый способ работы с читаемыми потоками, особенно при работе с асинхронными источниками данных или когда вы предпочитаете более последовательный, основанный на цикле подход к обработке данных.
Вот пример, демонстрирующий использование async-итераторов с читаемым потоком:
const fs = require('node:fs');
const { pipeline } = require('node:stream/promises');
async function main() {
await pipeline(
fs.createReadStream(__filename),
async function* (source) {
for await (let chunk of source) {
yield chunk.toString().toUpperCase();
}
},
process.stdout
);
}
main().catch(console.error);Этот код достигает того же результата, что и предыдущие примеры, без необходимости определять новый transform-поток. Ошибка из предыдущих примеров убрана ради краткости. Использована async-версия pipeline, и её следует обернуть в блок try...catch для обработки возможных ошибок.
По умолчанию потоки могут работать со строками, Buffer, TypedArray или DataView. Если в поток протолкнуть произвольное значение, отличное от этих (например, объект), будет выброшен TypeError. Однако работать с объектами можно, установив опцию objectMode в true. Это позволяет потоку работать с любым значением JavaScript, кроме null, который используется для сигнала о конце потока. Это значит, что вы можете push и read любое значение в читаемом потоке и write любое значение в записываемом потоке.
const { Readable } = require('node:stream');
const readable = Readable({
objectMode: true,
read() {
this.push({ hello: 'world' });
this.push(null);
},
});При работе в объектном режиме важно помнить, что опция highWaterMark относится к числу объектов, а не байтов.
При использовании потоков важно убедиться, что производитель (producer) не перегружает потребителя (consumer). Для этого во всех потоках в API Node.js используется механизм backpressure, и разработчики отвечают за поддержание этого поведения.
В любом сценарии, где буфер данных превысил highWaterMark или очередь записи в данный момент занята, .write() вернёт false.
Когда возвращается false, включается система backpressure. Она приостановит входящий поток Readable от отправки данных и подождёт, пока потребитель снова не будет готов. Как только буфер данных опустеет, будет эмитировано событие 'drain', чтобы возобновить входящий поток данных.
Для более глубокого понимания backpressure посмотрите руководство по backpressure.
Концепция потоков не эксклюзивна для Node.js. На самом деле у Node.js есть другая реализация концепции потоков — Web Streams, которая реализует стандарт WHATWG Streams. Хотя стоящие за ними концепции похожи, важно осознавать, что у них разные API и они напрямую несовместимы.
Web Streams реализуют классы ReadableStream, WritableStream и TransformStream, которые гомологичны потокам Node.js Readable, Writable и Transform.
Node.js предоставляет служебные функции для конвертации в/из Web Streams и потоков Node.js. Эти функции реализованы как методы toWeb и fromWeb в каждом классе потока.
Следующий пример в классе Duplex демонстрирует, как работать и с читаемыми, и с записываемыми потоками, сконвертированными в Web Streams:
const { Duplex } = require('node:stream');
const duplex = Duplex({
read() {
this.push('world');
this.push(null);
},
write(chunk, encoding, callback) {
console.log('writable', chunk);
callback();
},
});
const { readable, writable } = Duplex.toWeb(duplex);
writable.getWriter().write('hello');
readable
.getReader()
.read()
.then(result => {
console.log('readable', result.value);
});Вспомогательные функции полезны, если вам нужно вернуть Web Stream из модуля Node.js или наоборот. Для обычного потребления потоков async-итераторы обеспечивают бесшовное взаимодействие и с Node.js, и с Web Streams.
const { pipeline } = require('node:stream/promises');
async function main() {
const { body } = await fetch('https://nodejs.org/api/stream.html');
await pipeline(
body,
new TextDecoderStream(),
async function* (source) {
for await (const chunk of source) {
yield chunk.toString().toUpperCase();
}
},
process.stdout
);
}
main().catch(console.error);Учтите, что тело ответа fetch — это ReadableStream<Uint8Array>, и поэтому нужен TextDecoderStream, чтобы работать с чанками как со строками.
Эта работа основана на материалах, опубликованных Маттео Коллиной (Matteo Collina) в блоге Platformatic.