Lesson 17 / 25

Fan-out and Fan-in

Parallelising a slow stage and merging the results.

Many readers, one merged output

Fan-out means starting several goroutines that read from the same input channel; each value is received by exactly one of them, so work is spread automatically. Fan-in merges several result channels into one: start a forwarding goroutine per input, track them with a WaitGroup, and close the merged output after wg.Wait(). Fan-out loses ordering: results arrive in completion order. If order matters, attach an index to each item and reassemble at the end, or write results into a pre-sized slice by index. Fan-out helps when a stage is slow because of I/O or CPU work that can genuinely run in parallel.

Three workers fanned out, results fanned in

merge closes its output after every forwarder finishes.

package main

import (
	"context"
	"fmt"
	"sync"
)

func merge(ctx context.Context, ins ...<-chan string) <-chan string {
	out := make(chan string)
	var wg sync.WaitGroup
	for _, in := range ins {
		wg.Add(1)
		go func() {
			defer wg.Done()
			for v := range in {
				select {
				case out <- v:
				case <-ctx.Done():
					return
				}
			}
		}()
	}
	go func() { wg.Wait(); close(out) }()
	return out
}

func worker(id int, jobs <-chan int) <-chan string {
	out := make(chan string)
	go func() {
		defer close(out)
		for j := range jobs { // all workers share the jobs channel
			out <- fmt.Sprintf("worker %d did job %d", id, j)
		}
	}()
	return out
}

func main() {
	ctx := context.Background()
	jobs := make(chan int)
	go func() {
		defer close(jobs)
		for j := range 6 {
			jobs <- j
		}
	}()

	w1, w2, w3 := worker(1, jobs), worker(2, jobs), worker(3, jobs)
	for line := range merge(ctx, w1, w2, w3) {
		fmt.Println(line) // completion order, not job order
	}
}

More goroutines is not always faster

For CPU-bound work, more workers than GOMAXPROCS rarely helps. For I/O-bound work the limit is usually the remote system, so measure and bound the fan-out.

Quick check: What is lost when a stage is fanned out to several workers?

  • Channels can no longer be closed
  • Values may be delivered to two workers at once
  • The original ordering of the results
  • Context cancellation stops working
Answer

The original ordering of the results — Each value goes to exactly one receiver, but completion order varies.