Go

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. A chan int is assignable to <-chan int, but a []chan int cannot be passed as []<-chan int, so a caller with a slice must convert each element.
  • Capturing c in the goroutine is safe since Go 1.22, because each iteration has its own variable. Before 1.22 you had to pass c as an argument.
  • Only the goroutine that knows every sender has finished may close out. That is the job of the wg.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

All 36 Generics questions