← Concepts & practices
Concept Concurrency, scheduling, and delivery

Backpressure & queues

Decide what a full buffer does.

You already lean on buffers: a log call that returns straight away, a stream you read chunk by chunk, a job queue in front of a slow API. A buffer is a good idea for a burst. Let’s follow an app’s log lines into a buffer that keeps growing when the log service slows down, and decide what should happen when it’s full.

TypeScriptGo One log shipper, two implementations.

01 / The idea

A buffer for log lines is a fair start.

You’re building an app that ships its logs. Every log call adds a line to an in-memory buffer and returns straight away, and a shipper sends lines to a log service a few at a time. When a burst of requests writes a burst of lines, the buffer holds them and the shipper catches up.

Read the first bufferTypeScript · the version this lesson starts from
shipper.ts
// The first version: every log call appends, and the shipper sends from the front.
export class LineBuffer implements Buffer {
	private lines: Line[] = [];
	// Always accepted. A burst waits here until the shipper catches up.
	offer(line: Line): Decision {
		this.lines.push({ ...line });
		return 'accepted';
	}
	take(count: number): Line[] {
		return this.lines.splice(0, count);
	}
	snapshot(): Line[] {
		return this.lines.map((line) => ({ ...line }));
	}
}

Go’s version is the same buffer. Both languages meet again at BoundedBuffer in section 02.

Then the log service slows down. The app still writes 4 lines a second and only 1 ships. Every log call still succeeds, the buffer grows by 3 lines a second, and the line at the front keeps getting older. Nothing looks broken, and nothing tells the app anything is wrong. Memory just climbs.

A queue lets work wait. Backpressure is what happens when it’s full: a signal back to whoever produces the work, so they slow down, wait, or have the work refused. A buffer with no limit never sends that signal. Docker makes you choose: its default log delivery blocks the app, and in its non-blocking mode, “when the buffer is full, new messages will not be enqueued. Dropping messages is often preferred to blocking the log-writing process of an application.”

If you’ve read a streamed response in the browser, you’ve already used backpressure without writing it: response.body hands you the next chunk only when you call read(). Section 05 reads one, then builds a price board where the server can’t be slowed down at all.

02 / See the shape

Give the buffer a limit, and an answer for when it’s full.

The basic form is the bound: a full buffer answers wait or drop instead of growing. In the wild adds the app’s side, which has to obey that answer. At the call site runs one slowdown through all three buffers.

Both languages run the same model on a logical clock, with no network or timers, and pass the same 12 shared schedules.

The bound. A full buffer answers every offer with wait or drop, instead of growing to fit it.

TypeScriptReading
shipper.ts
// A bounded buffer answers every offer. When it's full, the answer is wait or drop.
export class BoundedBuffer implements Buffer {
	private lines: Line[] = [];
	private policy: 'wait' | 'drop';
	private capacity: number;
	constructor(policy: 'wait' | 'drop', capacity: number) {
		integer(capacity, 1, 16);
		this.policy = policy;
		this.capacity = capacity;
	}
	offer(line: Line): Decision {
		if (this.lines.length >= this.capacity) {
			return this.policy === 'wait' ? 'wait' : 'dropped';
		}
		this.lines.push({ ...line });
		return 'accepted';
	}
	take(count: number): Line[] {
		return this.lines.splice(0, count);
	}
	snapshot(): Line[] {
		return this.lines.map((line) => ({ ...line }));
	}
}
GoAlongside
shipper.go
// A bounded buffer answers every offer. When it's full, the answer is wait or drop.
type BoundedBuffer struct {
	lines    []Line
	policy   Policy
	capacity int
}

func NewBoundedBuffer(policy Policy, capacity int) *BoundedBuffer {
	bounded(capacity, 1, 16)
	return &BoundedBuffer{policy: policy, capacity: capacity}
}

func (b *BoundedBuffer) Offer(line Line) Decision {
	if len(b.lines) >= b.capacity {
		if b.policy == Wait {
			return Waiting
		}
		return Dropped
	}
	b.lines = append(b.lines, line)
	return Accepted
}
Reading the TypeScriptOne interface, a decision, no throw

LineBuffer and BoundedBuffer both implement Buffer, so Shipper doesn’t know which one it has. offer returns a Decision instead of throwing: a full buffer is a normal answer, not an error.

splice(0, count) removes lines from the front of the array, which copies the rest. That’s fine for eight lines; a buffer of thousands would use a ring buffer.

