
128 MiB 数据,10 毫秒发完,进程的 RSS 涨了 132.6 MiB。这不是哪里配错了,是 write() 的返回值没被看。写流不会替上游减速,它只在缓冲区满的时候把 write() 的返回值改成 false,接不接这个信号由调用方决定。
实验装置
消费者是一个自定义 Writable,每收一块就 sleep 0.15 毫秒,highWaterMark 设成 64 KiB。上游每块 64 KiB,一共 2048 块,合起来 128 MiB。脚本在 tools/verify-backpressure.mjs,三种写法各自 fork 一个子进程跑,内存数互不干扰:
node v26.8.1 / 每块 64 KiB,共 2048 块 = 128 MiB
消费者是一个每块 sleep 0.15 ms 的 Writable(highWaterMark 64 KiB)
no-drain = 忽略 write() 返回值 drain = 返 false 就 await drain file = 写真实文件
label issue buffered peakRSS done
----------------------------------------------------
no-drain 10ms 128MiB +132.6MiB 2601ms
drain 2607ms 0MiB +29.1MiB 2610ms
file 11ms 128MiB +132.6MiB 35ms
表里四列分别是:循环把 write() 发完用了多久、那一刻还在缓冲区里等着的数据量(stream.writableLength)、这一轮采样到的峰值 RSS 增量、整批数据写完的总耗时。
三行的差别在哪
第一行是没接返回值的那种写法:
for (const chunk of chunks) {
writable.write(chunk) // 返回值丢掉
}
循环本身就是同步的,10 毫秒把 2048 次 write 全部提交完,一次都没等。消费者再慢也拦不住它,1024 块之外的数据只能排队躺在内存里,writableLength 停在 128 MiB 这个数上,正好等于总量:全部数据都在堆上,一块也没被消化。RSS 增了 132.6 MiB,比 128 MiB 多出来的部分,是缓冲区这个结构自己的开销。
第二行把返回值用起来:
import { once } from 'node:events'
import { finished } from 'node:stream/promises'
for (const chunk of chunks) {
if (!writable.write(chunk)) {
await once(writable, 'drain')
}
}
writable.end()
await finished(writable)
write() 返回 false 表示缓冲区到了 highWaterMark,此刻往下等 drain 事件,等消费者把水位降下去再继续。代价写在表里:循环阶段从 10 毫秒变成 2607 毫秒,因为上游被迫跟着消费者的节奏走。换来的是 writableLength 停在 0,峰值 RSS 只涨 29.1 MiB。总耗时几乎一样(2610 毫秒对 2601 毫秒),因为瓶颈本来就是那个 sleep 0.15 毫秒的消费者,上游快慢不改变总时长,只改变这段等待发生在内存里还是发生在 await 上。
第三行是最容易骗人的一组。写真实文件时,同样丢掉返回值,缓冲区里同样堆了 128 MiB,峰值 RSS 同样涨 132.6 MiB,只是本地磁盘消费得快,整批 35 毫秒就结束了。内存占用的行为其实一模一样,差别全在消费者有多快。所以「往文件写大对象从来没出过事」这句话推不出「不接返回值没问题」,它只说明没撞上下游堵住的时候。
关于那 29.1 MiB
接了 drain 也不是零。采样方式是每次 drain 之后记一次 RSS,取最大值,所以这个数含着已经分配过、还没被回收的 Buffer。128 MiB 的数据流过去,中途有上百兆的 Buffer 被创建又丢弃,回收跟不上节奏,RSS 就不会回到起点。这一点与结论无关:真正的区别是水位有没有被控制住,是一百多兆还是几十兆,不是零和一。
什么时候必须看这个返回值
消费者比生产者慢的情况在服务端很常见,往网卡写会慢,做压缩会慢,下游服务限速也会慢,任何一种都会把速度差变成内存。这些场景就在手边:日志往远端推送,大文件往对象存储上传,导出的 CSV 边查边发。这些路径上的写操作如果都忽略返回值,一次慢下游就会把整个进程的内存抬上去,而且现象是延迟出现的,看起来像是别处的内存泄漏。
控制水位有现成的写法。Readable 与 Writable 之间用 pipeline() 连接,回压由框架处理,上游是文件流或者数据库游标时都可以放心接:
import { pipeline } from 'node:stream/promises'
import { createReadStream, createWriteStream } from 'node:fs'
await pipeline(
createReadStream('big.log'),
slowTransform,
createWriteStream('out.log'),
)
手工循环要留神一处细节:别写成每块都等一次回调。await new Promise((r) => writable.write(chunk, r)) 这种写法语法上没错,但它把每一次写入都串成一次完整的往返,吞吐会掉到比消费者自己还低,因为中间还夹了一个事件循环的往返。判断标准只有一个:写不进去才等。这两条路线我倾向于走 pipeline,手工循环只留给必须逐块改内容的场合。
结语
一次写入要花多少内存,取决于写的时候缓冲区里已经积了多少,跟单块大小没关系。128 MiB 与 132.6 MiB 这两个数之所以贴得这么近,就是因为没人拦着生产者。这个差值的量级不吓人,吓人的是它随数据量线性增长:把 2048 块换成 20000 块,第一行的 RSS 增量就会跟着涨到一 GB 以上,第三行的耗时还是几十毫秒。写流把选择权交给了调用方,返回值就是那个选择。