Implement a bounded, cancellable fan-out that returns the first error and does not leak goroutines.
Question 175HardGo 1.22 to 1.25
Requirements: limit concurrency, stop the remaining work at the first error, wait for every goroutine to finish, and never block a sender forever. Without errgroup:
func FanOut[T any](ctx context.Context, items []T, limit int,
fn func(context.Context, T) error) error {
ctx, cancel := context.WithCancelCause(ctx)
defer cancel(nil)
sem := make(chan struct{}, limit)
var wg sync.WaitGroup
loop:
for _, it := range items {
select {
case sem <- struct{}{}:
case <-ctx.Done():
break loop // stop launching new work
}
wg.Go(func() {
defer func() { <-sem }()
if err := fn(ctx, it); err != nil {
cancel(err) // first cause wins; later calls are no-ops
}
})
}
wg.Wait()
if err := context.Cause(ctx); err != nil && ctx.Err() != nil {
return err
}
return nil
}
Key points:
- The semaphore is acquired before
go. This bounds the number of live goroutines, not just the number doing work. WithCancelCauserecords the first error without a mutex.fnmust respectctx, otherwise cancellation cannot stop it.wg.Wait()guarantees there are no stray goroutines when the function returns.
In real code, errgroup.WithContext plus SetLimit does exactly this.
More on Goroutines & the Scheduler
- Q173Why can many CPU-bound goroutines hurt performance, and how do you size a worker pool?
- Q174How do goroutines interact with cgo, and why can cgo calls be expensive?
- Q176Is creating a goroutine per request (e.g., in net/http) a good design? What are the risks?
- Q177How are timers handled by the scheduler, and what changed with time.Timer in Go 1.23?
- Q178What are the scheduling implications of runtime.Goexit, blocking in init, and unbuffered sends from many goroutines?
- Q179When are the arguments of a go statement evaluated? What does this print?