Reading the GoA buffered channel does this for you

In everyday Go, make(chan Line, 8) is a bounded buffer. A send waits when it’s full, which is the wait policy, and a select with a default case turns it into drop. The example uses slices so every step stays on one goroutine and can be replayed exactly.

Sameer Ajmani’s pipelines post adds the part that’s easy to forget about waiting senders: “Goroutines are not garbage collected; they must exit on their own.”

03 / Follow the lines

Watch the log service slow down.

Five steps, every bar from running the shipper you just read. A bar is how many lines are in the buffer at the end of that second. A dashed line marks the limit, a mark under a second means the app waited, and a number under it counts dropped lines. Before each step, guess where the extra lines go.

In Try it, set the rates yourself, then switch between the three buffers on the same seconds.

Backpressure

One slow log service, three buffers.

No boundno limit

Ships 4 a second; the app writes 8 a second for two seconds

in the buffer
2
shipped
34
dropped
0
not yet written
84
oldest line waited
0s
seconds the app waited
0

No bound, no bound. Ships 4 a second; the app writes 8 a second for two seconds. Buffered lines each second: s1: 2, s2: 2, s3: 8, s4: 12, s5: 10, s6: 8, s7: 6, s8: 4, s9: 2, s10: 2, s11: 2, s12: 2. 2 buffered, 34 shipped, 0 dropped, 84 not yet written, the app waited 0 seconds. The buffer peaks at 12 lines and drains back once the burst ends. This is what a buffer is for.

01/ 05
Absorb a burst of log lines

A buffer absorbs a burst.

The app writes 8 lines a second for two seconds while the shipper sends 4. The buffer peaks at 12 and drains once the burst ends.

Reduced motion: choose a scene to see its completed state.

Read this scene

The app writes 8 lines a second for two seconds while the shipper sends 4. The buffer peaks at 12 and drains once the burst ends.

No bound, no bound. Ships 4 a second; the app writes 8 a second for two seconds. Buffered lines each second: s1: 2, s2: 2, s3: 8, s4: 12, s5: 10, s6: 8, s7: 6, s8: 4, s9: 2, s10: 2, s11: 2, s12: 2. 2 buffered, 34 shipped, 0 dropped, 84 not yet written, the app waited 0 seconds. The buffer peaks at 12 lines and drains back once the burst ends. This is what a buffer is for.

Watch restarts when you return. Step through keeps your selected step. Try it starts a fresh run of 120 lines each time you open it.

What a bounded buffer buys you

Now put names on what you just watched. These are the words you’ll hear in a design review, and each one points at something on this page.

Memory with a ceiling
With no bound, 34 lines wait after twelve seconds and the count keeps rising. With a bound of 8, the buffer never holds more, however long the slowdown lasts.
A signal that reaches the producer
Waiting slows the app to the shipping rate: 98 lines not yet written, against 72. The slowdown shows up where something can act on it.
Loss you can count
Dropping keeps the app at full speed and counts all 26 lost lines, instead of hiding the problem in a buffer that grows.
Delay that stops growing
In the bounded buffer the oldest line has waited 7 seconds, and that stays put. With no bound it’s 8 seconds and rising.
A decision written down
offer answers accepted, wait, or dropped. What a full buffer does is a line of code someone can review.

The review words are backpressure, a bounded queue, load shedding for dropping work on purpose, and flow control for a producer that slows to the consumer’s pace. Section 08 covers what they cost.

04 / Try a decision

A bound that nobody waits on.

Waiting on a log call slows requests down, so someone makes offer async: when the buffer is full, it returns a promise that settles once there’s room. The log call doesn’t await it, so requests never wait. The change is in parked.ts, and the lesson’s tests count what it holds.

After ten seconds, how many log lines is the process holding?

The buffer still holds at most 8 lines. For ten seconds the app logs 6 lines a second and the shipper sends 2.

05 / Give it a real job

Choose the answer for each kind of work.

In a real app, the buffer sits between request handlers and a shipper that sends batches to a log service. The limit is usually in bytes rather than lines, and the right answer for a full buffer depends on what the lines are for.

Debug logs

Drop and count

Losing a few is better than slowing every request.

Audit events

Wait or refuse

Losing one isn’t acceptable, so the caller has to slow down or fail.

Live values

Keep the latest

A newer value replaces one that hasn’t been sent yet.

