Go

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

All 38 Channels & select questions