Lesson 16 / 25
Pipelines
Stages connected by channels.
Each stage: receive, transform, send, close
A pipeline is a series of stages connected by channels. Each stage is a function that receives values from an inbound channel, does some work, and sends results on an outbound channel that it owns and closes when its input is exhausted. The first stage (the source) has only an output; the last (the sink) only consumes. Because each stage closes its output, range loops in downstream stages end naturally. To make early exit safe, every stage also takes a ctx and selects on ctx.Done() around its sends, so that when the sink stops reading, cancelling the context unwinds all stages instead of leaving them blocked.
Stages connected by channels
Pipelines, fan-out/fan-in and worker pools compose goroutines and channels into bounded, cancellable data flows.
Generate, square, print
A three-stage pipeline with cancellation.
package main
import (
"context"
"fmt"
)
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 square(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() // unwinds every stage if we stop early
for v := range square(ctx, gen(ctx, 1, 2, 3, 4)) {
fmt.Println(v) // 1 4 9 16, in order
}
}An assembly line
Each station on an assembly line does one job and passes the part along a conveyor. When the line is shut down, every station stops, rather than parts piling up at one station forever.
Quick check: In a pipeline, which stage should close a given channel?
- The main function, always
- The stage that receives from it
- The stage that sends on it (its owner)
- No one; channels must never be closed
Answer
The stage that sends on it (its owner) — Owners close their outputs so downstream range loops terminate.