Implement a generic worker pool with a fixed number of workers, context cancellation, and no goroutine leaks.
Question 237HardGo 1.22 to 1.25
func WorkerPool[In, Out any](
ctx context.Context,
workers int,
in <-chan In,
fn func(context.Context, In) Out,
) <-chan Out {
out := make(chan Out)
var wg sync.WaitGroup
for range workers {
wg.Go(func() {
for {
select {
case <-ctx.Done():
return
case v, ok := <-in:
if !ok {
return
}
select {
case out <- fn(ctx, v):
case <-ctx.Done():
return
}
}
}
})
}
go func() {
wg.Wait()
close(out) // only after ALL workers have exited
}()
return out
}
Key points:
- There are exactly
workersgoroutines, which bounds concurrency and memory. - Only the owner closes
out. A separate goroutine closes it afterwg.Wait(), so no worker ever sends on a closed channel. - Every blocking send and receive sits in a
selectwithctx.Done(), so an abandoned consumer never leaves workers blocked forever. - The producer that feeds
inmust also respectctxand closeinwhen it has no more work.
Results arrive out of order. To keep order, send (index, value) pairs or write into a pre-sized slice.
What the interviewer is looking for: who closes which channel, leak-free cancellation, and bounded goroutines. Also mention that for "N tasks, stop on first error" errgroup with SetLimit is simpler.
More on Concurrency Patterns & sync
- Q235Does this counter give the right answer with GOMAXPROCS=1? What about counter++ in general?
- Q236What does this print in Go 1.22+ vs before?
- Q238Implement fan-out / fan-in: a generic Merge that combines N channels into one.
- 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.