TT Lab
开始
学习 学习路径 课程

慢的不是一个请求,而是全部

亲手测量堆积的字节

在 TT Lab 中继续学习

目标

亲自测量在快速生产者与慢速消费者之间内存是如何增长的,并确认 write() 的返回值与 drain 事件挡住了什么。

为什么重要

被报告为“内存泄漏”的东西,很多并不是在泄漏,而是在排队。把每秒产生 10MB 的一方与每秒接收 2MB 的一方直接连起来,剩下的 8MB 总要堆在某处,而这个地方就是这个进程的堆。

Node 的流提供了一个挡住它的旋钮——write() 的返回值。但这个值放着不管,什么事也不会发生。读取并停下来,是写入一方的责任。 少了这一行,代码照样运行,测试也照样通过,一到生产环境,就会因为一个大文件而让进程崩溃。

而且,遵守背压似乎会变慢,但测量下来并不是这样。因为消费者接收的速度反正是一样的。几乎没有什么损失,挡住的却很多。

步骤

  1. 在 /root/work/backpressure/flow.mjs 中编写 makeSink(options),并把这个版本的默认 highWaterMark 测出来,写入 report.json 的 hwm。
  2. 用 fillUntilFalse(stream, chunk) 数出缓冲区变满的位置。
  3. 用 writeAll(stream, chunks) 遵守背压地写入。
  4. 用 floodNoWait(stream, chunks) 无视返回值灌入,写入 runs.flood。
  5. 用 writeAll 发送同样的量,写入 runs.paced,并把比率写入 backpressure。
  6. 用 pipeThrough(readable, writable) 做同样的事,写入 runs.pipeline。
  7. 用 estimateQueueBytes 和 secondsUntil 计算堆积速度。

参考

亲自测出这个版本的默认值

在 /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 分钟后就会碰到上限”,而不是“内存在往上涨”。容量规划就是从这里开始的。