Implement fan-out / fan-in: a generic Merge that combines N channels into one.
Question 238MediumGo 1.22 to 1.25
Fan-out means several goroutines read from the same input channel to spread the work. Fan-in means merging several channels into one.
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() {
for {
select {
case v, ok := <-c:
if !ok {
return
}
select {
case out <- v:
case <-ctx.Done():
return
}
case <-ctx.Done(): // also stop while idle on the input
return
}
}
})
}
go func() {
wg.Wait()
close(out)
}()
return out
}
// usage: fan-out to 3 workers reading the same input, then fan-in
// w1, w2, w3 := square(ctx, src), square(ctx, src), square(ctx, src)
// for v := range Merge(ctx, w1, w2, w3) { ... }
The design has one forwarding goroutine per input and a single closer goroutine that waits for all of them.
Gotchas:
- You cannot
close(out)from a forwarder, because the first one to finish would close the channel while the others are still sending. - Without the
ctxcases, a consumer that stops reading leaks every forwarder. Guarding only the send is not enough: a forwarder idle on an input that is never closed would still leak, which is why the receive is in aselecttoo. - Merge does not preserve ordering across inputs.
- For a fixed small N you can write a single-goroutine version with
selectand nil-channel disabling (see the next question).
More on Concurrency Patterns & sync
- Q236What does this print in Go 1.22+ vs before?
- Q237Implement a generic worker pool with a fixed number of workers, context cancellation, and no goroutine leaks.
- Q239How do nil channels behave in select, and how are they used to merge two channels?
- Q240Build a cancellable pipeline. What causes goroutine leaks in pipelines and how do you prevent them?
- Q241Spot the goroutine leak.
- Q242How do you close a channel safely when there are multiple senders?