Node.js封装Writable流代理pipe数据时报write after end错误
Writable流封装代理触发
write after end错误解决方案 问题现象
- 为Writable流创建封装代理后,通过
pipe()向代理传输数据时持续触发write after end错误。 - 测试发现代理文件流时,将代理的
_write方法转发调用原始流的_write方法可修复问题,但HTTP请求对象不存在_write方法,无法套用该方案。
可复现问题的完整示例
const http = require("http"); const fs = require("fs"); const { Writable, Readable } = require('stream'); class MyStream extends Readable { constructor() { super(); this._iterator = this.generator(); } *generator() { for(const c of "Hello") yield c; } _read(n) { const it = this._iterator.next(); if (it.done) this.push(null); else this.push(it.value); } } class Proxy extends Writable { constructor(req) { super(); this._node_req = req; } end(chunk, encoding, cb) { console.trace("end"); this._node_req.end(chunk, encoding, cb); return {}; // return something } _write(chunk, encoding, cb) { console.log("_write", chunk.toString("utf8")); return this._node_req.write(chunk, encoding, cb); } } const src = new MyStream() const req = http.request("http://httpbingo.org/post", { method: "POST" }, (res) => { res.pipe(process.stderr); }); src.pipe(new Proxy(req));
运行输出日志
sh$ node --version v14.17.1 sh$ node t.js _write H Trace: end at Proxy.end (/home/sylvain/Projects/getpro/t.js:33:13) at MyStream.onend (internal/streams/readable.js:665:10) at Object.onceWrapper (events.js:481:28) at MyStream.emit (events.js:375:28) at endReadableNT (internal/streams/readable.js:1317:12) at processTicksAndRejections (internal/process/task_queues.js:82:21) _write e events.js:352 throw er; // Unhandled 'error' event ^ Error [ERR_STREAM_WRITE_AFTER_END]: write after end at writeAfterEnd (_http_outgoing.js:694:15) at write_ (_http_outgoing.js:706:5) at ClientRequest.write (_http_outgoing.js:687:15) at Proxy._write (/home/sylvain/Projects/getpro/t.js:40:27) at doWrite (internal/streams/writable.js:377:12) at clearBuffer (internal/streams/writable.js:529:7) at onwrite (internal/streams/writable.js:430:7) at callback (internal/streams/writable.js:513:21) at afterWrite (internal/streams/writable.js:466:5) at onwrite (internal/streams/writable.js:446:7) Emitted 'error' event on ClientRequest instance at: at writeAfterEndNT (_http_outgoing.js:753:7) at processTicksAndRejections (internal/process/task_queues.js:83:21) { code: 'ERR_STREAM_WRITE_AFTER_END' }
问题根因
错误核心是违反了Node.js流的扩展规范:
- 自定义Writable流时,仅需要实现
_write、_final、_destroy这类下划线开头的内部方法,公共方法write、end由流基类实现,负责维护写入队列、背压控制、状态标记,直接重写公共end方法会完全跳过基类的缓冲区排空逻辑。 - 示例代码中重写
end后,源流结束时会立刻调用代理的end,进而立刻调用底层HTTP请求的end,但此时代理自身写入队列中缓存的剩余数据还没完成写入,后续队列消费时继续向已经关闭的HTTP请求写数据,就触发了write after end错误。 - 之前给文件流做代理时直接调用底层流
_write的写法属于歪打正着:文件流的内部_write方法不会校验流的结束状态,但这种写法不符合规范,遇到HTTP请求这类未暴露内部_write方法的流就完全无法使用。
正确实现方案
按照流的扩展规范实现代理即可:
- 不要重写任何公共的
write、end方法,所有自定义逻辑放在下划线开头的内部方法中。 - 流结束收尾逻辑放在
_final方法中实现,该方法会在代理自身所有缓冲的写入数据全部处理完成后才被调用,从根本上避免写完end还有数据待写入的问题。 _write方法中正常调用底层流的公共write方法即可,同时正确处理drain事件保证背压逻辑正常,额外透传底层流的错误事件避免未捕获异常。
修复后的Proxy实现代码:
class Proxy extends Writable { constructor(req) { super(); this._node_req = req; // 透传底层请求错误,避免未捕获异常 this._node_req.on('error', err => this.destroy(err)); } _write(chunk, encoding, cb) { console.log("_write", chunk.toString("utf8")); // 处理背压:底层缓冲区满时等待drain事件再回调 if (this._node_req.write(chunk, encoding)) { cb(); } else { this._node_req.once('drain', cb); } } _final(cb) { console.trace("所有数据写入完成,调用底层流end"); // 所有缓冲数据写完后才关闭底层流 this._node_req.end(cb); } }
该实现对所有符合Node.js流规范的Writable流都生效,不需要依赖底层流暴露内部_write方法,也不会出现write after end错误。
内容的提问来源于stack exchange,提问作者Sylvain Leroux
相关产品推荐
相关产品推荐

