The idea in one minute#
Almost all concurrent Go is built from a handful of shapes: a worker pool (bounded parallelism), a pipeline (stages connected by channels), fan-out / fan-in (spread work, collect results), a semaphore (limit how many at once), and an error group (first failure cancels the rest).
Two rules sit under all of them. Bound everything: every queue has a size and every pool a limit, so overload turns into waiting or rejection instead of memory growth. Every goroutine must have a way to end: know who stops it and how you would notice if it did not.
An analogy#
A post office. A fixed number of counters (worker pool) so staff are not overwhelmed. A queue with a rope — when the rope is full, the door says “come back later” (backpressure). Sorting, weighing, stamping in sequence (pipeline). And at closing time, every clerk finishes the customer in front of them and goes home; nobody is left locked in the back room (no leaks).
A picture#
flowchart LR SRC["Producer<br/>closes jobs when done"] -->|"jobs chan, bounded"| W1["worker 1"] SRC --> W2["worker 2"] SRC --> W3["worker N"] W1 -->|"results chan"| COL["Collector"] W2 --> COL W3 --> COL CTX["ctx: first error or<br/>caller cancel"] -.->|"Done()"| SRC CTX -.-> W1 CTX -.-> W2 CTX -.-> W3 WG["WaitGroup: when all workers exit,<br/>close(results)"] -.-> COL FULL["queue full"] -->|"block, drop or reject:<br/>a decision, not an accident"| SRC class SRC,COL neutral class W1,W2,W3 compute class CTX,WG queue class FULL warn
How it really works#
Worker pool#
N goroutines read from one jobs channel. Closing the channel ends them; a WaitGroup tells you
when they are done.
jobs := make(chan Job)
var wg sync.WaitGroup
for range n {
wg.Add(1)
go func() {
defer wg.Done()
for j := range jobs {
process(ctx, j)
}
}()
}How many workers? For CPU-bound work, runtime.GOMAXPROCS(0): more only adds switching.
For I/O-bound work, as many as the downstream system can take — a number to measure, not to
guess — and that limit is really a property of the downstream, so express it as a semaphore.
Semaphore: bounded concurrency without a pool#
sem := make(chan struct{}, 16)
for _, item := range items {
sem <- struct{}{} // blocks when 16 are in flight
go func() {
defer func() { <-sem }()
call(ctx, item)
}()
}Simpler than a pool when work arrives as a loop. errgroup.Group with SetLimit(16) is the
same thing with error handling.
Error group#
Run several tasks; the first error cancels the others and is returned.
g, ctx := errgroup.WithContext(ctx) // golang.org/x/sync/errgroup
g.SetLimit(8)
for _, u := range urls {
g.Go(func() error { return fetch(ctx, u) })
}
if err := g.Wait(); err != nil { /* ... */ }The program below implements this from the standard library so you can see there is no magic.
Pipeline#
Each stage is a goroutine (or a pool) that reads from an input channel, writes to an output channel, and closes its output when its input is exhausted. Closure propagates down the pipeline; cancellation propagates through the shared context.
read → tokenize → embed (pool of 4) → writeEvery send in a stage must be a select with ctx.Done(), or a downstream stage that stops
early leaves every upstream stage blocked forever.
Fan-in#
Merge several channels into one: one forwarding goroutine per input, a WaitGroup to close the
output after all have finished. When order matters, carry an index with each item and reorder
at the end, or give each job a slot in a preallocated results slice (IV.05).
Backpressure#
A bounded channel is backpressure for free: when consumers are slow, producers block. At the edge of the system — an HTTP handler — blocking is not acceptable, so make the choice explicit:
| Policy | Code | Use when |
|---|---|---|
| Block | queue <- job | Internal stages; the caller can wait |
| Reject | select { case queue <- job: default: return ErrBusy } → HTTP 429 or 503 | Request handlers: fail fast and let the client retry |
| Wait with a deadline | select with ctx.Done() | Requests with a latency budget |
| Drop oldest / newest | Explicit ring buffer | Telemetry, live streams |
Unbounded queues and unbounded go statements are the same bug: they convert overload into an
out-of-memory kill several minutes later.
Batching#
For AI workloads, grouping requests is the single biggest throughput lever. The shape:
collect until (batch is full) OR (the oldest item has waited maxWait)One goroutine owns the batch; requests arrive on a channel carrying a reply channel each. It is
a select over the input channel and a timer. Lesson VI.08 builds it.
Rate limiting#
time.Ticker for a fixed rate; golang.org/x/time/rate for a token bucket with bursts
(limiter.Wait(ctx)). For LLM APIs, limit tokens per minute, not only requests:
limiter.WaitN(ctx, estimatedTokens).
Goroutine leaks#
A leaked goroutine is one blocked forever. It holds its stack and everything reachable from it, and it never shows up as an error.
| Cause | Prevention |
|---|---|
| Sending a result nobody will receive | Buffer of 1, or select with ctx.Done() |
| Receiving from a channel nobody will close | The sender closes; the receiver also selects on ctx.Done() |
| A worker pool whose jobs channel is never closed | defer close(jobs) in the producer |
Waiting on a lock or Cond held by a goroutine that died | Short critical sections; defer Unlock |
| A ticker loop with no exit | select on ctx.Done(); defer ticker.Stop() |
Detecting them:
runtime.NumGoroutine()as an exported metric: a line that only climbs is a leak.- The goroutine profile (
/debug/pprof/goroutine?debug=1) groups goroutines by stack: ten thousand parked at the same line is the answer. - The goroutine leak profile (
/debug/pprof/goroutineleak; experimental in Go 1.26, generally available in 1.27) uses the garbage collector to find goroutines blocked on channels or locks that nothing else can reach — leaks that are provably permanent. - In tests:
testing/synctestmakes a test fail if goroutines started inside its bubble are still blocked when it ends;go.uber.org/goleakis the long-standing third-party check.
Graceful shutdown#
1. Stop accepting new work (close the listener; srv.Shutdown(ctx))
2. Let in-flight work finish (up to a deadline)
3. Cancel what remains (cancel the root context)
4. Wait for goroutines to exit (WaitGroup)
5. Flush and close resourcessignal.NotifyContext turns SIGTERM into a context cancellation, which is all that step 3 needs.
Code#
// pool.go — a bounded, cancellable worker pool with first-error cancellation and no leaks.
package main
import (
"context"
"errors"
"fmt"
"runtime"
"sync"
"sync/atomic"
"time"
)
// Group is a minimal errgroup: bounded concurrency, first error cancels, Wait returns it.
type Group struct {
ctx context.Context
cancel context.CancelCauseFunc
sem chan struct{}
wg sync.WaitGroup
}
func NewGroup(ctx context.Context, limit int) (*Group, context.Context) {
ctx, cancel := context.WithCancelCause(ctx)
return &Group{ctx: ctx, cancel: cancel, sem: make(chan struct{}, limit)}, ctx
}
// Go blocks while `limit` tasks are running: that is the backpressure.
func (g *Group) Go(f func(ctx context.Context) error) {
select {
case g.sem <- struct{}{}:
case <-g.ctx.Done():
return // already cancelled: do not start more work
}
g.wg.Add(1)
go func() {
defer g.wg.Done()
defer func() { <-g.sem }()
if err := f(g.ctx); err != nil {
g.cancel(err) // only the first cause is kept
}
}()
}
func (g *Group) Wait() error {
g.wg.Wait()
err := context.Cause(g.ctx)
g.cancel(nil)
return err
}
func embed(ctx context.Context, id int, inFlight, peak *atomic.Int64) error {
n := inFlight.Add(1)
defer inFlight.Add(-1)
for {
p := peak.Load()
if n <= p || peak.CompareAndSwap(p, n) {
break
}
}
select {
case <-time.After(5 * time.Millisecond): // the "work"
case <-ctx.Done():
return nil // stopped early: not this task's error
}
if id == 37 {
return fmt.Errorf("item %d: model returned NaN", id)
}
return nil
}
func main() {
before := runtime.NumGoroutine()
for _, failing := range []bool{false, true} {
var inFlight, peak, started atomic.Int64
g, _ := NewGroup(context.Background(), 8)
start := time.Now()
for id := 0; id < 100; id++ {
if !failing && id == 37 {
continue
}
g.Go(func(ctx context.Context) error {
started.Add(1)
return embed(ctx, id, &inFlight, &peak)
})
}
err := g.Wait()
fmt.Printf("failing=%-5v started %3d tasks, peak concurrency %d, %3d ms, err: %v\n",
failing, started.Load(), peak.Load(), time.Since(start).Milliseconds(), err)
}
// Backpressure at the edge: reject instead of queueing without bound.
queue := make(chan int, 4)
var ErrBusy = errors.New("busy")
submit := func(j int) error {
select {
case queue <- j:
return nil
default:
return ErrBusy
}
}
accepted, rejected := 0, 0
for j := 0; j < 10; j++ {
if submit(j) == nil {
accepted++
} else {
rejected++
}
}
fmt.Printf("burst of 10 into a queue of 4 with no consumer: %d accepted, %d rejected\n", accepted, rejected)
time.Sleep(20 * time.Millisecond)
fmt.Printf("goroutines before %d, after %d: nothing leaked\n", before, runtime.NumGoroutine())
}Remember this#
- Worker pool, semaphore, error group, pipeline, fan-in: five shapes cover most programs.
- Bound every queue and every pool; decide what happens when they are full.
- CPU-bound: about
GOMAXPROCSworkers. I/O-bound: the downstream’s limit. - Every send and receive that can block also selects on
ctx.Done(). - A goroutine leak is a goroutine blocked forever; watch
NumGoroutineand use the leak profile.
Try it#
- Run
pool.go. In the failing run, how many tasks started before cancellation took effect, and why not all 100? - Add
g.Gotasks that ignorectxand sleep 1 s. What doesWaitdo after the first error? What does that tell you about cooperative cancellation? - Build a three-stage pipeline (generate → square → sum) where the consumer stops after ten
values. Prove with
runtime.NumGoroutinethat every stage exits.
Check yourself#
- How do you choose the number of workers for CPU-bound and for I/O-bound work?
- What are the options when a bounded queue is full?
- What is a goroutine leak, and how would you detect one in production?