You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.28 07:48:56