Поток
Поток — это абстрактный интерфейс, реализованный различными объектами в Node.js. Например, запрос к HTTP-серверу является потоком, как и process.stdout. Потоки могут быть читаемыми, записываемыми или и тем, и другим. Все потоки являются экземплярами EventEmitter.
Вы можете загрузить базовые классы потоков, выполнив require('stream'). Предоставляются базовые классы для потоков Readable, Writable, Duplex и Transform.
Этот документ разделён на 3 раздела:
- Первый раздел объясняет части API, которые вам необходимо знать для использования потоков в ваших программах.
- Второй раздел объясняет части API, которые вам необходимо использовать, если вы сами реализуете свои пользовательские потоки. API разработан, чтобы сделать это для вас легко.
- Третий раздел более подробно рассматривает работу потоков, включая некоторые внутренние механизмы и функции, которые вам, вероятно, не следует изменять, если вы не уверены в том, что делаете.
API для потребителей потоков
Потоки могут быть читаемыми, записываемыми или обоими (Duplex).
Все потоки являются EventEmitters, но они также имеют другие пользовательские методы и свойства в зависимости от того, являются ли они читаемыми, записываемыми или дуплексными.
Если поток является и читаемым, и записываемым, то он реализует все методы и события. Таким образом, поток Duplex или Transform полностью описан этим API, хотя его реализация может быть несколько иной.
Для потребления потоков в ваших программах не нужно реализовывать интерфейсы Stream. Если вы реализуете интерфейсы потоков в своей программе, обратитесь также к API для разработчиков потоков.
Практически все программы Node.js, независимо от их сложности, используют потоки каким-либо образом. Вот пример использования потоков в программе Node.js:
const http = require('http');
var server = http.createServer( (req, res) => {
// req is an http.IncomingMessage, which is a Readable Stream
// res is an http.ServerResponse, which is a Writable Stream
var body = '';
// we want to get the data as utf8 strings
// If you don't set an encoding, then you'll get Buffer objects
req.setEncoding('utf8');
// Readable streams emit 'data' events once a listener is added
req.on('data', (chunk) => {
body += chunk;
});
// the end event tells you that you have entire body
req.on('end', () => {
try {
const data = JSON.parse(body);
// write back something interesting to the user:
res.write(typeof data);
res.end();
} catch (er) {
// uh oh! bad json!
res.statusCode = 400;
return res.end(`error: ${er.message}`);
}
});
});
server.listen(1337);
// $ curl localhost:1337 -d '{}'
// object
// $ curl localhost:1337 -d '"foo"'
// string
// $ curl localhost:1337 -d 'not json'
// error: Unexpected token o
Класс: stream.Duplex
Потоки Duplex — это потоки, которые реализуют как Readable, так и Writable интерфейсы.
Примеры потоков Duplex включают:
Класс: stream.Readable
Интерфейс потока Readable — это абстракция для источника данных, из которого вы читаете. Другими словами, данные выходят из потока Readable.
Поток Readable не начнёт отправлять данные, пока вы не укажете, что готовы их принять.
Потоки Readable имеют два «режима»: режим потока и режим приостановки. В режиме потока данные считываются из базовой системы и предоставляются вашей программе как можно быстрее. В режиме приостановки вы должны явно вызвать stream.read(), чтобы получить фрагменты данных. Потоки начинают работу в режиме приостановки.
Примечание: если обработчики событий данных не прикреплены, и нет stream.pipe() пунктов назначения, и поток переключен в режим потока, то данные будут потеряны.
Вы можете переключиться в режим потока, выполнив любое из следующих действий:
- Добавив обработчик событий
'data'для прослушивания данных. - Вызвав метод
stream.resume(), чтобы явно открыть поток. - Вызвав метод
stream.pipe(), чтобы отправить данные в Writable.
Вы можете вернуться в режим приостановки, выполнив любое из следующих действий:
- Если нет пунктов назначения подключения, вызвав метод
stream.pause(). - Если есть пункты назначения подключения, удалив все обработчики событий
'data'и удалив все пункты назначения подключения, вызвав методstream.unpipe().
Обратите внимание, что по причинам обратной совместимости удаление обработчиков событий 'data' не автоматически приостановит поток. Кроме того, если есть подключённые пункты назначения, то вызов stream.pause() не гарантирует, что поток останется приостановленным после того, как эти пункты назначения опорожнят буфер и запросят больше данных.
Примеры читаемых потоков включают:
- HTTP-ответы со стороны клиента
- HTTP-запросы на стороне сервера
- потоки чтения fs
- потоки zlib
- потоки crypto
- сокеты TCP
- stdout и stderr дочернего процесса
process.stdin
Событие: 'close'
Используется, когда поток и все его базовые ресурсы (например, дескриптор файла) закрыты. Событие указывает, что больше событий не будет отправлено, и дальнейшие вычисления не будут выполняться.
Не все потоки будут отправлять событие 'close', так как событие 'close' является необязательным.
Событие: 'data'
Прикрепление обработчика событий 'data' к потоку, который не был явно приостановлен, переведёт поток в режим потока. Данные будут переданы как только станут доступны.
Если вы просто хотите получить все данные из потока как можно быстрее, это лучший способ сделать это.
var readable = getReadableStreamSomehow();
readable.on('data', (chunk) => {
console.log('got %d bytes of data', chunk.length);
});
Событие: 'end'
Это событие срабатывает, когда больше нет данных для чтения.
Обратите внимание, что событие 'end' не будет сработать, если данные не будут полностью обработаны. Это можно сделать, переключив поток в режим потока или вызвав stream.read() повторно, пока не достигнете конца.
var readable = getReadableStreamSomehow();
readable.on('data', (chunk) => {
console.log('got %d bytes of data', chunk.length);
});
readable.on('end', () => {
console.log('there will be no more data.');
});
Событие: 'error'
Используется, если произошла ошибка при получении данных.
Событие: 'readable'
Когда из потока можно прочитать фрагмент данных, он отправит событие 'readable'.
В некоторых случаях прослушивание события 'readable' заставит некоторые данные быть прочитанными во внутренний буфер из базовой системы, если это ещё не было сделано.
var readable = getReadableStreamSomehow();
readable.on('readable', () => {
// there is some data to read now
});
После опорожнения внутреннего буфера событие 'readable' снова сработает, когда появятся новые данные.
Событие 'readable' не используется в режиме «потока», за исключением последнего, при достижении конца потока.
Событие 'readable' указывает, что в потоке есть новая информация: доступны новые данные или достигнут конец потока. В первом случае stream.read() вернёт эти данные. Во втором случае stream.read() вернёт null. Например, в следующем примере foo.txt — это пустой файл:
const fs = require('fs');
var rr = fs.createReadStream('foo.txt');
rr.on('readable', () => {
console.log('readable:', rr.read());
});
rr.on('end', () => {
console.log('end');
});
Вывод выполнения этого скрипта:
$ node test.js readable: null end
readable.isPaused()
- Возвращает: <Булево значение>
Этот метод возвращает, приостановлен ли readable явным образом кодом клиента (используя stream.pause() без соответствующего stream.resume()).
var readable = new stream.Readable readable.isPaused() // === false readable.pause() readable.isPaused() // === true readable.resume() readable.isPaused() // === false
readable.pause()
- Возвращает:
this
Этот метод заставит поток в режиме потока прекратить отправку событий 'data', переключившись из режима потока. Любые доступные данные останутся во внутреннем буфере.
var readable = getReadableStreamSomehow();
readable.on('data', (chunk) => {
console.log('got %d bytes of data', chunk.length);
readable.pause();
console.log('there will be no more data for 1 second');
setTimeout(() => {
console.log('now data will start flowing again');
readable.resume();
}, 1000);
});
readable.pipe(destination[, options])
-
destination<stream.Writable> Пункт назначения для записи данных -
options<Объект> Параметры подключения-
end<Булево значение> Завершить запись при завершении чтения. По умолчанию =true
-
Этот метод извлекает все данные из потока чтения и записывает их в предоставленное место назначения, автоматически управляя потоком, чтобы место назначения не было перегружено быстрым потоком чтения.
Можно подключить несколько мест назначения безопасно.
var readable = getReadableStreamSomehow();
var writable = fs.createWriteStream('file.txt');
// All the data from readable goes into 'file.txt'
readable.pipe(writable);
Эта функция возвращает поток назначения, поэтому вы можете создавать цепочки подключения, например:
var r = fs.createReadStream('file.txt');
var z = zlib.createGzip();
var w = fs.createWriteStream('file.txt.gz');
r.pipe(z).pipe(w);
Например, эмулируя команду Unix cat:
process.stdin.pipe(process.stdout);
По умолчанию stream.end() вызывается в пункте назначения, когда исходный поток отправляет событие 'end', чтобы destination больше не был доступен для записи. Передайте { end: false } как options для поддержания открытого потока назначения.
Это сохраняет writer открытым, чтобы в конце можно было записать «Goodbye».
reader.pipe(writer, { end: false });
reader.on('end', () => {
writer.end('Goodbye\n');
});
Обратите внимание, что process.stderr и process.stdout никогда не закрываются до завершения процесса, независимо от указанных параметров.
readable.read([size])
Метод read() извлекает данные из внутреннего буфера и возвращает их. Если данных нет, возвращает null.
Если вы передаёте аргумент size, то он вернёт указанное количество байт. Если size байтов недоступно, он вернёт null, за исключением случая, когда поток завершён, в этом случае возвращаются оставшиеся данные в буфере.
Если вы не указываете аргумент size, то возвращаются все данные из внутреннего буфера.
Этот метод следует вызывать только в режиме приостановки. В режиме потока этот метод вызывается автоматически до тех пор, пока внутренний буфер не будет исчерпан.
var readable = getReadableStreamSomehow();
readable.on('readable', () => {
var chunk;
while (null !== (chunk = readable.read())) {
console.log('got %d bytes of data', chunk.length);
}
});
Если этот метод возвращает фрагмент данных, он также вызовет событие 'data'.
Обратите внимание, что вызов stream.read([size]) после срабатывания события 'end' вернёт null. Ошибка выполнения не будет генерироваться.
readable.resume()
- Возвращает:
this
Этот метод заставляет читаемый поток возобновить отправку событий 'data'.
Этот метод переключает поток в режим потока. Если вы не хотите потреблять данные из потока, но хотите получить событие 'end', вы можете вызвать stream.resume(), чтобы открыть поток данных.
var readable = getReadableStreamSomehow();
readable.resume();
readable.on('end', () => {
console.log('got to the end, but did not read anything');
});
readable.setEncoding(encoding)
-
encoding<Строка> Кодировка для использования. - Возвращает:
this
Вызовите эту функцию, чтобы поток возвращал строки указанной кодировки вместо объектов Buffer. Например, если вы сделаете readable.setEncoding('utf8'), данные будут интерпретированы как данные UTF-8 и возвращены в виде строк. Если вы сделаете readable.setEncoding('hex'), данные будут закодированы в шестнадцатеричном формате.
Это корректно обрабатывает многобайтовые символы, которые в противном случае могли бы быть повреждены, если бы вы просто извлекли Buffer и вызвали buf.toString(encoding) на них. Если вы хотите читать данные как строки, всегда используйте этот метод.
Также можно отключить любую кодировку с помощью readable.setEncoding(null). Этот подход очень полезен, если вы работаете с двоичными данными или с большими многобайтовыми строками, распределёнными по нескольким частям.
var readable = getReadableStreamSomehow();
readable.setEncoding('utf8');
readable.on('data', (chunk) => {
assert.equal(typeof chunk, 'string');
console.log('got %d characters of string data', chunk.length);
});
readable.unpipe([destination])
-
destination<stream.Writable> Дополнительный конкретный поток для отмены подписки
Этот метод удалит хуки, созданные предыдущим вызовом stream.pipe().
Если пункт назначения не указан, все подписки отменяются.
Если пункт назначения указан, но для него не установлена подписка, это ничего не делает.
var readable = getReadableStreamSomehow();
var writable = fs.createWriteStream('file.txt');
// All the data from readable goes into 'file.txt',
// but only for the first second
readable.pipe(writable);
setTimeout(() => {
console.log('stop writing to file.txt');
readable.unpipe(writable);
console.log('manually close the file stream');
writable.end();
}, 1000);
readable.unshift(chunk)
Это полезно в определённых случаях, когда поток потребляется анализатором, которому необходимо «отменить» некоторые данные, которые он оптимистично извлёк из источника, чтобы поток можно было передать другой стороне.
Обратите внимание, что stream.unshift(chunk) не может быть вызвано после срабатывания события 'end'; будет поднята ошибка выполнения.
Если вы обнаружите, что часто используете stream.unshift(chunk) в своих программах, рассмотрите возможность реализации потока Transform вместо этого. (См. API для разработчиков потоков.)
// Pull off a header delimited by \n\n
// use unshift() if we get too much
// Call the callback with (error, header, stream)
const StringDecoder = require('string_decoder').StringDecoder;
function parseHeader(stream, callback) {
stream.on('error', callback);
stream.on('readable', onReadable);
var decoder = new StringDecoder('utf8');
var header = '';
function onReadable() {
var chunk;
while (null !== (chunk = stream.read())) {
var str = decoder.write(chunk);
if (str.match(/\n\n/)) {
// found the header boundary
var split = str.split(/\n\n/);
header += split.shift();
var remaining = split.join('\n\n');
var buf = new Buffer(remaining, 'utf8');
if (buf.length)
stream.unshift(buf);
stream.removeListener('error', callback);
stream.removeListener('readable', onReadable);
// now the body of the message can be read from the stream.
callback(null, header, stream);
} else {
// still reading the header.
header += str;
}
}
}
}
Обратите внимание, что в отличие от stream.push(chunk), stream.unshift(chunk) не завершит процесс чтения, не сбросив внутреннее состояние чтения потока. Это может привести к непредсказуемым результатам, если unshift() вызывается во время чтения (т. е. из реализации stream._read() в пользовательском потоке). Вызов unshift() с последующим немедленным вызовом stream.push('') корректно сбросит состояние чтения, однако лучше просто избегать вызова unshift() во время выполнения чтения.
readable.wrap(stream)
-
stream<Поток> Поток «старого стиля» для чтения
Версии Node.js до v0.10 имели потоки, которые не реализовывали весь API потоков в его современном виде. (См. Совместимость для получения дополнительной информации.)
Если вы используете старую библиотеку Node.js, которая отправляет события 'data' и имеет метод stream.pause(), который является только рекомендательным, то вы можете использовать метод wrap() для создания потока Readable, который использует старый поток как источник данных.
Вам очень редко придётся вызывать эту функцию, но она существует для удобства взаимодействия со старыми программами и библиотеками Node.js.
Например:
const OldReader = require('./old-api-module.js').OldReader;
const Readable = require('stream').Readable;
const oreader = new OldReader;
const myReader = new Readable().wrap(oreader);
myReader.on('readable', () => {
myReader.read(); // etc.
});
Класс: stream.Transform
Потоки преобразования — это потоки Duplex, где выходной сигнал каким-то образом вычисляется из входного. Они реализуют интерфейсы Readable и Writable.
Примеры потоков преобразования включают:
Класс: stream.Writable
Интерфейс потока Writable — это абстракция пункта назначения, в который вы записываете данные.
Примеры потоков Writable включают:
- HTTP-запросы на стороне клиента
- HTTP-ответы на стороне сервера
- потоки записи fs
- потоки zlib
- потоки crypto
- сокеты TCP
- stdin дочернего процесса
-
process.stdout,process.stderr
Событие: 'close'
Срабатывает, когда поток и любые связанные с ним ресурсы (например, дескриптор файла) закрыты. Событие указывает, что больше событий не будет отправлено, и дальнейшие вычисления не будут выполняться.
Не все потоки отправляют событие 'close', так как событие 'close' является необязательным.
Событие: 'drain'
Если вызов stream.write(chunk) возвращает false, то событие 'drain' укажет, когда можно начать запись дополнительных данных в поток.
// Write the data to the supplied writable stream one million times.
// Be attentive to back-pressure.
function writeOneMillionTimes(writer, data, encoding, callback) {
var i = 1000000;
write();
function write() {
var ok = true;
do {
i -= 1;
if (i === 0) {
// last time!
writer.write(data, encoding, callback);
} else {
// see if we should continue, or wait
// don't pass the callback, because we're not done yet.
ok = writer.write(data, encoding);
}
} while (i > 0 && ok);
if (i > 0) {
// had to stop early!
// write some more once it drains
writer.once('drain', write);
}
}
}
Событие: 'error'
Срабатывает, если при записи или передаче данных произошла ошибка.
Событие: 'finish'
Когда метод stream.end() был вызван, и все данные были отправлены в подлежащую систему, это событие отправляется.
var writer = getWritableStreamSomehow();
for (var i = 0; i < 100; i ++) {
writer.write('hello, #${i}!\n');
}
writer.end('this is the end\n');
writer.on('finish', () => {
console.error('all writes are now complete.');
});
Событие: 'pipe'
-
src<stream.Readable> поток-источник, который подключается к этому потоку для записи
Отправляется всякий раз, когда метод stream.pipe() вызывается на потоке для чтения, добавляя этот поток для записи в его набор пунктов назначения.
var writer = getWritableStreamSomehow();
var reader = getReadableStreamSomehow();
writer.on('pipe', (src) => {
console.error('something is piping into the writer');
assert.equal(src, reader);
});
reader.pipe(writer);
Событие: 'unpipe'
-
src<Readable Поток> поток-источник, который отменил подписку на этот поток для записи
Отправляется всякий раз, когда метод stream.unpipe() вызывается на потоке для чтения, удаляя этот поток для записи из его набора пунктов назначения.
var writer = getWritableStreamSomehow();
var reader = getReadableStreamSomehow();
writer.on('unpipe', (src) => {
console.error('something has stopped piping into the writer');
assert.equal(src, reader);
});
reader.pipe(writer);
reader.unpipe(writer);
writable.cork()
Принудительно буферизует все записи.
Буферизованные данные будут сброшены либо при вызове stream.uncork(), либо при вызове stream.end().
writable.end([chunk][, encoding][, callback])
Вызовите этот метод, когда больше данных в поток записываться не будет. Если указана, функция обратного вызова прикрепляется в качестве слушателя события 'finish'.
Вызов stream.write() после вызова stream.end() вызовет ошибку.
// write 'hello, ' and then end with 'world!'
var file = fs.createWriteStream('example.txt');
file.write('hello, ');
file.end('world!');
// writing more now is not allowed!
writable.setDefaultEncoding(encoding)
-
encoding<Строка> Новое кодирование по умолчанию
Устанавливает кодирование по умолчанию для потока на запись.
writable.uncork()
Очистить все данные, буферизованные с момента вызова stream.cork().
writable.write(chunk[, encoding][, callback])
Этот метод записывает данные в базовую систему и вызывает предоставленный обратный вызов, когда данные будут полностью обработаны. Если произойдет ошибка, обратный вызов может или не может быть вызван с ошибкой в качестве первого аргумента. Чтобы обнаружить ошибки записи, подпишитесь на событие 'error'.
Значение возврата указывает, следует ли продолжить запись прямо сейчас. Если данные должны были быть буферизованы внутри, то оно вернёт false. В противном случае оно вернёт true.
Это значение возврата является строго рекомендательным. Вы МОЖЕТЕ продолжить запись, даже если оно вернёт false. Однако записи будут буферизованы в памяти, поэтому лучше этого не делать чрезмерно. Вместо этого дождитесь события 'drain' перед записью дополнительных данных.
API для разработчиков потоков
Чтобы реализовать любой тип потока, шаблон одинаков:
- Расширьте соответствующий родительский класс в собственном подклассе. (Метод
util.inherits()особенно полезен для этого.) - Вызовите соответствующий конструктор родительского класса в конструкторе, чтобы убедиться, что внутренние механизмы настроены должным образом.
- Реализуйте один или несколько конкретных методов, как подробно описано ниже.
Класс для расширения и методы для реализации зависят от типа класса потока, который вы пишете:
| Сценарий использования | Класс | Реализуемые методы |
|---|---|---|
| Только чтение | ||
| Только запись | ||
| Чтение и запись | ||
| Обработать записанные данные, затем считать результат |
В вашем коде реализации очень важно никогда не вызывать методы, описанные в API для потребителей потоков. В противном случае вы можете потенциально вызвать побочные эффекты в программах, которые используют ваши потоковые интерфейсы.
Класс: stream.Duplex
Поток «дуплекс» — это поток, который одновременно является потоком на чтение и запись, например, соединение TCP-сокета.
Обратите внимание, что stream.Duplex — это абстрактный класс, предназначенный для расширения с помощью базовой реализации методов stream._read(size) и stream._write(chunk, encoding, callback), как вы делали бы с классом потока на чтение или запись.
Поскольку JavaScript не поддерживает множественное прототипное наследование, этот класс прототипно наследуется от Readable, а затем паразитно от Writable. Таким образом, от пользователя требуется реализовать как низкоуровневый метод stream._read(n), так и низкоуровневый метод stream._write(chunk, encoding, callback) в расширяемых классах дуплекса.
new stream.Duplex(options)
-
options<Объект> Передаётся как в конструктор Writable, так и в конструктор Readable. Также содержит следующие поля:-
allowHalfOpen<Булево> По умолчанию =true. Если установлено вfalse, поток автоматически завершит сторону чтения при завершении стороны записи и наоборот. -
readableObjectMode<Булево> По умолчанию =false. УстанавливаетobjectModeдля стороны чтения потока. Не имеет эффекта, еслиobjectModeравноtrue. -
writableObjectMode<Булево> По умолчанию =false. УстанавливаетobjectModeдля стороны записи потока. Не имеет эффекта, еслиobjectModeравноtrue.
-
В классах, которые расширяют класс Duplex, убедитесь, что вызываете конструктор, чтобы настройки буферизации могли быть должным образом инициализированы.
Класс: stream.PassThrough
Это тривиальная реализация потока Transform, который просто передает входные байты в выходной. Его назначение в основном для примеров и тестирования, но иногда есть случаи использования, когда он может пригодиться как строительный блок для новых типов потоков.
Класс: stream.Readable
stream.Readable — это абстрактный класс, предназначенный для расширения с помощью базовой реализации метода stream._read(size).
См. API для потребителей потоков, чтобы узнать, как потреблять потоки в своих программах. Далее следует объяснение того, как реализовать потоки Readable в ваших программах.
new stream.Readable([options])
-
options<Объект>-
highWaterMark<Число> Максимальное количество байтов для хранения во внутреннем буфере перед прекращением чтения из базового ресурса. По умолчанию =16384(16 КБ) или16для потоковobjectMode -
encoding<Строка> Если указано, буферы будут декодированы в строки с использованием указанного кодирования. По умолчанию =null -
objectMode<Булево> Указывает, должен ли этот поток вести себя как поток объектов. Это означает, чтоstream.read(n)возвращает одно значение вместо буфера размером n. По умолчанию =false -
read<Функция> Реализация методаstream._read().
-
В классах, которые расширяют класс Readable, убедитесь, что вызываете конструктор Readable, чтобы настройки буферизации могли быть должным образом инициализированы.
readable._read(size)
-
size<Число> Количество байтов для асинхронного чтения
Примечание: реализуйте этот метод, но не вызывайте его напрямую.
Этот метод имеет префикс подчеркивания, потому что он является внутренним для класса, который его определяет, и его должен вызывать только внутренние методы класса Readable. Все реализации потоков Readable должны предоставлять метод _read для извлечения данных из базового ресурса.
Когда вызывается _read(), если данные доступны из ресурса, реализация _read() должна начать помещать эти данные в очередь чтения, вызвав this.push(dataChunk). _read() должен продолжать читать из ресурса и помещать данные, пока push не вернёт false, в этот момент он должен прекратить чтение из ресурса. Только когда _read() будет вызван повторно после остановки, он должен начать чтение больше данных из ресурса и помещать эти данные в очередь.
Примечание: после вызова метода _read() он не будет вызван повторно до тех пор, пока не будет вызван метод stream.push().
Аргумент size рекомендательный. Реализации, где «чтение» — это единичный вызов, который возвращает данные, могут использовать его для определения количества данных для извлечения. Реализации, где это не имеет значения, такие как TCP или TLS, могут игнорировать этот аргумент и просто предоставлять данные по мере их появления. Например, нет необходимости «ждать», пока size байта не будут доступны перед вызовом stream.push(chunk).
readable.push(chunk[, encoding])
Примечание: Этот метод должен вызываться реализаторами Readable, а НЕ потребителями потоков Readable.
Если передано значение, отличное от null, метод push() добавляет чанк данных в очередь для последующего использования обработчиками потока. Если передано null, это сигнализирует об окончании потока (EOF), после чего больше данных записать нельзя.
Данные, добавленные с помощью push(), можно извлечь, вызвав метод stream.read(), когда произойдёт событие 'readable'.
Этот API разработан для максимальной гибкости. Например, вы можете обернуть источник более низкого уровня, имеющий механизм приостановки/возобновления и обратный вызов для данных. В этих случаях вы можете обернуть объект источника низкого уровня, сделав что-то вроде этого:
// source is an object with readStop() and readStart() methods,
// and an `ondata` member that gets called when it has data, and
// an `onend` member that gets called when the data is over.
util.inherits(SourceWrapper, Readable);
function SourceWrapper(options) {
Readable.call(this, options);
this._source = getLowlevelSourceObject();
// Every time there's data, we push it into the internal buffer.
this._source.ondata = (chunk) => {
// if push() returns false, then we need to stop reading from source
if (!this.push(chunk))
this._source.readStop();
};
// When the source ends, we push the EOF-signaling `null` chunk
this._source.onend = () => {
this.push(null);
};
}
// _read will be called when the stream wants to pull more data in
// the advisory size argument is ignored in this case.
SourceWrapper.prototype._read = function(size) {
this._source.readStart();
};
Пример: Поток подсчёта
Это базовый пример потока Readable. Он генерирует числа от 1 до 1 000 000 в порядке возрастания, а затем завершается.
const Readable = require('stream').Readable;
const util = require('util');
util.inherits(Counter, Readable);
function Counter(opt) {
Readable.call(this, opt);
this._max = 1000000;
this._index = 1;
}
Counter.prototype._read = function() {
var i = this._index++;
if (i > this._max)
this.push(null);
else {
var str = '' + i;
var buf = new Buffer(str, 'ascii');
this.push(buf);
}
};
Пример: SimpleProtocol v1 (неэффективный)
Это похоже на функцию parseHeader описанную здесь, но реализованную как пользовательский поток. Кроме того, обратите внимание, что эта реализация не конвертирует входящие данные в строку.
Однако, это было бы лучше реализовано как поток Transform. См. SimpleProtocol v2 для лучшей реализации.
// A parser for a simple data protocol.
// The "header" is a JSON object, followed by 2 \n characters, and
// then a message body.
//
// NOTE: This can be done more simply as a Transform stream!
// Using Readable directly for this is sub-optimal. See the
// alternative example below under the Transform section.
const Readable = require('stream').Readable;
const util = require('util');
util.inherits(SimpleProtocol, Readable);
function SimpleProtocol(source, options) {
if (!(this instanceof SimpleProtocol))
return new SimpleProtocol(source, options);
Readable.call(this, options);
this._inBody = false;
this._sawFirstCr = false;
// source is a readable stream, such as a socket or file
this._source = source;
source.on('end', () => {
this.push(null);
});
// give it a kick whenever the source is readable
// read(0) will not consume any bytes
source.on('readable', () => {
this.read(0);
});
this._rawHeader = [];
this.header = null;
}
SimpleProtocol.prototype._read = function(n) {
if (!this._inBody) {
var chunk = this._source.read();
// if the source doesn't have data, we don't have data yet.
if (chunk === null)
return this.push('');
// check if the chunk has a \n\n
var split = -1;
for (var i = 0; i < chunk.length; i++) {
if (chunk[i] === 10) { // '\n'
if (this._sawFirstCr) {
split = i;
break;
} else {
this._sawFirstCr = true;
}
} else {
this._sawFirstCr = false;
}
}
if (split === -1) {
// still waiting for the \n\n
// stash the chunk, and try again.
this._rawHeader.push(chunk);
this.push('');
} else {
this._inBody = true;
var h = chunk.slice(0, split);
this._rawHeader.push(h);
var header = Buffer.concat(this._rawHeader).toString();
try {
this.header = JSON.parse(header);
} catch (er) {
this.emit('error', new Error('invalid simple protocol data'));
return;
}
// now, because we got some extra data, unshift the rest
// back into the read queue so that our consumer will see it.
var b = chunk.slice(split);
this.unshift(b);
// calling unshift by itself does not reset the reading state
// of the stream; since we're inside _read, doing an additional
// push('') will reset the state appropriately.
this.push('');
// and let them know that we are done parsing the header.
this.emit('header', this.header);
}
} else {
// from there on, just provide the data to our consumer.
// careful not to push(null), since that would indicate EOF.
var chunk = this._source.read();
if (chunk) this.push(chunk);
}
};
// Usage:
// var parser = new SimpleProtocol(source);
// Now parser is a readable stream that will emit 'header'
// with the parsed header data.
Класс: stream.Transform
Поток "преобразования" — это дуплексный поток, где вывод каким-либо образом связан с вводом, например, поток zlib или crypto.
Нет требования, чтобы размер вывода был равен размеру ввода, количество чанков или время их прибытия совпадало. Например, поток Hash будет иметь только один чанк вывода, который предоставляется при завершении ввода. Поток zlib будет генерировать вывод, который может быть намного меньше или намного больше, чем его вход.
Вместо реализации методов stream._read() и stream._write(), классы Transform должны реализовать метод stream._transform() и могут также реализовать метод stream._flush() (см. ниже).
new stream.Transform([options])
-
options<Объект> Передаётся как конструкторам Writable, так и Readable. Также содержит следующие поля:-
transform<Функция> Реализация методаstream._transform(). -
flush<Функция> Реализация методаstream._flush().
-
В классах, расширяющих класс Transform, убедитесь, что вызывается конструктор, чтобы настройки буферизации были правильно инициализированы.
События: 'finish' и 'end'
События 'finish' и 'end' соответственно из родительских классов Writable и Readable. Событие 'finish' срабатывает после вызова stream.end() и обработки всех чанков методом stream._transform(), а 'end' срабатывает после вывода всех данных, что происходит после вызова обратного вызова в stream._flush().
transform._flush(callback)
-
callback<Функция> Вызовите эту функцию (при необходимости с аргументом ошибки), когда вы закончите сброс оставшихся данных.
Примечание: Этот метод НЕ должен вызываться напрямую. Он может быть реализован дочерними классами и вызывается только внутренними методами класса Transform.
В некоторых случаях ваша операция преобразования может потребовать выдать немного больше данных в конце потока. Например, поток сжатия Zlib будет хранить некоторое внутреннее состояние, чтобы оптимально сжать вывод. Однако в конце он должен сделать все возможное с оставшимися данными, чтобы данные были полными.
В этих случаях вы можете реализовать метод _flush(), который будет вызываться в самом конце, после потребления всех записанных данных, но до выдачи события 'end' для сигнализации о завершении читаемой части. Как и с stream._transform(), вызовите transform.push(chunk) ноль или более раз, как необходимо, и вызовите callback при завершении операции сброса.
Этот метод имеет префикс подчеркивания, потому что он внутренний для определяющего его класса и не должен вызываться напрямую пользовательскими программами. Однако вы обязаны переопределить этот метод в своих расширяемых классах.
transform._transform(chunk, encoding, callback)
-
chunk<Буфер> | <Строка> Чанк, который необходимо преобразовать. Будет всегда буфером, если опцияdecodeStringsне была установлена наfalse. -
encoding<Строка> Если чанк — строка, то это тип кодировки. Если чанк — буфер, то это специальное значение — 'buffer', игнорируйте его в этом случае. -
callback<Функция> Вызовите эту функцию (при необходимости с аргументом ошибки и данными), когда вы закончите обработку предоставленного чанка.
Примечание: Этот метод НЕ должен вызываться напрямую. Он должен быть реализован дочерними классами и вызываться только внутренними методами класса Transform.
Все реализации потоков Transform должны предоставить метод _transform() для приема входных данных и вывода результатов.
_transform() должен выполнить всё необходимое в этом конкретном классе Transform, чтобы обработать записанные байты и передать их в читаемую часть интерфейса. Выполняйте асинхронные операции ввода-вывода, обрабатывайте данные и т. д.
Вызовите transform.push(outputChunk) ноль или более раз, чтобы сгенерировать вывод из этого входного чанка, в зависимости от того, сколько данных вы хотите вывести в результате этого чанка.
Вызовите функцию обратного вызова только тогда, когда текущий чанк полностью обработан. Обратите внимание, что в результате любого чанка входных данных может быть или не быть вывода. Если вы передадите второй аргумент в обратный вызов, он будет передан в метод push. Другими словами, следующие варианты эквивалентны:
transform.prototype._transform = function (data, encoding, callback) {
this.push(data);
callback();
};
transform.prototype._transform = function (data, encoding, callback) {
callback(null, data);
};
Этот метод имеет префикс подчеркивания, потому что он внутренний для определяющего его класса и не должен вызываться напрямую пользовательскими программами. Однако вы обязаны переопределить этот метод в своих расширяемых классах.
Пример: SimpleProtocol парсер v2
Пример здесь простого парсера протокола можно реализовать, просто используя более высокий уровень класса потока Transform, аналогично примерам parseHeader и SimpleProtocol
v1.
В этом примере, вместо предоставления входных данных в качестве аргумента, они передаются в парсер, что является более привычным подходом к потокам Node.js.
const util = require('util');
const Transform = require('stream').Transform;
util.inherits(SimpleProtocol, Transform);
function SimpleProtocol(options) {
if (!(this instanceof SimpleProtocol))
return new SimpleProtocol(options);
Transform.call(this, options);
this._inBody = false;
this._sawFirstCr = false;
this._rawHeader = [];
this.header = null;
}
SimpleProtocol.prototype._transform = function(chunk, encoding, done) {
if (!this._inBody) {
// check if the chunk has a \n\n
var split = -1;
for (var i = 0; i < chunk.length; i++) {
if (chunk[i] === 10) { // '\n'
if (this._sawFirstCr) {
split = i;
break;
} else {
this._sawFirstCr = true;
}
} else {
this._sawFirstCr = false;
}
}
if (split === -1) {
// still waiting for the \n\n
// stash the chunk, and try again.
this._rawHeader.push(chunk);
} else {
this._inBody = true;
var h = chunk.slice(0, split);
this._rawHeader.push(h);
var header = Buffer.concat(this._rawHeader).toString();
try {
this.header = JSON.parse(header);
} catch (er) {
this.emit('error', new Error('invalid simple protocol data'));
return;
}
// and let them know that we are done parsing the header.
this.emit('header', this.header);
// now, because we got some extra data, emit this first.
this.push(chunk.slice(split));
}
} else {
// from there on, just provide the data to our consumer as-is.
this.push(chunk);
}
done();
};
// Usage:
// var parser = new SimpleProtocol();
// source.pipe(parser)
// Now parser is a readable stream that will emit 'header'
// with the parsed header data.
Класс: stream.Writable
stream.Writable — это абстрактный класс, предназначенный для расширения с реализацией метода stream._write(chunk, encoding, callback).
См. API для потребителей потоков, чтобы узнать, как использовать потоки Writable в ваших программах. Следующее — объяснение того, как реализовать потоки Writable в ваших программах.
new stream.Writable([options])
-
options<Объект>-
highWaterMark<Число> Уровень буферизации, когдаstream.write()начинает возвращатьfalse. По умолчанию =16384(16 КБ) или16для потоковobjectMode. -
decodeStrings<Логическое> Нужно ли декодировать строки в буферы перед передачей их вstream._write(). По умолчанию =true. -
objectMode<Логическое> Является лиstream.write(anyObj)валидной операцией. Если установлено, можно писать произвольные данные, а не только данныеBuffer/String. По умолчанию =false. -
write<Функция> Реализация методаstream._write(). -
writev<Функция> Реализация методаstream._writev().
-
В классах, расширяющих класс Writable, убедитесь, что вызывается конструктор, чтобы параметры буферизации были правильно инициализированы.
writable._write(chunk, encoding, callback)
-
chunk<Буфер> | <Строка> Часть данных для записи. Всегда будет буфером, если параметрdecodeStringsне был установлен вfalse. -
encoding<Строка> Если chunk — строка, это тип кодировки. Если chunk — буфер, это специальное значение — 'buffer', его следует игнорировать. -
callback<Функция> Вызовите эту функцию (при необходимости с аргументом ошибки), когда обработка части данных завершена.
Все реализации потоков Writable должны предоставлять метод stream._write() для отправки данных на подлежащий ресурс.
Примечание: не следует вызывать эту функцию напрямую. Она должна быть реализована дочерними классами и вызываться только внутренними методами класса Writable.
Вызовите колбэк используя стандартный шаблон callback(error), чтобы сообщить о успешном завершении записи или возникновении ошибки.
Если флаг decodeStrings установлен в параметрах конструктора, то chunk может быть строкой, а не буфером, и encoding укажет тип строки. Это для поддержки реализаций с оптимизированной обработкой определённых кодировок строк. Если вы явно не установили параметр decodeStrings в false, то можете безопасно пропустить аргумент encoding, и предположить, что chunk всегда будет буфером.
Этот метод имеет префикс подчеркивания, потому что является внутренним для класса, его определяющего, и не должен вызываться непосредственно пользовательскими программами. Однако вы обязаны переопределить этот метод в своих дочерних классах.
writable._writev(chunks, callback)
Примечание: не следует вызывать эту функцию напрямую. Она может быть реализована дочерними классами и вызываться только внутренними методами класса Writable.
Эта функция полностью необязательна для реализации. В большинстве случаев она не нужна. Если реализована, она будет вызвана со всеми частями данных, которые находятся в очереди записи.
Упрощённый API конструктора
В простых случаях теперь есть возможность создать поток без наследования.
Это можно сделать, передав соответствующие методы в качестве параметров конструктора:
Примеры:
Duplex
var duplex = new stream.Duplex({
read: function(n) {
// sets this._read under the hood
// push data onto the read queue, passing null
// will signal the end of the stream (EOF)
this.push(chunk);
},
write: function(chunk, encoding, next) {
// sets this._write under the hood
// An optional error can be passed as the first argument
next()
}
});
// or
var duplex = new stream.Duplex({
read: function(n) {
// sets this._read under the hood
// push data onto the read queue, passing null
// will signal the end of the stream (EOF)
this.push(chunk);
},
writev: function(chunks, next) {
// sets this._writev under the hood
// An optional error can be passed as the first argument
next()
}
});
Readable
var readable = new stream.Readable({
read: function(n) {
// sets this._read under the hood
// push data onto the read queue, passing null
// will signal the end of the stream (EOF)
this.push(chunk);
}
});
Transform
var transform = new stream.Transform({
transform: function(chunk, encoding, next) {
// sets this._transform under the hood
// generate output as many times as needed
// this.push(chunk);
// call when the current chunk is consumed
next();
},
flush: function(done) {
// sets this._flush under the hood
// generate output as many times as needed
// this.push(chunk);
done();
}
});
Writable
var writable = new stream.Writable({
write: function(chunk, encoding, next) {
// sets this._write under the hood
// An optional error can be passed as the first argument
next()
}
});
// or
var writable = new stream.Writable({
writev: function(chunks, next) {
// sets this._writev under the hood
// An optional error can be passed as the first argument
next()
}
});
Потоки: Под капотом
Буферизация
Потоки Writable и Readable буферизуют данные во внутренней структуре, к которой можно получить доступ из _writableState.getBuffer() или _readableState.buffer соответственно.
Объём буферизуемых данных зависит от параметра highWaterMark , переданного в конструктор.
Буферизация в Readable потоках происходит, когда реализация вызывает stream.push(chunk). Если потребитель потока не вызывает stream.read(), данные будут находиться во внутренней очереди до момента потребления.
Буферизация в Writable потоках происходит при многократном вызове stream.write(chunk), даже если он возвращает false.
Цель потоков, особенно с методом stream.pipe(), заключается в ограничении буферизации данных приемлемыми уровнями, чтобы источники и назначения с различной скоростью не перегружали доступную память.
Совместимость со старыми версиями Node.js
В версиях Node.js до v0.10 интерфейс потока Readable был проще, но также менее мощный и менее полезный.
- Вместо ожидания вызова метода
stream.read(), события'data'начинали излучаться немедленно. Если вам нужно было выполнить ввод-вывод, чтобы решить, как обработать данные, то вам нужно было сохранить части данных в некотором типе буфера, чтобы они не потерялись. - Метод
stream.pause()был рекомендательным, а не гарантированным. Это означало, что вам всё ещё нужно было быть готовыми к получению событий'data', даже когда поток был в приостановленном состоянии.
В Node.js v0.10 был добавлен класс Readable. Для обратной совместимости со старыми программами Node.js, потоки Readable переключаются в «режим потока» при добавлении обработчика события 'data' или вызове метода stream.resume(). В результате, даже если вы не используете новый метод stream.read() и событие 'readable', вам больше не нужно беспокоиться о потере частей данных 'data'.
Большинство программ будут продолжать работать нормально. Однако это вводит особый случай в следующих условиях:
- Обработчик события
'data'не добавлен. - Метод
stream.resume()никогда не вызывается. - Поток не перенаправлен ни на какое приемное место назначения.
Например, рассмотрим следующий код:
// WARNING! BROKEN!
net.createServer((socket) => {
// we add an 'end' method, but never consume the data
socket.on('end', () => {
// It will never get here.
socket.end('I got your message (but didnt read it)\n');
});
}).listen(1337);
В версиях Node.js до v0.10 входящие данные сообщения просто отбрасывались. Однако в Node.js v0.10 и выше сокет останется приостановленным навсегда.
Решение в этом случае — вызвать метод stream.resume(), чтобы начать поток данных:
// Workaround
net.createServer((socket) => {
socket.on('end', () => {
socket.end('I got your message (but didnt read it)\n');
});
// start the flow of data, discarding it.
socket.resume();
}).listen(1337);
В дополнение к переключению потоков Readable в режим потока, потоки в стиле до v0.10 могут быть обернуты в класс Readable с помощью метода stream.wrap().
Режим объектов
Обычно потоки работают исключительно со строками и буферами.
Потоки в режиме объектов могут излучать другие общие значения JavaScript, отличные от буферов и строк.
Readable поток в режиме объектов всегда возвращает единственный элемент при вызове stream.read(size), независимо от аргумента размера.
Writable поток в режиме объектов всегда игнорирует аргумент encoding для stream.write(data, encoding).
Специальное значение null по-прежнему сохраняет своё значение в режиме объектов. То есть, для Readable потоков в режиме объектов, значение null в качестве возвращаемого значения от stream.read() указывает, что больше нет данных, и stream.push(null) будет сигнализировать об окончании данных потока (EOF).
Ни один поток в ядре Node.js не является потоком в режиме объектов. Эта модель используется только внешними библиотеками потоков.
Вы должны установить objectMode в конструкторе своего класса потока в объекте опций. Установка objectMode во время работы потока небезопасна.
Для Duplex потоков objectMode может быть установлено исключительно для стороны чтения или записи с помощью readableObjectMode и writableObjectMode соответственно. Эти параметры могут быть использованы для реализации парсеров и сериализаторов с потоками Transform.
const util = require('util');
const StringDecoder = require('string_decoder').StringDecoder;
const Transform = require('stream').Transform;
util.inherits(JSONParseStream, Transform);
// Gets \n-delimited JSON string data, and emits the parsed objects
function JSONParseStream() {
if (!(this instanceof JSONParseStream))
return new JSONParseStream();
Transform.call(this, { readableObjectMode : true });
this._buffer = '';
this._decoder = new StringDecoder('utf8');
}
JSONParseStream.prototype._transform = function(chunk, encoding, cb) {
this._buffer += this._decoder.write(chunk);
// split on newlines
var lines = this._buffer.split(/\r?\n/);
// keep the last partial line buffered
this._buffer = lines.pop();
for (var l = 0; l < lines.length; l++) {
var line = lines[l];
try {
var obj = JSON.parse(line);
} catch (er) {
this.emit('error', er);
return;
}
// push the parsed object out to the readable consumer
this.push(obj);
}
cb();
};
JSONParseStream.prototype._flush = function(cb) {
// Just handle any leftover
var rem = this._buffer.trim();
if (rem) {
try {
var obj = JSON.parse(rem);
} catch (er) {
this.emit('error', er);
return;
}
// push the parsed object out to the readable consumer
this.push(obj);
}
cb();
};
stream.read(0)
В некоторых случаях вам необходимо вызвать обновление механизмов чтения базового потока, не потребляя при этом никаких данных. В этом случае вы можете вызвать stream.read(0), которая всегда вернёт null.
Если внутренний буфер чтения находится ниже highWaterMark, и поток в данный момент не читает, то вызов stream.read(0) вызовет вызов низкого уровня stream._read().
Практически никогда нет необходимости делать это. Однако в некоторых частях внутренней работы Node.js, особенно в классах потоков Readable, встречаются такие случаи.
stream.push('')
Передача нулевого байтового символа или буфера (если не в режиме объектного режима) имеет интересный побочный эффект. Поскольку это вызов stream.push(), он завершит процесс reading. Однако он не добавляет никаких данных в буфер чтения, поэтому для пользователя нет доступных данных.
Очень редко есть случаи, когда у вас нет данных для предоставления сейчас, но потребитель вашего потока (или, возможно, другая часть вашего кода) знает, когда нужно снова проверить, вызвав stream.read(0). В таких случаях вы можете вызвать stream.push('').
На данный момент единственный случай использования этой функциональности — в классе tls.CryptoStream, который устарел в Node.js/io.js v1.0. Если вам необходимо использовать stream.push(''), рассмотрите другой подход, поскольку это практически всегда указывает на серьёзную ошибку.
© Joyent, Inc. and other Node contributors
Licensed under the MIT License.
Node.js is a trademark of Joyent, Inc. and is used with its permission.
We are not endorsed by or affiliated with Joyent.
https://nodejs.org/dist/latest-v4.x/docs/api/stream.html