Write a generic fan-in Merge for channels. What are the concurrency pitfalls?
Question 146HardGo 1.22 to 1.25
func Merge[T any](ctx context.Context, chans ...<-chan T) <-chan T {
out := make(chan T)
var wg sync.WaitGroup
for _, c := range chans {
wg.Go(func() { // Go 1.25; before that: wg.Add(1); go func() { defer wg.Done(); ... }()
for v := range c {
select {
case out <- v:
case <-ctx.Done():
return
}
}
})
}
go func() {
wg.Wait()
close(out)
}()
return out
}
a, b := make(chan int), make(chan int)
go func() { a <- 1; close(a) }()
go func() { b <- 2; close(b) }()
sum := 0
for v := range Merge(context.Background(), a, b) {
sum += v
}
fmt.Println(sum) // 3
Points to mention:
- The parameter type is
...<-chan T. Achan intis assignable to<-chan int, but a[]chan intcannot be passed as[]<-chan int, so a caller with a slice must convert each element. - Capturing
cin the goroutine is safe since Go 1.22, because each iteration has its own variable. Before 1.22 you had to passcas an argument. - Only the goroutine that knows every sender has finished may close
out. That is the job of thewg.Wait()goroutine. - Without the
ctx.Done()case, a consumer that stops reading early leaves the senders blocked forever: a goroutine leak. - Output order across inputs is not deterministic.
More on Generics
- Q144What is an "instantiation cycle"? Why does this recursive generic function fail to compile?
- Q145What does this print? (
%Tand reflection on generic types) - Q147How do you get a pointer to a literal value generically? Compare a
Ptr[T]helper with Go 1.26new(expr).