Advanced Go Concurrency Patterns

Table of Contents
Go's concurrency model, built around goroutines and channels, provides a powerful and elegant way to build highly concurrent systems. While the basics of starting a goroutine (go func()) and communicating via channels are well understood by most Go developers, mastering advanced concurrency patterns is crucial for building robust, scalable, and high-performance applications. In this deep dive, we will explore several advanced Go concurrency patterns, dissecting their implementation and discussing when and why to use them.
1. The Pipeline Pattern
The pipeline pattern is a cornerstone of concurrent data processing in Go. It consists of a series of stages connected by channels, where each stage is a group of goroutines running the same function. In each stage, the goroutines:
- Receive values from upstream via inbound channels.
- Perform some function on that data, usually producing new values.
- Send values downstream via outbound channels.
Implementing a Pipeline
A robust pipeline requires careful management of goroutines, especially concerning cancellation and error handling. Let's construct a multi-stage pipeline that generates numbers, squares them, and then filters for even results.
package main
import (
"context"
"fmt"
"sync"
)
// Stage 1: Generator
func generator(ctx context.Context, nums ...int) <-chan int {
out := make(chan int)
go func() {
defer close(out)
for _, n := range nums {
select {
case out <- n:
case <-ctx.Done():
return
}
}
}()
return out
}
// Stage 2: Squarer
func squarer(ctx context.Context, in <-chan int) <-chan int {
out := make(chan int)
go func() {
defer close(out)
for n := range in {
select {
case out <- n * n:
case <-ctx.Done():
return
}
}
}()
return out
}
// Stage 3: Filter
func filterEven(ctx context.Context, in <-chan int) <-chan int {
out := make(chan int)
go func() {
defer close(out)
for n := range in {
if n%2 == 0 {
select {
case out <- n:
case <-ctx.Done():
return
}
}
}
}()
return out
}
func main() {
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
nums := []int{1, 2, 3, 4, 5, 6, 7, 8, 9, 10}
genChan := generator(ctx, nums...)
sqChan := squarer(ctx, genChan)
filterChan := filterEven(ctx, sqChan)
for val := range filterChan {
fmt.Println(val)
}
}
This pipeline ensures that if the context is canceled, all goroutines gracefully exit, preventing resource leaks. The explicit use of context.Context is a best practice for advanced Go concurrency.
2. Fan-Out and Fan-In
When a single stage in a pipeline is computationally expensive, it can become a bottleneck. The fan-out/fan-in pattern solves this by distributing the work across multiple goroutines (fan-out) and then multiplexing their results back onto a single channel (fan-in).
Fan-Out
Fan-out simply involves starting multiple goroutines to read from the same channel. The channel acts as a work queue, distributing items among the workers.
Fan-In
Fan-in requires a multiplexer function that takes multiple channels and returns a single channel. It uses a sync.WaitGroup to track when all input channels have been closed.
func merge(ctx context.Context, cs ...<-chan int) <-chan int {
var wg sync.WaitGroup
out := make(chan int)
output := func(c <-chan int) {
defer wg.Done()
for n := range c {
select {
case out <- n:
case <-ctx.Done():
return
}
}
}
wg.Add(len(cs))
for _, c := range cs {
go output(c)
}
go func() {
wg.Wait()
close(out)
}()
return out
}
By combining pipeline stages with fan-out/fan-in, we can dramatically increase throughput for CPU-bound tasks.
3. The Worker Pool Pattern
The worker pool pattern is essential for limiting concurrency, which is vital when interacting with resources that have strict rate limits or when trying to avoid unbounded memory consumption.
Unlike unbounded goroutine creation, a worker pool maintains a fixed number of goroutines that consume tasks from a shared channel.
type Task struct {
ID int
Value string
}
type Result struct {
TaskID int
Output string
}
func worker(id int, tasks <-chan Task, results chan<- Result) {
for t := range tasks {
// Simulate work
output := fmt.Sprintf("Worker %d processed Task %d: %s", id, t.ID, t.Value)
results <- Result{TaskID: t.ID, Output: output}
}
}
func main() {
numWorkers := 3
numTasks := 10
tasks := make(chan Task, numTasks)
results := make(chan Result, numTasks)
// Start workers
for w := 1; w <= numWorkers; w++ {
go worker(w, tasks, results)
}
// Send tasks
for j := 1; j <= numTasks; j++ {
tasks <- Task{ID: j, Value: fmt.Sprintf("Data %d", j)}
}
close(tasks)
// Collect results
for a := 1; a <= numTasks; a++ {
<-results
}
}
This pattern controls the maximum number of concurrent operations, providing stability and predictable performance under load.
4. Context for Cancellation and Timeouts
In any sophisticated concurrent system, the ability to cancel in-flight operations is critical. The context package is the standard way to propagate cancellation signals and deadlines across API boundaries and between goroutines.
When an operation exceeds a specified timeout, or when a user aborts a request, the context should be canceled. Every goroutine involved in the operation should periodically check ctx.Done() or use a select statement that includes <-ctx.Done().
func fetchResource(ctx context.Context, url string) (string, error) {
req, err := http.NewRequestWithContext(ctx, http.MethodGet, url, nil)
if err != nil {
return "", err
}
client := http.DefaultClient
resp, err := client.Do(req)
if err != nil {
return "", err
}
defer resp.Body.Close()
// process response...
return "Success", nil
}
This ensures that network requests do not hang indefinitely and resources are promptly released when no longer needed.
5. The Semaphore Pattern
Sometimes you need to control access to a specific resource, limiting the number of goroutines that can access it concurrently, independent of the number of active workers. While a sync.Mutex allows only one goroutine to access a resource at a time, a semaphore allows up to N goroutines.
In Go, a buffered channel is commonly used as a semaphore.
func main() {
concurrencyLimit := 5
semaphore := make(chan struct{}, concurrencyLimit)
var wg sync.WaitGroup
for i := 0; i < 20; i++ {
wg.Add(1)
go func(id int) {
defer wg.Done()
// Acquire semaphore
semaphore <- struct{}{}
// Critical section
fmt.Printf("Goroutine %d executing\n", id)
time.Sleep(time.Second) // Simulate work
// Release semaphore
<-semaphore
}(i)
}
wg.Wait()
}
This pattern is highly effective for throttling requests to third-party APIs or managing connections to legacy databases that cannot handle high concurrency.
6. The Or-Done Channel Pattern
When dealing with multiple channels or when wrapping functions that you don't control, you often need a way to ensure that you aren't leaking goroutines if a cancellation occurs before a channel operation completes. The or-done pattern elegantly solves this by wrapping the read operation.
func orDone(ctx context.Context, c <-chan interface{}) <-chan interface{} {
valStream := make(chan interface{})
go func() {
defer close(valStream)
for {
select {
case <-ctx.Done():
return
case v, ok := <-c:
if !ok {
return
}
select {
case valStream <- v:
case <-ctx.Done():
}
}
}
}()
return valStream
}
This allows you to write cleaner loops when consuming channels, abstracting away the boilerplate of checking ctx.Done() on every read.
7. The Tee Channel Pattern
Sometimes you need to split a stream of data originating from a single channel into two separate streams, allowing independent processing paths. This is known as a tee channel (analogous to the tee command in Unix).
func tee(ctx context.Context, in <-chan interface{}) (<-chan interface{}, <-chan interface{}) {
out1 := make(chan interface{})
out2 := make(chan interface{})
go func() {
defer close(out1)
defer close(out2)
for val := range orDone(ctx, in) {
// Create local copies of channels to track which ones we've sent to
var out1Local, out2Local = out1, out2
for i := 0; i < 2; i++ {
select {
case <-ctx.Done():
return
case out1Local <- val:
out1Local = nil // Prevent sending to this channel again
case out2Local <- val:
out2Local = nil // Prevent sending to this channel again
}
}
}
}()
return out1, out2
}
The tee channel guarantees that each value received from the in channel is sent to both out1 and out2 before the next value is processed. Setting the local channel reference to nil after a successful send is a clever trick in Go; sending to a nil channel blocks forever, meaning the select statement will only evaluate the remaining non-nil cases for the rest of the loop iteration. This pattern is exceptionally useful for duplicating event streams, such as sending telemetry data to a metrics service while simultaneously logging it to a file.
8. ErrGroup for Coordinated Tasks
While sync.WaitGroup is excellent for waiting on a collection of goroutines, it falls short when you need to handle errors or cancel remaining tasks if one task fails. The golang.org/x/sync/errgroup package provides synchronization, error propagation, and Context cancellation for groups of goroutines working on subtasks of a common task.
package main
import (
"context"
"fmt"
"golang.org/x/sync/errgroup"
)
func main() {
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
g, ctx := errgroup.WithContext(ctx)
urls := []string{
"http://example.com",
"http://example.net",
"http://example.org",
}
for _, url := range urls {
// Launch a goroutine to fetch the URL.
url := url // https://golang.org/doc/faq#closures_and_goroutines
g.Go(func() error {
// Use the context returned by errgroup.WithContext.
// If any goroutine returns an error, this context will be canceled.
// return fetchResource(ctx, url)
fmt.Printf("Fetching %s\n", url)
return nil
})
}
// Wait for all HTTP fetches to complete.
if err := g.Wait(); err != nil {
fmt.Printf("Failed to fetch all URLs: %v\n", err)
} else {
fmt.Println("Successfully fetched all URLs.")
}
}
Using errgroup is significantly cleaner than managing a separate error channel and cancellation logic manually. It's the recommended approach whenever you have multiple concurrent tasks that are part of a larger, logical operation where the success of the whole depends on the success of its parts.
Conclusion
Mastering advanced Go concurrency requires moving beyond the basic go keyword and thinking deeply about data flow, resource management, and lifecycle control. By employing patterns like pipelines, fan-out/fan-in, worker pools, robust context cancellation, and semaphores, you can design systems that are not only highly concurrent but also resilient, predictable, and maintainable. These patterns form the architectural bedrock of high-performance Go applications in production.
You Might Also Like
Free In-Browser Developer Tools
Clean AI CLI logs, build cron expressions, decode JWTs, and calculate chmod permissions offline.
Related Articles

gRPC vs ConnectRPC: Modern Microservices and Browser-Native Protobuf
An architectural evaluation of gRPC vs ConnectRPC in TypeScript and Go. Explore HTTP/1.1 vs HTTP/2 streaming, browser clients without Envoy proxies, and p99 RPC latency.
Read more
Quantum Computing for Developers
A developer guide to quantum computing: write quantum algorithms with Qiskit, understand quantum gates, and simulate circuits on classical hardware.
Read more
Distributed Tracing with OpenTelemetry
Instrument microservices with OpenTelemetry distributed tracing: trace cross-service context propagation, latency bottlenecks, and export to Jaeger.
Read more