Below the API

Concurrency Patterns

Intermediate Advanced 1h 5m Difficulty 4/5

Prerequisites 01–05

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) → write

Every 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:

PolicyCodeUse when
Blockqueue <- jobInternal stages; the caller can wait
Rejectselect { case queue <- job: default: return ErrBusy } → HTTP 429 or 503Request handlers: fail fast and let the client retry
Wait with a deadlineselect with ctx.Done()Requests with a latency budget
Drop oldest / newestExplicit ring bufferTelemetry, 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.

CausePrevention
Sending a result nobody will receiveBuffer of 1, or select with ctx.Done()
Receiving from a channel nobody will closeThe sender closes; the receiver also selects on ctx.Done()
A worker pool whose jobs channel is never closeddefer close(jobs) in the producer
Waiting on a lock or Cond held by a goroutine that diedShort critical sections; defer Unlock
A ticker loop with no exitselect 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/synctest makes a test fail if goroutines started inside its bubble are still blocked when it ends; go.uber.org/goleak is 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 resources

signal.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 GOMAXPROCS workers. 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 NumGoroutine and use the leak profile.

Try it#

  1. Run pool.go. In the failing run, how many tasks started before cancellation took effect, and why not all 100?
  2. Add g.Go tasks that ignore ctx and sleep 1 s. What does Wait do after the first error? What does that tell you about cooperative cancellation?
  3. Build a three-stage pipeline (generate → square → sum) where the consumer stops after ten values. Prove with runtime.NumGoroutine that every stage exits.

Check yourself#

  1. How do you choose the number of workers for CPU-bound and for I/O-bound work?
  2. What are the options when a bounded queue is full?
  3. What is a goroutine leak, and how would you detect one in production?

↑↓ navigate ↵ open