TT Lab
Get started
Learn Learning paths Courses

It Wasn't One Request - Everything Got Slow

Nothing Is Leaking - A Queue Is Forming

Continue in TT Lab

In one line

If you simply connect a fast producer to a slow consumer, the difference left over must pile up somewhere, and in Node that somewhere is this process's heap. The return value of write() is the handle for preventing that, but reading it and stopping is the writer's job.

Why this was needed

A good share of investigations that begin with "it looks like there is a memory leak" are not leaks. If you connect something that produces 10MB per second to something that takes 2MB per second, the leftover 8MB has to be somewhere, and that somewhere is the stream's internal buffer. If you take a heap snapshot, Buffers take up a huge amount, yet nowhere in the code is anything holding them. Nobody made anything leak; the line just got long.

The reason this incident is found especially late is that it never shows up with small input. With 2MB of test data, the job ends before the buffer even fills. It first shows up on the day a single 800MB file arrives in production, and that day the process hits its memory limit and dies. On Kubernetes it shows up as OOMKilled, it restarts, and the same request comes in again.

So this incident has to be caught in code review. There is one place to catch it: is the writer looking at the return value of write()? If it is not, then when that code blows up depends only on the input size, and we do not choose that size.

How it works

Node's Writable stream keeps a queue internally and has a baseline called highWaterMark. After putting a chunk in the queue, write() returns true if the amount queued is below the baseline, and false otherwise. To quote the official documentation, false means "this chunk was accepted, but please wait for the drain event before writing more" (stream documentation).

What matters is that false does not reject the write. If you ignore it and keep writing, the stream keeps accepting and queuing. The limit blocks nothing. Stopping is entirely the writer's responsibility, and that is what these two lines do.

for (const chunk of chunks) {
  if (!stream.write(chunk)) {
    await new Promise((resolve) => stream.once("drain", resolve));
  }
}

The default value of the baseline has varied between versions, so it is better to measure than to memorize. Measured on Node 22.11.0 in the lab image, both Writable and Readable are 65536 bytes (64KiB), and object mode is 16 objects. If you carry the 16384 from your old memory into your calculations, they will be off.

The difference is dramatic when you measure. On the same image, pouring in 2000 chunks of 64KB (125MB) while ignoring the return value leaves the full 125MB piled up in the buffer, and RSS grew by 131MB. Sending the same amount while waiting for drain, the maximum piled up was 65536 bytes, exactly the baseline. That is a difference of thousands of times.

And what about the total time? The side that respected backpressure was actually a little faster, because the speed at which the consumer takes data is the same either way. The intuition that "waiting makes it slower" does not hold here: if you do not wait, that difference goes into memory instead of time.

What it looks like in the field

Hand-written backpressure has one more blank: who cleans up when something fails. readable.pipe(writable) respects backpressure but does not propagate errors. If an error occurs on the consumer side, the producer does not know about it and keeps reading, and the error blows up in a place where no one is listening. Leftover file handles and memory come as a bonus.

The stream/promises function pipeline fills that gap. When one side fails, it cleans up the rest and propagates the error to the caller. When I actually measured with a failing consumer attached, pipeline rejected with that error and also cleaned up the producer, while the one joined with pipe waited forever without any word. That is why searching for pipe( in a code review is worth as much as searching for Sync.

Finally, the most useful thing in practice is that this story can be turned into arithmetic. If you know the production rate and the consumption rate, the amount that accumulates per second is their difference, and the time left until the limit is the limit divided by that difference. If you can say "at this rate we hit the limit in 4 minutes" instead of "memory is going up a bit," the conversation changes from guessing to planning.

What you will do in the next lab

You will start by measuring this version's default baselines yourself and writing them down. You will count the point where write() returns false, build a write that waits for drain, and send the same amount in two ways, measuring the buffered bytes, RSS, and total time side by side.

Then you will do the same job with pipeline, attach a consumer that fails on purpose, and check that the error propagates and the producer is cleaned up. Finally you will build a function that calculates the accumulation rate and the time left until the limit. The grader recomputes the ratios you wrote from the two records and compares them, and runs the function you built to see whether the same properties appear.