亲手测量堆积的字节
目标
亲自测量在快速生产者与慢速消费者之间内存是如何增长的,并确认 write() 的返回值与 drain 事件挡住了什么。
为什么重要
被报告为“内存泄漏”的东西,很多并不是在泄漏,而是在排队。把每秒产生 10MB 的一方与每秒接收 2MB 的一方直接连起来,剩下的 8MB 总要堆在某处,而这个地方就是这个进程的堆。
Node 的流提供了一个挡住它的旋钮——write() 的返回值。但这个值放着不管,什么事也不会发生。读取并停下来,是写入一方的责任。 少了这一行,代码照样运行,测试也照样通过,一到生产环境,就会因为一个大文件而让进程崩溃。
而且,遵守背压似乎会变慢,但测量下来并不是这样。因为消费者接收的速度反正是一样的。几乎没有什么损失,挡住的却很多。
步骤
- 在
/root/work/backpressure/flow.mjs中编写makeSink(options),并把这个版本的默认 highWaterMark 测出来,写入report.json的hwm。 - 用
fillUntilFalse(stream, chunk)数出缓冲区变满的位置。 - 用
writeAll(stream, chunks)遵守背压地写入。 - 用
floodNoWait(stream, chunks)无视返回值灌入,写入runs.flood。 - 用
writeAll发送同样的量,写入runs.paced,并把比率写入backpressure。 - 用
pipeThrough(readable, writable)做同样的事,写入runs.pipeline。 - 用
estimateQueueBytes和secondsUntil计算堆积速度。
参考
makeSink({hwm, delayMs})是把接收的量统计到receivedBytes、receivedChunks的 Writable。delayMs为 0 时用setImmediate,大于 0 时则延迟相应的时间再调用回调。writeAll和floodNoWait都返回{peakBuffered}。stream.writableLength就是那一刻堆积的字节数。- 第 4 步和第 5 步要用相同的量测量才能比较。不要把数据块预先做成数组,而要每次发送时重新生成——如果预先生成,还没开始测量就已经把内存用完了,而如果反复使用同一个缓冲区,只是堆积引用,内存不会增长。
rssDeltaMB是把process.memoryUsage().rss的差值以 MB 写下来。- 常见错误:用
readable.pipe(writable)就结束。它会遵守背压,但不会把错误向上传,所以即使消费者死了,生产者也会继续读取。
亲自测出这个版本的默认值
在 /root/work/backpressure/flow.mjs 中导出 makeSink({hwm, delayMs}),并把 node 以及 hwm.writableDefault、hwm.objectMode、hwm.readableDefault 测出来后写入 /root/work/backpressure/report.json。
请打印 new Writable({write(c, e, cb) { cb(); }}).writableHighWaterMark。它可能与你记得的数字不同。
makeSink 把接收的字节数和数据块数统计到 receivedBytes、receivedChunks,当 delayMs 大于 0 时,延迟相应的时间再调用回调。
write() 返回 false 的位置
导出 fillUntilFalse(stream, chunk)。调用 write() 直到出现 false,返回包含那次返回 false 的调用在内的次数。
write() 返回的是“把这个数据块收下放入之后,累积的量是否小于上限”。从与上限相等的那一刻起,就是 false。
所以在 64KB 的上限上写一次 64KB,第一次调用就已经是 false 了。把这个数字弄明白,下一步的等待规则才会显得自然。
停下来,等待 drain
导出 writeAll(stream, chunks)。按顺序全部写入,但当 write() 为 false 时,先等待 drain 再继续写,并返回 {peakBuffered}。
await new Promise(r => stream.once("drain", r)) 这一行就是背压的全部。
请用 once 而不是 on。如果每次都用 on 来挂,监听器会堆积起来并触发警告,而那最终也是内存。
无视返回值会堆积到什么程度
导出 floodNoWait(stream, chunks),无视返回值灌入 2000 个 64KB 的数据块,把 chunks、chunkBytes、hwm、peakBuffered、bufferedMB、rssDeltaMB、wallMs、receivedChunks 写入 runs.flood。
peakBuffered 是 stream.writableLength 的最大值。灌进去多少就原样堆积多少——上限什么也挡不住。
如果把数据块预先做成数组,还没开始测量,内存就已经用完了。请用生成器每次发送时再生成。如果 rssDeltaMB 几乎是 0,就是掉进了那个陷阱。
遵守背压发送同样的量
用 writeAll 发送相同的数据块个数和大小,把同样的项目写入 runs.paced,并在 backpressure 中写入 bufferRatio(flood 与 paced 的 peakBuffered 之比,第一位小数)、rssRatio(第一位小数)、timeRatio(paced 与 flood 的 wallMs 之比,第二位小数)。
堆积量相差几千倍,而 timeRatio 却在 1 附近。
因为消费者接收的速度无论哪种方式都相同。等待并不是慢——如果不等待,这个差距只是堆积在内存里而已。
失败时由谁来收拾
导出 pipeThrough(readable, writable)。使用 stream/promises 的 pipeline,全部流完则完成,任何一方失败则以那个错误拒绝。流过同样的量,把 chunks、chunkBytes、hwm、peakBuffered、wallMs、receivedChunks 写入 runs.pipeline。
pipe() 也会遵守背压。但消费者死了,什么都不会发生,所以生产者继续读取,错误在没有人听的地方炸开。
pipeline 在一方失败时会清理其余部分,并把错误一直传给调用方。评分器会故意接上一个会失败的消费者来确认这一点。
几秒之后会爆炸
导出 estimateQueueBytes({producerBps, consumerBps, seconds}) 和 secondsUntil({producerBps, consumerBps, limitBytes})。当消费快于或等于生产时,分别是 0 和 Infinity。如果传入的不是数字,就抛出异常。
堆积的速度是 max(0, 생산 - 소비)(占位符依次为生产速度、消费速度),距离上限的剩余时间是 한도 / 그 속도(占位符依次为上限、堆积速度)。算术就是全部。
有了这两行,就能说“照现在的速度,4 分钟后就会碰到上限”,而不是“内存在往上涨”。容量规划就是从这里开始的。