Explain the pipeline pattern and how to cancel it properly.
Question 209HardGo 1.22 to 1.25
A pipeline is a chain of stages. Each stage receives from an inbound channel, transforms the values, and sends on an outbound channel it owns and closes. The rules: every stage closes its output when its input is exhausted, and every send and receive also watches a cancellation signal so downstream consumers can quit early.
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 when main returns early
for v := range sq(ctx, sq(ctx, gen(ctx, 1, 2, 3, 4))) {
fmt.Println(v)
if v > 20 { break }
}
}
Fan-out: several goroutines read from the same input channel. Fan-in: merge their outputs using a WaitGroup-guarded close. The most common mistake is a stage that forgets the ctx.Done() case on its send, which leaks the whole upstream chain when the consumer stops.
More on Channels & select
- Q207What does this print? (len and cap of channels)
- Q208Channels or mutexes: how do you decide?
- Q210What happens if you send on a closed channel inside a
selectwith adefault? - Q211Does an unbuffered channel give a happens-before guarantee? What does the memory model say about channels?
- Q212Implement a generic "first response wins" (hedged request / race) function.
- Q213What does this print? (request/response over a channel of channels)