Fluent Bit, a log shipper, does the waiting version. When an input’s buffered data passes mem_buf_limit, it pauses that input and resumes once enough has been flushed. Its documentation warns that some inputs can still lose data while paused, and suggests buffering on disk when that matters.

The example has wait and drop; the frontend below adds keep-the-latest. It leaves out bytes, batching, retries, and disk. None of those change the question of what a full buffer does.

Build UIs?Every streamed response you read already has backpressure built in, and one day a WebSocket makes you choose what to throw away.

Where it already is in your components

A fetch response body is a ReadableStream. Each read() gives you the next chunk, and until you call it again the stream’s internal queue fills. MDN describes what happens then: the stream “sends a signal backwards through the chain to tell earlier transform streams (or the original source) to slow down delivery.”

The textbook panes read a streamed chat reply that way, through TextDecoderStream, which every major browser has supported since 2022. You don’t write the backpressure; you read one chunk at a time and let the stream set the pace.

When you have to own it

Now it’s a live price board over a WebSocket. MDN is direct: the WebSocket API “doesn’t support backpressure”, so when messages arrive faster than the page handles them, it buffers them, pegs the CPU, or both. You can’t make the server wait, so you choose what to keep.

Prices replace each other, so the board keeps only the newest price for each watched symbol and draws once per animation frame. The map never holds more entries than there are symbols, and browsers pause requestAnimationFrame in background tabs without anything piling up. For a chat or an order book, where every message matters, keep-the-latest would be the wrong answer.

latest.ts
// Prices replace each other, so a screen that can't keep up only needs the newest one per symbol.
export class LatestPrices {
	private pending = new Map<string, number>();
	private replaced = 0;

	// Called for every message. A newer price replaces one the screen hasn't shown yet.
	set(symbol: string, price: number): void {
		if (this.pending.has(symbol)) this.replaced++;
		this.pending.set(symbol, price);
	}

	// Called once per frame: the newest price for each symbol that changed since the last frame.
	take(): Map<string, number> {
		const changed = this.pending;
		this.pending = new Map();
		return changed;
	}

	get skipped(): number {
		return this.replaced;
	}
}

A chat reply read from a streamed response one chunk at a time, so the stream sets the pace. React and Svelte.

ReactAlready in your code
StreamedReply.tsx
import { useEffect, useState } from 'react';

type Status = 'reading' | 'done' | 'failed';

export function StreamedReply({ prompt }: { prompt: string }) {
	const [text, setText] = useState('');
	const [status, setStatus] = useState<Status>('reading');

	useEffect(() => {
		const controller = new AbortController();
		setText('');
		setStatus('reading');

		async function read() {
			const response = await fetch('/api/reply', {
				method: 'POST',
				body: JSON.stringify({ prompt }),
				signal: controller.signal
			});
			if (!response.ok || !response.body) throw new Error(`HTTP ${response.status}`);
			// Each read() pulls the next chunk. Until we ask again, the stream's queue fills
			// and it signals back up the chain to slow down.
			const reader = response.body.pipeThrough(new TextDecoderStream()).getReader();
			for (;;) {
				const { done, value } = await reader.read();
				if (done) return;
				setText((previous) => previous + value);
			}
		}

		read().then(
			() => setStatus('done'),
			() => {
				// Leaving the page aborts the fetch; that isn't a failure to show.
				if (!controller.signal.aborted) setStatus('failed');
			}
		);
		return () => controller.abort();
	}, [prompt]);

	return (
		<>
			<p aria-busy={status === 'reading'}>{text}</p>
			{status === 'failed' && <p role="alert">The reply stopped. Try again.</p>}
		</>
	);
}

06 / Recognize it elsewhere

Anywhere a buffer sits between something fast and something slow.

You’ve met all of these. For each one, find the limit and what a full buffer does.

Familiar buffers, their limits, and what happens when full
Where you’ve seen itThe bufferWhen it’s full
Docker container logsThe non-blocking log buffer, max-buffer-size, 1 MB by defaultNew messages are dropped. In the default blocking mode, the app’s writes wait.
Node.js streamsA writable stream’s bufferwrite() returns false, and the source pauses until 'drain'.
Go channelsmake(chan T, n)A send waits for room; a select with default skips instead.
Browser streamsA stream’s internal queueIt signals back up the chain to slow down.
A WebSocketThe page’s own memoryNothing tells the server. You choose what to keep.

