Pipelines, Select, and Cancel
Pipelines, Select, and Cancel
Overview
Pipelines compose stages connected by channels. The hard part is not leaking goroutines when a consumer stops early, and merging multiple streams correctly.
Leak: early exit without cancel
func rangeGen(start, stop int) <-chan int {
out := make(chan int)
go func() {
for i := start; i < stop; i++ {
out <- i // blocks if nobody receives
}
close(out)
}()
return out
}
// Consumer breaks early → producer stuck on send forever.
for v := range rangeGen(41, 46) {
fmt.Println(v)
if v == 42 {
break
}
}sequence (top → bottom):
actors: rangeGen, consumer
rangeGen --> consumer : 41
rangeGen --> consumer : 42
consumer --> consumer : break
note: LEAKED G
Goroutines are cheap but not free; leaks accumulate.
Cancel channel + select
func rangeGen(cancel <-chan struct{}, start, stop int) <-chan int {
out := make(chan int)
go func() {
defer close(out)
for i := start; i < stop; i++ {
select {
case out <- i:
case <-cancel:
return
}
}
}()
return out
}
func main() {
cancel := make(chan struct{})
defer close(cancel) // signal on exit
for v := range rangeGen(cancel, 41, 46) {
fmt.Println(v)
if v == 42 {
break
}
}
}select rules (teaching model):
1. run a case that can proceed
2. if several ready → choose at random (fairness)
3. if none ready → wait (or take default if present)
regions: select
flow:
[out left-arrow i] --consumer reading--> [RunA]
[left-arrow cancel] --cancel closed--> [Exit]
Cancel vs done
| Name | Direction | Meaning |
|---|---|---|
| done | worker → waiter | “I finished” |
| cancel | waiter → worker | “Stop now” |
Both are often named done in the wild—be explicit in your APIs.
Merge sequentially (slow)
// Reads in1 fully, then in2 — second producer stalls.
for v := range in1 {
out <- v
}
for v := range in2 {
out <- v
}Merge concurrently (WaitGroup)
func merge(in1, in2 <-chan int) <-chan int {
out := make(chan int)
var wg sync.WaitGroup
forward := func(in <-chan int) {
defer wg.Done()
for v := range in {
out <- v
}
}
wg.Add(2)
go forward(in1)
go forward(in2)
go func() {
wg.Wait()
close(out)
}()
return out
}Merge with select + nil disable
func mergeSelect(in1, in2 <-chan int) <-chan int {
out := make(chan int)
go func() {
defer close(out)
for in1 != nil || in2 != nil {
select {
case v, ok := <-in1:
if !ok {
in1 = nil // disable case
continue
}
out <- v
case v, ok := <-in2:
if !ok {
in2 = nil
continue
}
out <- v
}
}
}()
return out
}Why nil?
closed channel is always ready with ok=false
without nil, select spins forever on closed cases
nil channel cases never select → clean exit when both nil
Fan-out / fan-in pipeline sketch
producer ──chan──► stage ──chan──► consumer
│ │
└── cancel closed early ──► select exits (no leak)
Each stage: func(in <-chan T) <-chan U starting its own goroutine, closing out on exit, accepting cancel when needed.
Runnable example
go mod init example && go run .package main
import "fmt"
func rangeGen(cancel <-chan struct{}, start, stop int) <-chan int {
out := make(chan int)
go func() {
defer close(out)
for i := start; i < stop; i++ {
select {
case out <- i:
case <-cancel:
return
}
}
}()
return out
}
func main() {
cancel := make(chan struct{})
defer close(cancel)
for v := range rangeGen(cancel, 41, 46) {
fmt.Println(v)
if v == 42 {
break
}
}
fmt.Println("consumer done; producer exits via cancel")
}What to notice: Early break no longer leaks the generator.
Try next: Implement mergeSelect for three channels by generalizing the nil pattern.