~/writing
Applied Go concurrency patterns

I gave a talk at Ankara Gophers #11 in December 2022. Rather than walk through the patterns one at a time as isolated snippets, I built a batch image processor to demonstrate the concurrency patterns in Go, and added one pattern per step, letting each new requirement pull the next pattern in. The code is on GitHub at cengizcan/go-concurrency-gophers-11 and you can find the slides here: Applied Go Concurrency Patterns.pdf.
Concurrency is hard
Two terms worth keeping separate, because they get used interchangeably despite meaning different things: race condition and data race.
A race condition occurs when the timing or order of events affects the correctness of the code. A data race occurs when one goroutine accesses a mutable object while another goroutine is writing to it.
You can have a race condition with no data race at all. The common root of both is sharing state which makes concurrency hard. Go’s answer is two primitives: Goroutines and Channels.
Goroutines are independently executing functions, kinda like green threads, and much cheaper than an OS thread: they start with tiny growable stacks (~2KB) instead of fixed multi-MB ones, and Go schedules them in user space instead of paying for kernel context switches. Channels are the communication structure between them, and they exist specifically so you don’t have to share state.
The idea of channels, adopted by Go, comes from Tony Hoare’s Communicating Sequential Processes, 1978: sequential processes that communicate data rather than share it.
Famous proverb comes from the same place:
Do not communicate by sharing memory; instead, share memory by communicating.
c := make(chan int, 3)// Send. Sender is blocked if the buffer is full or there is no receiverc <- 1// Receive. Receiver is blocked if the buffer is empty or there is no senderval := <-c // Pass the copy of the valueChannels are not bulletproof, though. They deadlock when goroutines end up waiting on each other inifinitely. They copy data, which can cost you. Passing pointers through them puts you right back into data-race territory. And sending on a closed channel panics.
The project: A batch image processor
Build a batch image processor that scans a file system for images, then for each one resizes it, assigns it an incremental document ID, labels it with multiple AI algorithms, compresses it via an external lossless compression service, and saves it, finally scanning multiple sources concurrently.
Every step below adds one of these, and each one needs a pattern the previous step didn’t.
Step 1: Scanning, and the generator pattern
Walk the file system recursively; when you hit an image, queue it for processing. We can follow two options here: collect every path into a slice and return it at the end, or use the generator pattern, a function that starts a goroutine and returns a receive-only channel down which it sends each path the moment it’s found.
func generator() <-chan int { out := make(chan int) go func() { for i := 0; i < 5; i++ { out <- i } }() return out}Let’s apply the generator pattern to our image scanner, so the processor,the consumer of the channel here, starts working on the first file before the walk has finished:
func scan(dir string) <-chan string { out := make(chan string)
go func() { files, err := ioutil.ReadDir(dir) if err != nil { panic(err) }
for _, f := range files { path := fmt.Sprintf("%s/%s", dir, f.Name()) if f.IsDir() { // walk through the child directory for p := range scan(path) { out <- p } } else if strings.HasSuffix(f.Name(), ".jpg") { out <- path } } close(out) }()
return out}Beware: skip close(out) and every range on that channel blocks forever.
Step 2: Sequential versus a goroutine per image
Now the actual work: read, resize, generate a document id, save.
func ProcessImageS2(name string, num numerator.Sequential) *ProcessResponse { img := image.ReadImage(name) r := len(img.Bytes) img = image.Resize(img) img.Id = num.Next() image.WriteImage(img)
return &ProcessResponse{ Name: img.Name, Meta: img.Meta, Id: img.Id, Read: r, Written: len(img.Bytes), }}The sequential version is a loop. The concurrent version starts a goroutine per image
and waits on a WaitGroup:
func ProcessConcurrent(cnt int) { wg := new(sync.WaitGroup) for p := range image.ScanImages(cnt) { wg.Add(1) go func(path string) { defer wg.Done() ProcessImageS2(path, numerator.NewV2()) }(p) } wg.Wait()}The numbers:
| Images | Sequential | Goroutines | Faster by |
|---|---|---|---|
| 10 | 7.37 ms | 0.73 ms | 10× |
| 100 | 74.6 ms | 1.4 ms | 50× |
| 1,000 | 724 ms | 9 ms | 80× |
| 10,000 | 7.27 s | 0.09 s | 80× |
Note where it stops improving. The speedup climbs to 80× and then flattens, because the serial part of the program doesn’t go away no matter how many goroutines you add.
Amdahl’s law, briefly: the serial part sets the ceiling. If 95% of the program can be parallel the most you can ever get is 20 times faster execution, however many cores you add.
Step 3: Limit the number of Goroutines
A goroutine per image works well at ten thousand images and not at a million. The fan out pattern puts a fixed number of worker goroutines on one input channel.
func FanOut(input <-chan string, workerCnt int, fn ProcessImage) <-chan *ProcessResponse { num := numerator.NewV2()
output := make(chan *ProcessResponse) var wg sync.WaitGroup
for i := 0; i < workerCnt; i++ { wg.Add(1) go func(ind int) { defer wg.Done() for { path, ok := <-input if !ok { return // return when closed } output <- fn(path, num) } }(i) } go func() { wg.Wait() close(output) }() return output}The wg.Wait() sits in its own goroutine so it can close output once every worker has
returned and the receiver’s range then terminates on its own.
| Images | Unbounded | 100 workers | 500 workers | 1,000 workers |
|---|---|---|---|---|
| 10,000 | 0.2 s | 1 s | 0.25 s | 0.2 s |
| 500,000 | 11 s | 47.5 s | 10.1 s | 9.0 s |
| 750,000 | 28 s | 71.4 s | 16.1 s | 15.5 s |
At ten thousand images a worker pool buys us nothing (it ties the unbounded version at best) and 100 workers is five times worse; we’ve throttled. At 750,000 the unbounded version takes 28 seconds and 1,000 workers takes 15.5. The right worker count is a function of the workload, and picking it badly is worse than not limiting at all.
Step 4: Compression, and taking the first response
Next, we will send the image to several external compression services. They return roughly the same
result at different speeds, so take whichever answers first and discard the rest, the
same shape as JavaScript’s Promise.any().
func Compress(src []byte) *CompResult { result := make(chan *CompResult)
for _, comp := range services { go func(compress Compressor) { res := compress(src)
select { case result <- res: default: } }(comp) } return <-result}The select with a default is the important part. The losing goroutines find nobody
receiving, fall to default, and exit instead of blocking forever on a send. Without
it, three of the four goroutines leak on every single image.
There’s an alternative way to write it, called the quit channel pattern:
func compressWithDone(quit chan bool, compress Compressor, src []byte, result chan<- *CompResult) { response := make(chan *CompResult) go func() { response <- compress(src) }()
select { case result <- <-response: case <-quit: return }}Step 5: Labelling, and the timeout pattern
We will use several downstream services to satisfy requirements like labeling the objects in an image and scanning them for compliance. We want all the responses from those services, but we won’t wait forever, so we take everything that arrives before the timeout.
func Label(src []byte) map[string]string { timeout := make(chan struct{}, 1) results := make(chan *LabelResult, len(algos))
go func() { time.Sleep(15 * time.Millisecond) timeout <- struct{}{} }()
for _, t := range algos { go t(src, timeout, results) }
resMap := make(map[string]string) for { select { case r := <-results: resMap[r.AlgoName] = r.Value if len(resMap) == len(algos) { return resMap } case <-timeout: return resMap } }}results is buffered to len(algos) intentionally because an unbuffered channel would leave every late goroutine blocked on
its send forever.
Step 6: A worker pool
Now, fan out works, but it’s not flexible enough, welded to one job. A worker pool, however, is the same idea made generic, so it can finally be written once for any result type.
type WorkerPool[T any] struct { workerCnt int Tasks chan Executable[T] Results chan T}
type Executable[T any] interface { Execute() T}
func worker[T any](ctx context.Context, tasks <-chan Executable[T], results chan<- T, wg *sync.WaitGroup) { defer wg.Done() for { select { case task, ok := <-tasks: if !ok { return } results <- task.Execute() case <-ctx.Done(): return } }}The case <-ctx.Done(): is what fan out was missing. Cancellation now propagates: every
worker is watching the same context, and the whole pool unwinds when it’s cancelled.
Step 7: Multiple sources, and the fan in pattern
Finally, we’ve arrived at the last requirement: feed the pool from several scanners at once. Fan in merges many channels into one.
func fanIn(channels ...<-chan string) <-chan string { var wg sync.WaitGroup multiplexedStream := make(chan string)
wg.Add(len(channels)) for _, ch := range channels { go func(c <-chan string) { defer wg.Done() for i := range c { multiplexedStream <- i } }(ch) }
go func() { wg.Wait() close(multiplexedStream) }()
return multiplexedStream}Same closing trick as fan out: one goroutine waits on the group and closes the merged
channel, so the consumer’s range ends cleanly.
Bonus: Mutex versus atomic
We’ve done a lot, but we haven’t been able to remove all the shared state: the document numerator is the one genuinely shared piece of state in the program. Let’s see how we can gracefully handle this. There are three implementation alternatives. The first, the naive one, isn’t graceful at all:
type sequenceV1 struct { current int}
func (s *sequenceV1) Next() int { s.current++ return s.current}s.current++ is a read, an add, and a write. Two goroutines interleaving there will
hand out the same id twice, but why? We’ll come to this later. The second alternative uses a mutex:
type sequenceV2 struct { mu sync.Mutex current int}
func (s *sequenceV2) Next() int { s.mu.Lock() defer s.mu.Unlock() s.current++ return s.current}A mutex locks a critical section of the code so there’s only one execution at a time. Now we’re safe from duplicate IDs, but as our code grows we’d risk deadlocks caused by multiple mutexes.
And finally, with an atomic:
type sequenceV3 struct { current int32}
func (s *sequenceV3) Next() int { return int(atomic.AddInt32(&s.current, 1))}Under the hood, a plain current++ in our naive solution is a three-step process for the CPU: load, add, store. It fetches the value into a register, increments it, then writes it back, three separate steps another core can interleave with. Atomics collapse those into a single, hardware-locked CPU instruction, so they’re thread-safe by construction and faster than reaching for a mutex.
In order of preference: design so you don’t need a lock at all; if you do, measure with the mutex contention profiler; then try an atomic.
Wrapping up
All the patterns we covered:
- Generator: return a channel instead of a slice
- Fan out: fixed set of workers on one input
- First-response-wins:
selectwith adefault(or a quit channel), take whichever answers first - Timeout: take what arrived, move on
- Worker pool: fan out, generalised, with cancellation
- Fan in: merge many channels into one
References
- Go Concurrency Patterns — Rob Pike, Google I/O 2012
- Advanced Go Concurrency Patterns — Sameer Ajmani, Google I/O 2013
- Concurrency is not parallelism
- Pipelines and cancellation
- Timing out, moving on
- Communicating Sequential Processes — Hoare, 1978
- Amdahl’s law