It Wasn't One Request - Everything Got Slow
Measure the Bytes That Pile Up
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
- In
/root/work/backpressure/flow.mjs, createmakeSink(options), and measure this version's default highWaterMark and write it toreport.jsonunderhwm. - Count the point where the buffer fills with
fillUntilFalse(stream, chunk). - Write while respecting backpressure with
writeAll(stream, chunks). - Pour in data while ignoring the return value with
floodNoWait(stream, chunks)and record it inruns.flood. - Send the same amount with
writeAll, record it inruns.paced, and write the ratios inbackpressure. - Do the same job with
pipeThrough(readable, writable)and record it inruns.pipeline. - Calculate the accumulation rate with
estimateQueueBytesandsecondsUntil.
Notes
makeSink({hwm, delayMs})is a Writable that counts what it receives inreceivedBytes·receivedChunks. IfdelayMsis 0, it calls the callback withsetImmediate; if greater than 0, it calls the callback that much later.writeAllandfloodNoWaitboth return{peakBuffered}.stream.writableLengthis the number of bytes queued at that moment.- Steps 4 and 5 must be measured with the same amount to be comparable. Do not build the chunks into an array ahead of time; create them fresh each time you send. If you build them ahead, everything is already used before you start measuring, and if you reuse the same buffer, only references pile up and memory does not grow.
rssDeltaMBis the difference ofprocess.memoryUsage().rsswritten in MB.- A common mistake: stopping at
readable.pipe(writable). It respects backpressure but does not propagate errors, so even if the consumer dies, the producer keeps reading.
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.