A buffer that always has room isn’t using any of this. It becomes backpressure when a full buffer changes what the producer does.

07 / Already in your toolbox

Your tools already bound their buffers.

Three places to look. For each one, find the limit and who gets told.

Node.js · Backpressuring in streams

The guide to write() returning false and waiting for 'drain'. In its benchmark, the same compression job used about 88 MB of memory with backpressure and about 1.5 GB without it.

Read the guide ↗

Go blog · Pipelines and cancellation

Sameer Ajmani’s 2014 post on stages joined by channels. It calls choosing a buffer size from the number of values you expect “fragile”, and uses a done channel so downstream stages can tell senders to stop.

Read the post ↗

Fluent Bit · Backpressure

A production log shipper’s own account of pausing an input at mem_buf_limit, resuming after a flush, and when to buffer on disk instead.

Read the documentation ↗
A useful counterexample: a checkout requestWhen neither waiting nor dropping is an answer

An order isn’t a log line. If the order queue is full, dropping loses a sale, and waiting holds the customer’s request open until something times out.

Refuse it clearly instead: answer 503 with a Retry-After, so the client knows the order wasn’t taken and when to try again. That’s backpressure across the network, and Retry, backoff & idempotency covers the client’s side.

08 / The parts to watch

A limit moves the problem somewhere you can see it.

These are the places it still goes wrong.

Waiting moves the backlog upstream

A waiting app writes only as fast as lines ship. If that app is serving requests, they slow down too. Docker’s logging docs warn that apps “are likely to fail in unexpected ways when STDERR or STDOUT streams block.”

Parked writers are a queue too

Section 04’s promises, a goroutine blocked on a send, a request held open: each one keeps its work in memory. A bound on the buffer means little if the writers waiting on it aren’t bounded.

Dropped work has to be counted

A silent drop looks like a quiet system. Count what you drop and report it, or the gap in the logs is a mystery during the next incident.

Shipping again isn’t caught up

When shipping returns to the rate the app writes at, the backlog stays where it was. It only shrinks with spare capacity. Step 5 shows it: 16 lines while shipping matches writing, 8 after four seconds with 2 a second to spare.

David Yanacek’s Amazon Builders’ Library article puts the risk plainly: queueing meant to improve availability “can dramatically increase the recovery time after an outage.”

Depth hides age

Eight buffered lines look the same whether they arrived a second ago or a minute ago. Watch how long the oldest line has waited as well as how many are waiting.

Lines aren’t all the same size

A limit of 8 lines doesn’t bound memory if one line is a megabyte. Real buffers limit bytes, like Docker’s max-buffer-size.

09 / Make the call

What would you have to change tomorrow?

Give both designs a plausible change and follow the work it creates.

How a change affects an unbounded buffer and a bounded one
The changeNo limitA bounded buffer
A two-second burst of requestsAbsorbs it: 12 lines at the peak, back to its usual 2 five seconds later.A limit of 8 makes the app wait or drop lines in the busiest second.
The log service slows to 1 line a second34 lines after twelve seconds, and rising.Holds at 8. The app waits, or lines are dropped and counted.
The slowdown lasts an hourMemory grows the whole time.Memory stays put. The cost shows up as slower writes or counted drops.
Every line must be keptKept in memory, and all lost together if the process dies.Wait, or refuse the work. Never drop.
Someone makes the log call async without awaitingNo change: it never waited anyway.The bound stops working: 42 lines held after ten seconds.

Reach for a bounded buffer whenever the producer can stay faster than the consumer for longer than a burst, and decide in code what happens when it’s full. A log service that slows for a minute is the moment.

Keep the plain buffer for a short burst you know ends, like a batch job with a fixed number of lines.

The question I’d leave beside the code is: when this is full, who finds out, and what do they do?

10 / Take the idea with you

Explain the shipper without saying “backpressure.”

“The buffer holds up to eight lines. When it’s full, either the app waits until a line ships or the new line is thrown away and counted, and we chose which on purpose.” In a review, the words are backpressure, bounded queue, load shedding, and flow control.

Before moving on, jot down why the unbounded buffer looked healthy while it grew, why the async log call held 42 lines, and one buffer in your own code, a queue, a channel, or an array of pending work, along with what happens when it’s full.

Connections to follow nextRelated lessons

Take the shipper into your editor. Limit the buffer by bytes instead of lines, and see how one very long line changes what fits.

Back to Concepts & practices →