If a Go pipeline's consumer stops receiving halfway, what happens to the upstream stage goroutines?
answer
- nobody receives, so nobody proceeds
- a blocked goroutine is never collected
- the send needs an escape hatch
- select the send against cancellation
- defer cancel in the consumer
basics
~20 sThey block forever on their next send and never return, pinning every value they hold. The fix: thread a context through every stage and write each send as a select against cancellation, so an abandoned stage exits.
solid answer
~50 sA stage that finishes a value and evaluates `out <- v` blocks until something receives. If the consumer breaks out of its range loop, nothing ever receives: the last stage is parked on that send, the stage before it is parked on its own send, and every stage goroutine — plus everything it references — stays alive for the life of the process. Nothing panics and nothing is collected, because a blocked goroutine is a GC root. The construction that fixes it is to pass a `context.Context` (or a `done chan struct{}`) into every stage and write the send as `select { case out <- v: case <-ctx.Done(): return }`. The consumer does `defer cancel()`, so abandoning the chain closes `ctx.Done()`, each stage's select takes the return branch, its `defer close(out)` runs, and the whole chain unwinds.
code
go · 14 linesfunc resize(ctx context.Context, in <-chan Frame) <-chan Frame {
out := make(chan Frame)
go func() {
defer close(out)
for f := range in {
select {
case out <- shrink(f):
case <-ctx.Done():
return // abandoned: exit instead of parking on the send
}
}
}()
return out
}go deeper
Know that a channel send waits until a receiver takes the value, and that a goroutine parked on a send never finishes on its own. That is the whole reason a pipeline needs a way to be told to stop.
Be ready to write the cancellable stage: the send as one case of a select, ctx.Done() as the other, defer close of the output above it, and defer cancel in the consumer. Explain why cancelling unwinds the chain from either end.
Show how you find this in a live process — a goroutine profile filling with chan send frames at one stage's line — and how you keep it out of the codebase with a test that abandons the pipeline and asserts every stage returned.
Decide what cancellation means for the pipelines your teams expose: whether an abandoned consumer must drain in-flight work or may drop it, and whether stages take a context per call rather than storing one, so every service follows the same rule.
## What "abandoned" means A pipeline is built so that values are pulled through it: each stage can only make progress when the stage after it receives. The consumer at the end is what keeps that motion going. When the consumer stops early — a `break` on the first corrupt frame, an HTTP handler whose request was cancelled, a `return` on an error — the pull stops, and everything upstream stalls at the point where it was pushing. Concretely, with three stages over an unbuffered chain: 1. the encode stage evaluates `out <- v` and parks, because nobody is receiving; 2. the resize stage finishes its next frame, evaluates its own send, and parks; 3. the decode stage does the same one item later. All three goroutines are now blocked forever. Nothing has failed: a send with no receiver is not an error in Go, it is simply a wait. ## Why the runtime cannot clean this up A goroutine blocked on a channel operation is a root for the garbage collector, not garbage. The runtime has no way to prove that no future receive will arrive — the channel might still be referenced by code that has not run yet — so the goroutine, its stack, and every value reachable from it stay live. In a photo pipeline that means the decoded pixel buffers in flight are pinned too, so the leak is measured in megabytes, not just in goroutine count. Also worth stating clearly, because it is a common wrong guess: goroutines have no parent-child relationship. Returning from the function that started a stage does nothing to that stage. The runtime's deadlock detector will not help either — it only fires when **every** goroutine is asleep, and a real service always has others running. ## The fix: a cancellation signal every stage watches The stage takes a `context.Context` as its first parameter and selects its send against cancellation: ```go select { case out <- v: case <-ctx.Done(): return } ``` `ctx.Done()` is a channel that is closed when the context is cancelled, so every stage watching it is released at once — closing is a broadcast. On the return branch, the stage's `defer close(out)` fires, which also terminates the range loop of whatever is downstream, so the chain unwinds cleanly from any point. The caller owns the signal: ```go ctx, cancel := context.WithCancel(context.Background()) defer cancel() ``` The `defer cancel()` is not optional bookkeeping — it is the entire mechanism. A pipeline wired with a context whose `cancel` is never called leaks exactly as badly as one with no context at all, and so does a stage that accepts a `ctx` but never selects on `ctx.Done()` at the send. A `done chan struct{}` closed by the caller works identically and predates `context` in this pattern; use a context when the pipeline sits under a request or under some other cancellable operation, which is usually. ## Where to watch, and where not to bother The **send** always needs the select, because that is where a stage parks when its consumer is gone. The **receive** needs it only when the stage's input might never be closed — a generator reading from a network connection or a long-lived queue, say. When the input is another stage that always closes its output, `for range in` terminates on its own as soon as the upstream stage exits. A note on capacity: giving the channel between two stages a buffer does not fix this. It lets the producer run ahead by the buffer's capacity, then the send parks exactly as before. The block moves a few items later; the goroutine still never returns. ## Seeing it and proving it In a running process, the goroutine profile is the direct evidence: it dumps every live goroutine with its state and stack, so a leak shows as a growing population of goroutines sitting in `chan send` at one particular stage's line. A CPU profile shows nothing at all, because blocked goroutines burn no CPU, and the race detector shows nothing either — a goroutine politely blocked on a send is not a data race. The test worth writing is an abandonment test: build the chain, consume two frames, cancel, and then assert that every stage actually returned — wait for each stage's output channel to be closed and drained, with a timeout, so a stage that never exits fails the test deterministically. That is far better than sleeping and sampling a goroutine count, which is timing-dependent and flaky.
- Does giving the channel between two stages a buffer fix the leak?No. A buffer lets the producing stage run ahead by the channel's capacity, so the block happens a few items later; once the buffer is full the send parks exactly as before and the goroutine still never returns. Cancellation is what ends a stage; capacity only changes when the stage notices it has been abandoned.
- Does a stage watch for cancellation on its send, its receive, or both?The send always, because that is where a stage parks when its consumer is gone. Watch the receive too when the input might never be closed — a generator reading a network stream or a long-lived queue. When the input is another stage that always closes its output, `for f := range in` terminates on its own once that stage exits.
- How do you prove in a test that an abandoned pipeline actually stops?Build the chain, consume two frames, cancel, then assert every stage returned: wait for each stage's output channel to be closed and drained, with a timeout so a stuck stage fails rather than hangs. Avoid sleeping and sampling a goroutine count — that is timing-dependent and turns a real leak into a flaky test.
saying these in an interview costs you the question
- Says the blocked stage goroutines are garbage collected
- Thinks abandoning the consumer panics the senders
- Adds a bigger channel buffer and calls the leak fixed
- Passes a context but never selects on ctx.Done() at the send
- Creates a cancellable context and never calls cancel