Go

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 ctx cases, 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 a select too.
  • Merge does not preserve ordering across inputs.
  • For a fixed small N you can write a single-goroutine version with select and nil-channel disabling (see the next question).

More on Concurrency Patterns & sync

All 38 Concurrency Patterns & sync questions