Build a cancellable pipeline. What causes goroutine leaks in pipelines and how do you prevent them?
Question 240HardGo 1.22 to 1.25
func gen(ctx context.Context, nums ...int) <-chan int {
out := make(chan int)
go func() {
defer close(out)
for _, n := range nums {
select {
case out <- n:
case <-ctx.Done():
return
}
}
}()
return out
}
func sq(ctx context.Context, in <-chan int) <-chan int {
out := make(chan int)
go func() {
defer close(out)
for n := range in {
select {
case out <- n * n:
case <-ctx.Done():
return
}
}
}()
return out
}
func main() {
ctx, cancel := context.WithCancel(context.Background())
defer cancel() // unblocks every stage if we exit early
for v := range sq(ctx, sq(ctx, gen(ctx, 1, 2, 3))) {
fmt.Println(v) // 1 16 81
}
}
The pipeline rules:
- Each stage owns and closes its own output channel.
- Each stage ranges over its input until the input is closed.
- Every send is guarded by
ctx.Done().
Leak cause: a downstream consumer stops reading (early return, break, or an error) while an upstream goroutine is blocked on out <- v. That goroutine and everything it references stay alive forever.
Detection: go.uber.org/goleak in tests, runtime.NumGoroutine in metrics, and the /debug/pprof/goroutine endpoint. Go 1.26 also added an experimental goroutine-leak profile.
Go 1.23 range-over-func iterators (iter.Seq) are often a simpler choice for sequential pipelines, because they need no goroutines at all.
More on Concurrency Patterns & sync
- Q238Implement fan-out / fan-in: a generic Merge that combines N channels into one.
- Q239How do nil channels behave in select, and how are they used to merge two channels?
- Q241Spot the goroutine leak.
- Q242How do you close a channel safely when there are multiple senders?
- Q243How do you implement a semaphore in Go? Compare a buffered channel with golang.org/x/sync/semaphore.
- Q244Explain errgroup: WithContext, SetLimit, TryGo. What are its semantics and gotchas?