不是在漏,是在排队
一句话总结
把快的生产者和慢的消费者直接连起来,剩下的差额就一定会堆积在某个地方,而在 Node 里,这个“某个地方”就是这个进程的堆。write() 的返回值就是挡住它的旋钮,但读取并停下来,是写入一方的责任。
为什么需要它
许多以“好像有内存泄漏”为开头的调查,其实并不是泄漏。把每秒产生 10MB 的一方与每秒接收 2MB 的一方连起来,剩下的 8MB 总得待在某处,而那个地方就是流的内部缓冲区。抓一个堆快照,会看到一大堆 Buffer,可代码里却没有任何地方抓着它们。没有人让它泄漏,只是队伍变长了。
这起事故之所以发现得特别晚,是因为小输入下永远不会出现。用 2MB 的测试数据,缓冲区还没满就结束了。直到生产环境里来了一个 800MB 的文件那天,它才第一次暴露,那天进程撞上内存上限而崩溃。如果是 Kubernetes,会被标记为 OOMKilled,重启,同样的请求再次进来。
所以这起事故要在代码评审中抓住。要抓的地方只有一个——写入方有没有在看 write() 的返回值? 如果没有,那段代码什么时候爆炸,就只取决于输入大小,而这个大小不是我们定的。
工作原理
Node 的 Writable 流内部有一个队列,并有一条叫 highWaterMark 的高水位线。write() 把数据块放进队列之后,如果累积的量少于水位线就返回 true,否则返回 false。照官方文档的说法,false 的意思是“这个数据块我收下了,但在继续写之前,请先等待 drain 事件”(流文档)。
重要的是,false 并不是拒绝写入。如果无视它继续写,流会继续接收并放进队列。上限什么也挡不住。停下来完全是写入一方的责任,这就是这两行。
for (const chunk of chunks) {
if (!stream.write(chunk)) {
await new Promise((resolve) => stream.once("drain", resolve));
}
}
水位线的默认值曾经随版本而变,所以不要背,最好去测。在实验镜像的 Node 22.11.0 上测量,Writable 和 Readable 都是 65536 字节(64KiB),对象模式是 16 个。如果沿用旧记忆里的 16384 来计算,就会出偏差。
差别一测就很惊人。在同一个镜像上,无视返回值,把 2000 个 64KB 的数据块(125MB)灌进去,缓冲区里就原封不动堆了 125MB,RSS 增加了 131MB。同样的量,如果一边等待 drain 一边发送,累积的最大值是 65536 字节,也就是正好停在水位线上。差了几千倍。
那么总耗时又如何呢?遵守背压的一方反而略快一点。因为无论哪种方式,消费者接收的速度都是一样的。“等待就会变慢”这个直觉在这里不成立——如果不等待,这个差距只是不体现在时间上,而是转到了内存上。
在现场相遇的样子
手写的背压还有一个空白。失败时由谁来收拾。readable.pipe(writable) 会遵守背压,但不会把错误向上传。如果消费者一侧出了错,生产者对此一无所知,继续读取,错误会在没有人听的地方炸开。文件句柄和内存残留下来,只是额外的附赠。
stream/promises 的 pipeline 填补了这个空白。一方失败,就清理其余部分,并把错误一直传给调用方。实际测量,接上一个会失败的消费者时,pipeline 会以那个错误拒绝,同时连生产者也一并清理,而用 pipe 接起来的一方,则无声无息地永远等下去。在代码评审中找 pipe(,与找 Sync 一样有价值,原因就在这里。
最后,这件事可以转换成算术,这一点在实务中最有用。知道了生产速度和消费速度,每秒堆积的量就是它们的差,距离上限的剩余时间,就是上限除以这个差。如果能说“照现在的速度,4 分钟后就会碰到上限”,而不是“内存在往上涨”,对话就从猜测变成了计划。
下一项实验要做什么
从亲手测出这个版本的默认水位线并写下来开始。数一数 write() 返回 false 的位置,写出等待 drain 的写入,用两种方式发送同样的量,并排测量累积字节数、RSS 和总耗时。
接着用 pipeline 做同样的事,并故意接上一个会失败的消费者,确认错误会向上传、生产者会被清理。最后写出计算堆积速度和距离上限剩余时间的函数。评分器会从两份记录里重新计算你写下的比率并对照,还会直接运行你写的函数,看是否出现同样的性质。