ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Node.js Streams2 流 API 完全解析:v0.10 流接口重构的设计、兼容性与实现原理

Node.js Streams2 流 API 完全解析:v0.10 流接口重构的设计、兼容性与实现原理 Node.js Streams2 流 API 完全解析v0.10 流接口重构的设计、兼容性与实现原理【免费下载链接】nodejs.orgThe Node.js® Website项目地址: https://gitcode.com/GitHub_Trending/no/nodejs.org导读本文基于 nodejs.org 官方博客在 2012 年 12 月发布的经典技术文章《A New Streaming API for Node v0.10》收录于本仓库 apps/site/pages/en/blog/feature/streams2.md系统讲解 Node.js 历史上最重要的一次流接口重构——Streams2。它彻底改变了data事件立即发射与pause()仅具建议性两大设计缺陷引入了read()/push()/unshift()等拉取式pull-based读取模型并为 Readable、Writable、Duplex、Transform、PassThrough 五大基类奠定了沿用至今的架构。读完本文你将完整掌握 Streams2 的 API 细节、向后兼容策略与全部代码示例并从仓库中的 v0.9.4、v0.10.0 等发布公告里看到这一重构从预览到稳定落地的真实过程以及它如何影响今天 Node.js 的流实现包括现代 Web Streams API。背景v0.8 时代流接口的四大痛点从 Node 诞生之初团队就一直在逐步迭代理想的基于事件的 API。这一过程最终演化为贯穿 Node 核心模块与大量 npm 模块的Stream 接口。统一接口大大提升了程序与库的可移植性和可靠性从领域特有的事件和方法走向统一的流接口是巨大的进步。但截至 v0.8Node 的流仍然存在四个根本性问题pause()并不会真正暂停。它只是建议性的advisory-only实现上虽然简化了内部逻辑但对用户而言语义混乱做的并不是名字所表达的事。data事件会立即发射不管你是否准备好。这让先加载用户会话、再决定如何处理请求之类的常见任务变得异常困难——数据可能在处理逻辑就绪前就已到达。无法消费指定数量的字节再把剩余部分留给程序的其它部分处理。自己实现流极其困难。要把 pause、resume、写入缓冲、data事件的种种微妙之处全部做对门槛极高缺乏共享基类意味着每个人都要反复解决同样的问题、犯同样的错误。Common simple tasks should be easy, or we arent doing our job.常见简单任务本应容易否则就是我们失职。—— Isaac Z. Schlueter尽管 Node 常被评价为比其他平台更擅长流处理作者认为这与其说是褒奖不如说是对整个软件行业现状的控诉比隔壁家强还不够我们必须做到想象中最好。为什么拖到 v0.10 才动手兼容性约束修复的意愿一直存在但 Node 社区已积累了数年爆发式增长任何改动都必须极其谨慎——如果在 0.10 里弄坏了所有 Node 程序就没人愿意升级一切努力都将白费。这个讨论在 0.4、0.6、0.8 时代各发生过一次结论始终是工作量太大、难以向后兼容且总有更紧迫的问题待解决。到 0.10团队终于咬了子弹bit the bullet对流实现做出重大变更。这就是社区在 Twitter、IRC、邮件列表中热议的streams2。Streams2 概览与发布节奏Streams2 的核心结论如下即原文档的 tl;drNode 流很棒除了所有那些让人抓狂的地方0.10 将带来新的 Stream 实现绰号streams2Readable 流拥有read()方法返回一个 buffer 或 null下文有完整文档data事件、pause()、resume()依然照常工作只不过这次它们真的会按照你预期的方式工作旧程序几乎总是无需修改即可运行但流初始处于暂停状态paused state必须被读取才会被消费⚠️ 警告如果从不添加data事件处理器、也不调用resume()流将永远处于暂停状态永远不会发射end。版本落地时间线仓库佐证原文档指出第一个包含该变更的预览版是0.9.4并强烈建议开发者试用并反馈。仓库中的 v0.9.4 发布公告2012-12-21证实了这一点streams: Update all streaming interfaces to use new classes (isaacs)紧随其后的 v0.9.5 发布公告 则记录了http: Performance enhancements for http under streams2 (isaacs)可见 streams2 在预览阶段就持续针对 http 模块做性能增强。而 v0.10.0 稳定版发布公告2013-03-11进一步披露了落地细节Streams2 API 在开发过程中就被用于 npm 注册表中的模块当时已有37 个已发布模块依赖readable-stream库readable-stream这个 npm 包允许你在v0.8 旧代码库中使用新的 Stream 接口所有 Node 核心流都基于同一套易于扩展的基类构建行为一致性大幅提升用户态程序中创建流接口也比以往更容易性能方面Streams2 最初落地时曾造成明显性能回退尤其是 http 模块团队坚持Node 不允许在我们的主要用例上变慢的原则经过数月优化v0.10 最终在 HTTP 基准上总体更快大字符串消息场景约慢 1-5%其余持平或更快而fs.ReadStream吞吐量大幅提升、且几乎不再受 chunk size 影响如 buf size1024 时吞吐提升约 1726%。Stream 完整 API 文档原文收录稳定性2 - 不稳定Unstable流stream是 Node 中由各种对象实现的抽象接口。例如HTTP 服务器的请求是一个流stdout 也是。流可以是只读readable、只写writable或两者兼具。所有流都是 EventEmitter 的实例。通过require(stream)可以加载流基类其中提供了 Readable、Writable、Duplex、Transform 四类基类。兼容性Compatibility在更早版本的 Node 中Readable 流接口更简单但也更弱、更没用data事件不会等你调用read()就立即发射。如果需要在处理数据前先做 I/O就必须把块暂存到某种 buffer 中以防丢失。pause()是建议性的而非保证性的。即使流处于暂停状态你依然要准备好接收data事件。Node v0.10 新增了下面描述的 Readable 类。为向后兼容旧程序当添加data事件处理器、或调用pause()/resume()时Readable 流会切换进 old mode旧模式。其效果是即使你不使用新的read()方法和readable事件也无需再担心丢失data数据块。大多数程序将照常运行。但以下边缘情况需要特别留意没有添加data事件处理器从未调用pause()和resume()。考虑下面这段代码原文档明确标注为BROKEN! 错误示例// WARNING! BROKEN! net .createServer(function (socket) { // we add an end method, but never consume the data socket.on(end, function () { // It will never get here. socket.end(I got your message (but didnt read it)\n); }); }) .listen(1337);在 v0.10 之前的版本中传入的消息数据会被直接丢弃因此end能正常触发但在 v0.10 及之后socket 将永远保持暂停状态end处理器永远不会执行。解决方法是在此场景下调用resume()触发 old mode 行为// Workaround net .createServer(function (socket) { socket.on(end, function () { 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 之前的旧式流还可以通过wrap()方法包装进 Readable 类详见下文。Class: stream.ReadableReadable 流包含以下方法、成员与事件。注意stream.Readable是一个抽象类设计上要求子类实现底层的_read(size)方法见下文。new stream.Readable([options])options{Object}highWaterMark{Number} 内部缓冲区在停止从底层资源读取前可存储的最大字节数。默认值 16kbencoding{String} 若指定则缓冲区数据将按该编码解码为字符串。默认值 nullobjectMode{Boolean} 是否将流作为对象流处理即stream.read(n)返回单个值而不是大小为 n 的 Buffer。默认值 false在扩展 Readable 类的子类中务必调用构造函数以便正确初始化缓冲设置。readable._read(size)size{Number} 异步读取的字节数注意此函数不应被直接调用。它应由子类实现、且仅由内部 Readable 类方法调用。所有 Readable 流实现都必须提供_read方法来从底层资源获取数据。方法名以下划线开头因为它对定义它的类而言是内部方法不应被用户程序直接调用但你应当在自己的扩展类中覆写它。当数据可用时通过调用readable.push(chunk)把它放入读取队列。如果push返回 false则应停止读取当_read再次被调用时再开始推送更多数据。size参数是建议性的advisory。对于一次调用即返回数据的实现可以用它得知应获取多少数据而对于 TCP、TLS 这类场景则无需理会该参数只需在数据可用时提供即可——例如没有必要非得等到凑够size字节才调用stream.push(chunk)。readable.push(chunk)chunk{Buffer | null | String} 推入读取队列的数据块返回 {Boolean} 是否应继续执行更多推送注意此函数应由 Readable 的实现者调用而不是 Readable 子类的消费者调用。_read()在至少一次push(chunk)调用之前不会被再次调用。如果没有数据可用你可以调用push()空字符串来允许未来再次触发_read同时不向队列添加任何数据。Readable 类的工作原理是把数据放入读取队列待readable事件触发后由read()方法取出。push()显式向读取队列插入数据如果以null调用则表示数据结束EOF。在包装一个具有 pause/resume 机制和 data 回调的底层数据源时可以这样做// 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. var stream new Readable(); source.ondata function (chunk) { // if push() returns false, then we need to stop reading from source if (!stream.push(chunk)) source.readStop(); }; source.onend function () { stream.push(null); }; // _read will be called when the stream wants to pull more data in // the advisory size argument is ignored in this case. stream._read function (n) { source.readStart(); };readable.unshift(chunk)chunk{Buffer | null | String} 要前插到读取队列的数据块返回 {Boolean} 是否应继续执行更多推送这是readable.push(chunk)的对应物不是把数据放到读取队列末尾而是放到读取队列前端。这在流被解析器消费、需要把乐观读取的多余数据退回去的场景中非常有用。以下是一个简单数据协议的解析器示例——header 是一个 JSON 对象后跟 2 个\n字符然后是消息体原文档注明用 Transform 流实现更简单见下文// 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. See below. function SimpleProtocol(source, options) { if (!(this instanceof SimpleProtocol)) return new SimpleProtocol(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; var self this; source.on(end, function () { self.push(null); }); // give it a kick whenever the source is readable // read(0) will not consume any bytes source.on(readable, function () { self.read(0); }); this._rawHeader []; this.header null; } SimpleProtocol.prototype Object.create(Readable.prototype, { constructor: { value: SimpleProtocol }, }); SimpleProtocol.prototype._read function (n) { if (!this._inBody) { var chunk this._source.read(); // if the source doesnt have data, we dont 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); // 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.这段代码展示了 Streams2 的三大关键手法read(0)触发缓冲区刷新但不消费字节、push()在无数据时允许未来再次调用_read、unshift(b)把多余字节退回队列前端。readable.wrap(stream)stream{Stream} 一个 旧式old style可读流如果你在使用一个旧式 Node 库——它发射data事件、pause()只是建议性的——那么可以用wrap()方法创建以该旧流为数据源的 Readable 流var OldReader require(./old-api-module.js).OldReader; var oreader new OldReader(); var Readable require(stream).Readable; var myReader new Readable().wrap(oreader); myReader.on(readable, function () { myReader.read(); // etc. });Event: readable当有数据可供消费时触发。事件发射后应调用read()方法来消费数据。Event: end当流收到 EOF在 TCP 术语中即 FIN时触发表示不再会有data事件。如果流同时也是可写的可能仍可以继续写入。Event: datadata事件默认发射Buffer若调用了setEncoding()则发射字符串。注意添加data事件监听器会把 Readable 流切换进 old mode——数据一旦可用就立即发射而不是等你调用read()来消费。Event: error接收数据时发生错误则触发。Event: close当底层资源例如底层的文件描述符被关闭时触发。并非所有流都会发射此事件。readable.setEncoding(encoding)让data事件发射字符串而非Buffer。encoding可以是utf8、utf16leucs2、ascii或hex。编码也可以在构造函数中通过encoding字段指定。readable.read([size])size{Number | null} 可选要读取的字节数返回 {Buffer | String | null}注意此函数应由 Readable 流的使用者调用。在readable事件发射后调用它以消费数据。size参数设定你感兴趣的最小字节数若不设置则返回内部缓冲区的全部内容。如果无数据可消费或内部缓冲区中的字节数少于size参数则返回null并在有更多数据可用时再次发射readable事件。调用stream.read(0)总是返回null但会触发内部缓冲区刷新除此之外无其它作用。readable.pipe(destination, [options])destination{Writable Stream}options{Object} 可选end{Boolean}默认值 true把此可读流连接到destination写入流本流上的数据会写入目标。它正确管理背压back-pressure避免慢速目标被快速可读流压垮。此函数返回destination流。模拟 Unixcat命令process.stdin.pipe(process.stdout);默认情况下当源流发射end时会在目标上调用end()使destination不再可写。传入{ end: false }可保持目标流打开reader.pipe(writer, { end: false }); reader.on(end, function () { writer.end(Goodbye\n); });注意无论指定何种选项process.stderr和process.stdout在进程退出前都不会被关闭。readable.unpipe([destination])destination{Writable Stream} 可选撤销之前建立的pipe()。若未提供 destination则移除所有已建立的管道。readable.pause()将可读流切换进 old mode——通过data事件发射数据而不是缓冲起来供read()方法消费。暂停数据流暂停状态下不会发射data事件。readable.resume()将可读流切换进 old mode——通过data事件发射数据而不是缓冲起来供read()方法消费。在pause()之后恢复data事件。Class: stream.WritableWritable 流包含以下方法、成员与事件。注意stream.Writable是抽象类设计上要求子类实现底层的_write(chunk, encoding, cb)方法见下文。new stream.Writable([options])options{Object}highWaterMark{Number}write()开始返回 false 时的缓冲区水位。默认值 16kbdecodeStrings{Boolean} 是否在把字符串传给_write()前解码为 Buffer。默认值 true在扩展 Writable 类的子类中务必调用构造函数以便正确初始化缓冲设置。writable._write(chunk, encoding, callback)chunk{Buffer | String} 要写入的数据块。除非decodeStrings选项设为false否则始终是 Bufferencoding{String} 若 chunk 是字符串则为编码类型若 chunk 是 Buffer 则忽略。注意除非显式将decodeStrings设为falsechunk始终是 Buffercallback{Function} 处理完给定 chunk 后调用可带错误参数所有 Writable 流实现都必须提供_write方法来向底层资源发送数据。注意此函数绝对不能直接调用。它应由子类实现、且仅由内部 Writable 类方法调用。使用标准的callback(error)模式回调以表明写入成功完成或出错。若构造函数选项中设置了decodeStringschunk 可能是字符串而非 Bufferencoding会指明字符串类型——这是为了支持对特定字符串编码有优化处理的实现。如果你没有显式把decodeStrings设为false就可以安全地忽略encoding参数并假定 chunk 始终是 Buffer。writable.write(chunk, [encoding], [callback])chunk{Buffer | String} 要写入的数据encoding{String} 可选。若 chunk 是字符串编码默认为utf8callback{Function} 可选。此 chunk 成功写入后调用返回 {Boolean}把chunk写入流。数据已刷新到底层资源时返回true返回false表示缓冲区已满数据将在未来发送drain事件会指示缓冲区何时再次清空。write()何时返回 false由构造函数提供的highWaterMark决定。writable.end([chunk], [encoding], [callback])chunk{Buffer | String} 可选最后要写入的数据encoding{String} 可选。若 chunk 是字符串编码默认为utf8callback{Function} 可选。最后的 chunk 成功写入后调用调用此方法以表示写入流的数据结束。Event: drain当流的写入队列清空、可以安全地无缓冲写入时触发。当stream.write()返回false时应监听它。Event: close当底层资源例如底层文件描述符被关闭时触发。并非所有流都会发射此事件。Event: finish调用end()且没有更多 chunk 要写入时触发。Event: pipesource{Readable Stream}当流被传给某个可读流的 pipe 方法时触发。Event: unpipesource{Readable Stream}当使用源 Readable 流的unpipe()方法撤销之前建立的pipe()时触发。Class: stream.Duplex双工duplex流同时是 Readable 和 Writable例如 TCP socket 连接。注意stream.Duplex是抽象类设计上要求子类像实现 Readable / Writable 那样实现底层的_read(size)和_write(chunk, encoding, callback)方法。由于 JavaScript 没有多重原型继承此类原型继承自 Readable再从 Writable 寄生继承。因此扩展双工类时需要同时实现底层_read(n)方法与_write(chunk, encoding, cb)方法。new stream.Duplex(options)options{Object} 传给 Writable 和 Readable 构造函数。另有以下字段allowHalfOpen{Boolean}默认值 true。若设为false则可写侧结束时自动结束可读侧反之亦然。在扩展 Duplex 类的子类中务必调用构造函数以便正确初始化缓冲设置。Class: stream.Transform转换transform流是一种输出与输入存在因果联系的双工流例如 zlib 流或 crypto 流。输出不要求与输入等大小、等块数或同时到达。例如Hash 流在输入结束时只会产生一个输出块zlib 流的输出可能远小于或远大于其输入。与实现_read()和_write()不同Transform 类必须实现_transform()方法并可选择实现_flush()方法见下文。new stream.Transform([options])options{Object} 传给 Writable 和 Readable 构造函数。在扩展 Transform 类的子类中务必调用构造函数以便正确初始化缓冲设置。transform._transform(chunk, encoding, callback)chunk{Buffer | String} 要转换的数据块。除非decodeStrings设为false否则始终是 Bufferencoding{String} 若 chunk 是字符串则为编码类型若 chunk 是 Buffer 则忽略callback{Function} 处理完给定 chunk 后调用可带错误参数注意此函数绝对不能直接调用。它应由子类实现、且仅由内部 Transform 类方法调用。所有 Transform 流实现都必须提供_transform方法以接收输入并产生输出。_transform应执行该 Transform 类特有的处理把写入的字节处理好并交给接口的可读部分——可以做异步 I/O、处理数据等。调用transform.push(outputChunk)0 次或多次根据你想从该输入块产生多少输出来决定。只有当当前 chunk 被完全消费时才调用 callback。注意特定输入块可能产生输出也可能不产生。transform._flush(callback)callback{Function} 完成剩余数据刷新后调用可带错误参数注意此函数绝对不能直接调用。子类可以MAY实现它若实现则仅由内部 Transform 类方法调用。某些情况下你的转换操作可能需要在流结束时再多发射一些数据。例如Zlib压缩流会存储内部状态以优化输出压缩率结束时它需要尽最大努力处理剩余数据使输出完整。此时可实现_flush方法它会在所有已写数据被消费之后、可读侧发射end之前被调用。与_transform一样酌情调用transform.push(chunk)0 次或多次并在刷新操作完成时调用callback。示例用 Transform 重写 SimpleProtocol 解析器上文基于unshift()的协议解析器用更高层的 Transform 流可以大幅简化。此版本不再把输入作为构造参数传入而是把数据pipe 进解析器——这是更地道的 Node 流用法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 Object.create(Transform.prototype, { constructor: { value: SimpleProtocol }, }); 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(b); } } else { // from there on, just provide the data to our consumer as-is. this.push(b); } done(); }; var parser new SimpleProtocol(); source.pipe(parser); // Now parser is a readable stream that will emit header // with the parsed header data.对比两种实现可以看出 Streams2 的架构理念底层_read/_write处理字节搬运的复杂细节高层_transform只需聚焦输入块 → 输出块的纯转换逻辑。Class: stream.PassThrough这是Transform流的一个平凡实现——把输入字节原样传给输出。它的主要用途是示例和测试偶尔也会在实际场景中派上用场。从仓库看 Streams2 的历史回响Streams2 的设计并非终点而是现代 Node.js 流的起点。仓库中的后续发布公告可以佐证其深远影响v0.10.0 发布公告 明确写道所有 Node 核心流都基于同一套易于扩展的基类构建且 Streams2 API 在开发期间就已被 npm 生态中的模块使用并通过readable-stream包回移植到 v0.8 老代码库——这正是它兼容旧世界、面向新世界设计哲学的实证。近十年后的 v18 发布公告 中Web Streams APIReadableStream、WritableStream、TransformStream及各类 controller、queuing strategy随 V8 升级进入 Node——这是流思想在标准平台层的延续。v21 发布公告 则记载了 streams 团队持续优化 Writable 和 Readable 流的工作包括由 streams 维护者 Robert Nagy 主导的、通过移除冗余检查等方式对流的进一步优化——证明从 streams2 确立的 Readable/Writable 抽象至今仍是 Node 性能优化与 API 演进的核心主线。如果你希望完整阅读原始 API 文档含全部参数与示例可直接查看仓库中的 streams2.md想了解该 API 最终随稳定版发布的表现与基准数据可阅读 v0.10.0 发布公告。小结Streams2 给今天的我们留下的实践要点拉取优于推送readableread()的拉取式消费让读多少、何时读由消费者掌控read(0)可刷新缓冲区而不消费字节read(size)支持按需读取最小字节数。暂停语义要真实可靠pause()/resume()从建议性变为保证性这是流正确性的基石。兼容性靠显式模式切换添加data监听或调用pause()/resume()会切回 old modewrap()可包装旧式流。但务必牢记——既不监听data也不调用resume()的流会永久暂停、永不发射end。背压由框架管理pipe()会正确管理背压防止慢目标被快源压垮write()返回 false 时应等待drain再继续写。实现自定义流只需覆写底层钩子Readable 覆写_read、Writable 覆写_write、Transform 覆写_transform可选_flush复杂细节全部由基类处理——这正是 streams2 解决人人重复造轮子问题的方式。理解 Streams2 的这段历史与 API 设计不仅有助于阅读今天 Node.js 的流文档与源码也能帮助你写出更符合流式思维、更善用背压的高质量 Node 程序。【免费下载链接】nodejs.orgThe Node.js® Website项目地址: https://gitcode.com/GitHub_Trending/no/nodejs.org创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表