TT Lab
Get started
Learn Learning paths Courses

It Wasn't One Request - Everything Got Slow

Measure the Bytes That Pile Up

Continue in TT Lab

Goal

Measure for yourself how memory grows between a fast producer and a slow consumer, and confirm what the return value of write() and the drain event prevent.

Why it matters

A good share of what gets reported as a "memory leak" is not leaking but a line that is standing. If you simply connect something that produces 10MB per second to something that takes 2MB per second, the leftover 8MB has to pile up somewhere, and that somewhere is this process's heap.

Node's streams give you one handle to prevent this: the return value of write(). But if you leave that value alone, it does nothing. Reading it and stopping is the writer's job. When that one line is missing, the code runs fine, the tests pass, and in production a single large file kills the process.

You might think respecting backpressure would make things slower, but when you measure, it does not. The speed at which the consumer takes data is the same anyway. You lose almost nothing and what you block is large.

Steps

  1. In /root/work/backpressure/flow.mjs, create makeSink(options), and measure this version's default highWaterMark and write it to report.json under hwm.
  2. Count the point where the buffer fills with fillUntilFalse(stream, chunk).
  3. Write while respecting backpressure with writeAll(stream, chunks).
  4. Pour in data while ignoring the return value with floodNoWait(stream, chunks) and record it in runs.flood.
  5. Send the same amount with writeAll, record it in runs.paced, and write the ratios in backpressure.
  6. Do the same job with pipeThrough(readable, writable) and record it in runs.pipeline.
  7. Calculate the accumulation rate with estimateQueueBytes and secondsUntil.

Notes

Measure this version's defaults yourself

In /root/work/backpressure/flow.mjs, export makeSink({hwm, delayMs}), and in /root/work/backpressure/report.json write node and hwm.writableDefault·hwm.objectMode·hwm.readableDefault, measured rather than recalled.

Print new Writable({write(c, e, cb) { cb(); }}).writableHighWaterMark. It may differ from the number you remember.

makeSink counts the received bytes and number of chunks in receivedBytes·receivedChunks, and if delayMs is greater than 0 it calls the callback that much later.

The point where write() returns false

Export fillUntilFalse(stream, chunk). Until false comes back, call write(), and return the count including the call that returned false.

write() returns "after taking this chunk in, is the amount queued less than the limit?" From the moment it equals the limit, it is false.

So if you write 64KB once with a 64KB limit, the first call is already false. You need to get a feel for this number so that the waiting rule in the next step feels natural.

Stop and wait for drain

Export writeAll(stream, chunks). Write everything in order, but when write() returns false, wait for drain and then continue writing, and return {peakBuffered}.

The single line await new Promise(r => stream.once("drain", r)) is the whole of backpressure.

Do not use on; use once. If you attach with on every time, listeners pile up and a warning appears, and in the end that is memory too.

How far does it pile up if you ignore the return value

Export floodNoWait(stream, chunks), pour in 2000 chunks of 64KB while ignoring the return value, and record the following in runs.flood: chunks·chunkBytes·hwm·peakBuffered·bufferedMB·rssDeltaMB·wallMs·receivedChunks.

peakBuffered is the maximum of stream.writableLength. It piles up exactly as much as you pour in; the limit blocks nothing.

If you build the chunks into an array ahead of time, all the memory is already used before you start measuring. Create them with a generator each time you send. If rssDeltaMB is almost 0, you have fallen into that trap.

Send the same amount while respecting backpressure

Send the same number and size of chunks with writeAll, write the same fields to runs.paced, and in backpressure write bufferRatio (peakBuffered of flood/paced, to one decimal place), rssRatio (to one decimal place), and timeRatio (wallMs of paced/flood, to two decimal places).

The amount piled up differs by thousands of times, yet timeRatio is near 1.

That is because the speed at which the consumer takes data is the same either way. Waiting is not slow: if you do not wait, that difference simply piles up in memory.

Who cleans up when something fails

Export pipeThrough(readable, writable). Use the stream/promises function pipeline, so that it completes when everything has flowed through and rejects with that error if either side fails. Send the same amount through and record the following in runs.pipeline: chunks·chunkBytes·hwm·peakBuffered·wallMs·receivedChunks.

pipe() also respects backpressure. But when the consumer dies, nothing happens, so the producer keeps reading and the error blows up in a place no one is listening.

When one side fails, pipeline cleans up the rest and propagates the error to the caller. The grader attaches a consumer that fails on purpose to confirm this.

In how many seconds does it blow up

Export estimateQueueBytes({producerBps, consumerBps, seconds}) and secondsUntil({producerBps, consumerBps, limitBytes}). If consumption is faster than or equal to production, they return 0 and Infinity respectively. If a non-numeric value is passed in, throw an exception.

The accumulation rate is max(0, 생산 - 소비) (production minus consumption), and the time left until the limit is 한도 / 그 속도 (the limit divided by that rate). It is all arithmetic.

With these two lines you can say "at this rate we hit the limit in 4 minutes" instead of "memory is going up a bit." Capacity planning starts here.