Go for Infrastructure Engineers›06 · Concurrency you'll actually use

Lesson 06 of 9 · Part 2 — Tools that talk to systems

Concurrency you'll actually use

The concurrency patterns infrastructure code actually needs: goroutines and WaitGroups, worker pools with bounded parallelism, collecting results and errors, cancellation with context, errgroup, protecting shared state, and finding races and leaks.

Practitioner
Key wordsgoroutineschannelssync.WaitGroupworker poolcontextcancellationerrgroupSetLimitsync.Mutexrace detectorgoroutine leakselect

Goroutines and WaitGroups

A goroutine is a function running concurrently; start one with go. A sync.WaitGroup waits for a group of them:

package main

import (
    "fmt"
    "sync"
    "time"
)

func check(node string) string {
    time.Sleep(100 * time.Millisecond) // pretend to call the node
    return node + ": ok"
}

func main() {
    nodes := []string{"n1", "n2", "n3", "n4"}
    results := make([]string, len(nodes))

    var wg sync.WaitGroup
    for i, n := range nodes {
        wg.Add(1)
        go func() {
            defer wg.Done()
            results[i] = check(n) // each goroutine writes its own index: no lock needed
        }()
    }
    wg.Wait()
    fmt.Println(results)
}

(Since Go 1.22 each loop iteration has its own i and n, so capturing them in the closure is safe.)

Goroutines are extra pairs of hands. A WaitGroup is the clipboard where each helper ticks "done". Channels are conveyor belts between helpers. A context is the whistle: when it blows, everyone stops what they're doing and goes home.

Bounded parallelism with errgroup

Unbounded goroutines can overload APIs, file descriptors or the cluster you're talking to. errgroup (golang.org/x/sync/errgroup) bounds parallelism and handles the first error:

package main

import (
    "context"
    "fmt"
    "strings"
    "time"

    "golang.org/x/sync/errgroup"
)

func audit(ctx context.Context, cluster string) (string, error) {
    select {
    case <-time.After(200 * time.Millisecond): // pretend to call the API
    case <-ctx.Done():
        return "", ctx.Err()
    }
    if strings.HasSuffix(cluster, "-broken") {
        return "", fmt.Errorf("audit %s: api unreachable", cluster)
    }
    return cluster + ": 0 findings", nil
}

func main() {
    clusters := []string{"dev-1", "dev-2", "stg-1", "prd-1", "prd-2-broken", "prd-3"}
    results := make([]string, len(clusters))

    g, ctx := errgroup.WithContext(context.Background())
    g.SetLimit(3) // at most 3 audits at once
    for i, c := range clusters {
        g.Go(func() error {
            r, err := audit(ctx, c)
            if err != nil {
                return err // cancels ctx for the others
            }
            results[i] = r
            return nil
        })
    }
    if err := g.Wait(); err != nil {
        fmt.Println("error:", err)
    }
    for _, r := range results {
        if r != "" {
            fmt.Println(r)
        }
    }
}

If you want all results even when some fail (a report across clusters), don't return the error from g.Go; store it next to the result instead.

Channels: when work flows between stages

Channels connect producers and consumers. A worker pool reading jobs from one channel and writing results to another:

package main

import (
    "fmt"
    "sync"
)

func main() {
    jobs := make(chan string)
    results := make(chan string)

    var wg sync.WaitGroup
    for w := 1; w <= 3; w++ {
        wg.Add(1)
        go func() {
            defer wg.Done()
            for node := range jobs { // ends when jobs is closed
                results <- fmt.Sprintf("worker %d drained %s", w, node)
            }
        }()
    }

    go func() { // close results once every worker is done
        wg.Wait()
        close(results)
    }()

    go func() {
        for _, n := range []string{"n1", "n2", "n3", "n4", "n5"} {
            jobs <- n
        }
        close(jobs) // the sender closes
    }()

    for r := range results {
        fmt.Println(r)
    }
}

Rules: the sender closes a channel; receivers range until it's closed; a send to a channel nobody reads blocks forever (a leak). Prefer errgroup or plain WaitGroups when you don't need a pipeline.

Context: deadlines and cancellation everywhere

Every function that does I/O or waits takes ctx context.Context as its first parameter and gives up when ctx.Done() fires. Create deadlines at the top (context.WithTimeout) and pass the context down; never store it in structs.

Shared state: mutex or ownership

When goroutines must update shared data, protect it with a sync.Mutex (or use sync/atomic for counters). Better still, give each goroutine its own slot (as in the first example) or send results over a channel to one owner. Then run tests with -race: the race detector finds unsynchronised access at run time.

Try it: audit many clusters at once

  1. Run the errgroup program; change SetLimit to 1 and 6 and time it.
  2. Change it to collect all results and errors instead of stopping at the first error.
  3. Add a 1-second overall context.WithTimeout and make one audit slow; confirm everything stops on time.
  4. Remove the per-index writes and append to a shared slice from all goroutines; run with -race and read the report; fix it with a mutex.
  5. Create a leak on purpose (send to an unbuffered channel with no reader), find it with runtime.NumGoroutine() or pprof, then fix it.

Going deeper: concurrency in long-running programs

  • Every goroutine needs a clear owner and a way to stop (context or a closed channel).
  • Expose net/http/pprof on a local port in services and controllers to inspect goroutines and memory in production.
  • For rate-limited APIs, combine bounded parallelism with a token-bucket limiter.

Recap

  • Goroutines + WaitGroup for simple fan-out; write results to separate slots.
  • errgroup with SetLimit for bounded parallelism and first-error cancellation.
  • Channels for pipelines and worker pools: the sender closes.
  • context in every I/O function; mutex or ownership for shared state; test with -race.

This site is a public version of my personal engineering knowledge hub. It intentionally excludes confidential company information and internal operational details.