Go in 2 Weeks #06: Goroutines, Buffered vs Unbuffered Channels, select and Worker Pools
Key takeaways
Goroutines are cheap functions scheduled by the Go runtime onto a small pool of OS threads; channels pass values between them and synchronise at the same time. This post covers unbuffered vs buffered channels, select, worker pools and pipelines, plus the deadlocks and goroutine leaks C++ developers hit first.
Series overview
Go in 2 Weeks #06
This post covers Days 10–11 of the two-week Go curriculum for C++ developers.
Previous: #05 Error handling | Next: #07 Modules & testing
Introduction: why goroutines feel different from std::thread
In C++, std::thread maps one-to-one to an OS thread. Each one gets a kernel scheduling entity and a stack sized up front (commonly 8 MB of reserved virtual memory on Linux, 1 MB on Windows by default). Creating thousands is possible but wasteful, so C++ servers use thread pools and hand work to them through queues guarded by mutexes and condition variables.
A goroutine is a function running under the Go runtime’s scheduler. The runtime keeps a small pool of OS threads, and roughly GOMAXPROCS of them (defaulting to the number of available CPUs) run Go code at any moment. Goroutines start with a small stack (a few kilobytes) that the runtime grows by copying when needed. When a goroutine blocks on a channel, a mutex, a timer or network I/O, the scheduler parks it and runs another one on the same thread. The practical result is that “one goroutine per connection” or “one goroutine per task” is the normal design, not an anti-pattern.
Cheap does not mean free, though. Each goroutine holds its stack and everything it references until it returns. The runtime never kills a goroutine for you, so a goroutine blocked forever is a memory leak. Most of the bugs in this post come from that single fact.
Goroutines and WaitGroup
C++ vs Go
#include <iostream>
#include <thread>
#include <vector>
void worker(int id) { std::cout << "Worker " << id << " running\n"; }
int main() {
std::vector<std::thread> threads;
for (int i = 0; i < 10; i++) threads.emplace_back(worker, i);
for (auto& t : threads) t.join();
}
package main
import (
"fmt"
"sync"
)
func worker(id int) { fmt.Printf("Worker %d running\n", id) }
func main() {
var wg sync.WaitGroup
for i := 0; i < 10000; i++ {
wg.Add(1)
go func() {
defer wg.Done()
worker(i)
}()
}
wg.Wait()
}
There is no join() on a goroutine and no handle to one. go f() returns immediately, and if main returns, the program exits with every other goroutine still running. sync.WaitGroup is the usual way to wait: Add before starting, Done when finished (via defer so it runs on every return path), Wait to block until the counter reaches zero.
Two details that trip people up:
- Call
wg.Add(1)beforego, not inside the goroutine. OtherwiseWaitcan run before any goroutine has incremented the counter and return immediately. - Loop variables. Since Go 1.22 each loop iteration gets a fresh
i, so the closure above is correct. In older Go versions all closures shared one variable and typically printed the final value; you will still seego func(id int) { ... }(i)in older code for that reason.
Go 1.25 added wg.Go(f), which does the Add/Done bookkeeping for you:
var wg sync.WaitGroup
var mu sync.Mutex
total := 0
for i := 1; i <= 100; i++ {
wg.Go(func() {
mu.Lock()
total += i
mu.Unlock()
})
}
wg.Wait()
fmt.Println(total) // 5050
Channels: communicating instead of sharing
C++ mutex vs Go channel
std::mutex mtx;
std::vector<int> results;
void worker(int id) {
int result = id * 2;
std::lock_guard<std::mutex> lock(mtx);
results.push_back(result);
}
func worker(id int, ch chan<- int) {
ch <- id * 2
}
func main() {
ch := make(chan int)
for i := 0; i < 10; i++ {
go worker(i, ch)
}
results := make([]int, 0, 10)
for i := 0; i < 10; i++ {
results = append(results, <-ch)
}
fmt.Println(len(results))
}
In the Go version only main touches results, so there is nothing to lock. Workers hand values over the channel and the channel does the synchronisation. That is the meaning of the Go proverb “do not communicate by sharing memory; share memory by communicating”. It is guidance, not a rule: for a simple counter or a cache map, sync.Mutex or sync/atomic is shorter and faster, and the standard library uses mutexes heavily. Use channels when you are transferring ownership of data or signalling events between goroutines.
Channel basics and close
ch := make(chan int, 1)
ch <- 42
close(ch)
v, ok := <-ch
fmt.Println(v, ok) // 42 true (buffered values are still delivered after close)
v, ok = <-ch
fmt.Println(v, ok) // 0 false (closed and drained)
close means “no more values will be sent”. It is not required for garbage collection; an unreferenced channel is collected whether or not it was closed. Close a channel only when receivers need that signal, typically to end a for v := range ch loop.
Directional types
func sender(ch chan<- int) { ch <- 1; close(ch) }
func receiver(ch <-chan int) {
for v := range ch {
fmt.Println(v)
}
}
chan<- int is send-only and <-chan int is receive-only. A bidirectional channel converts implicitly to either, and the compiler rejects a receive on a send-only channel. Put these in function signatures: they document who owns which end and catch mistakes such as a consumer closing the channel.
Unbuffered vs buffered channels
// Unbuffered: the send completes only when a receiver takes the value
ch := make(chan int)
go func() { ch <- 1 }()
<-ch
// Buffered: up to 3 sends complete with no receiver
b := make(chan int, 3)
b <- 1
b <- 2
b <- 3
// a 4th send would block until someone receives
An unbuffered channel is a handoff: sender and receiver meet, and when the send returns you know the receiver has the value. That guarantee is often exactly what you want, for example to signal “I have started” or “I have finished”.
A buffered channel decouples the two sides up to its capacity. It smooths out bursts when producer and consumer run at different rates and lets a producer continue without waiting for each item. It does not add throughput if the consumer is permanently slower; the buffer just fills and you are back to blocking.
The trap is reaching for a buffer to “fix” a deadlock. If a program only works with make(chan T, 10), it usually means the design relies on nobody ever sending an 11th value. That holds until the input grows. Choose a capacity for a reason you can state (one slot per worker, a result slot so a sender never blocks), not by trial and error.
select: waiting on several channels
select blocks until one of its cases can proceed; if several are ready it picks one at random, which prevents one busy channel from starving another.
select {
case msg := <-ch1:
fmt.Println("ch1:", msg)
case msg := <-ch2:
fmt.Println("ch2:", msg)
}
Timeouts
select {
case result := <-ch:
fmt.Println("Received:", result)
case <-time.After(1 * time.Second):
fmt.Println("Timeout!")
}
Non-blocking operations with default
select {
case ch <- 1:
fmt.Println("Sent")
default:
fmt.Println("Would block")
}
A select with a default case never blocks. That is useful for “drop the message if the queue is full”, but a for { select { ... default: } } loop with nothing blocking inside it is a busy loop that burns a full CPU core.
A nil channel blocks forever in both directions, and inside select that makes its case permanently not ready. Setting a channel variable to nil after it is closed is the idiomatic way to disable one case while continuing to read the others.
Concurrency patterns
Worker pool
A fixed number of goroutines read from a shared jobs channel. Here the pool size bounds concurrency (for example, the number of simultaneous database connections), and the closing sequence is the important part:
func worker(id int, jobs <-chan int, results chan<- int, wg *sync.WaitGroup) {
defer wg.Done()
for job := range jobs {
results <- job * 2
}
}
func main() {
jobs := make(chan int)
results := make(chan int)
var wg sync.WaitGroup
for w := 1; w <= 3; w++ {
wg.Add(1)
go worker(w, jobs, results, &wg)
}
go func() {
for j := 1; j <= 9; j++ {
jobs <- j
}
close(jobs) // workers' range loops end
}()
go func() {
wg.Wait()
close(results) // only after every worker has stopped sending
}()
sum := 0
for r := range results {
sum += r
}
fmt.Println("sum:", sum) // sum: 90
}
Closing results directly from a worker would panic when another worker sends afterwards. The separate “wait then close” goroutine is the standard way to close a channel that has many senders.
Pipeline
Each stage owns its output channel, closes it when done, and returns it as receive-only:
func generator(nums ...int) <-chan int {
out := make(chan int)
go func() {
defer close(out)
for _, n := range nums {
out <- n
}
}()
return out
}
func square(in <-chan int) <-chan int {
out := make(chan int)
go func() {
defer close(out)
for n := range in {
out <- n * n
}
}()
return out
}
func main() {
for v := range square(generator(1, 2, 3, 4)) {
fmt.Println(v) // 1 4 9 16
}
}
The weakness of this simple form is early exit: if main stops reading after the first value, both stage goroutines block forever on their sends. Real pipelines pass a context.Context (or a done channel) into every stage and select on it next to each send.
Cancellation with context
func slowSquare(ctx context.Context, n int) (int, error) {
select {
case <-time.After(time.Duration(n) * 100 * time.Millisecond):
return n * n, nil
case <-ctx.Done():
return 0, ctx.Err()
}
}
func main() {
ctx, cancel := context.WithTimeout(context.Background(), 250*time.Millisecond)
defer cancel()
for _, n := range []int{1, 2, 3} {
v, err := slowSquare(ctx, n)
fmt.Println(n, v, err)
}
}
1 1 <nil>
2 0 context deadline exceeded
3 0 context deadline exceeded
The deadline covers the whole call chain: after the first call uses 100 ms, the second needs 200 ms more and misses the 250 ms budget. Always defer cancel(); it releases the context’s timer even when the deadline is never reached. The context cancellation guide goes deeper.
Errors you will actually see
Deadlock. When every goroutine is blocked, the runtime detects it and aborts:
func main() {
ch := make(chan int)
ch <- 1 // no receiver, ever
}
fatal error: all goroutines are asleep - deadlock!
goroutine 1 [chan send]:
main.main()
The goroutine state in brackets (chan send, chan receive, semacquire for a WaitGroup or mutex) tells you what it was waiting on. The detector only fires when all goroutines are stuck. In a server with an HTTP listener running, a stuck goroutine never triggers it and simply leaks.
Send on a closed channel panics:
panic: send on closed channel
It almost always means a receiver or one of several senders closed the channel. Move the close to the single owner, or to a goroutine that waits for all senders.
Data races do not produce an error at all by default: the program just gives wrong results occasionally. Run tests with go test -race (it requires cgo, so a C compiler, on some platforms) to detect unsynchronised access to shared variables.
I find goroutine leaks the most insidious of these, because nothing crashes. The typical shape is the timeout pattern with an unbuffered result channel: the caller gives up after the timeout and returns, and the worker, finishing later, blocks forever on result <- value because nobody will ever receive. Each timed-out request leaves one goroutine behind, and the process’s memory climbs slowly until someone looks at runtime.NumGoroutine() or a goroutine profile from net/http/pprof and finds thousands parked on the same line. Giving the result channel a buffer of one, as in Exercise 3 below, is the fix.
The other habit I brought from C++ that hurt was treating close like a destructor that must always be called. Closing channels “for cleanliness” from the receiving side is how you get the send-on-closed-channel panic in production, usually on a rare timing path that tests never hit.
Exercises
Exercise 1: parallel downloads
package main
import (
"fmt"
"io"
"net/http"
"sync"
)
func download(url string, wg *sync.WaitGroup, results chan<- string) {
defer wg.Done()
resp, err := http.Get(url)
if err != nil {
results <- fmt.Sprintf("%s: error - %v", url, err)
return
}
defer resp.Body.Close()
body, err := io.ReadAll(resp.Body)
if err != nil {
results <- fmt.Sprintf("%s: read error - %v", url, err)
return
}
results <- fmt.Sprintf("%s: %d bytes", url, len(body))
}
func main() {
urls := []string{"https://go.dev", "https://github.com", "https://stackoverflow.com"}
var wg sync.WaitGroup
results := make(chan string, len(urls))
for _, url := range urls {
wg.Add(1)
go download(url, &wg, results)
}
go func() {
wg.Wait()
close(results)
}()
for result := range results {
fmt.Println(result)
}
}
Results arrive in completion order, not in the order of urls. If order matters, send an index with each result and write into a pre-sized slice. For real use, set a timeout: the default http.Client has none.
Exercise 2: rate limiter
func rateLimiter(requests <-chan int, rate time.Duration, done chan<- struct{}) {
ticker := time.NewTicker(rate)
defer ticker.Stop()
for req := range requests {
<-ticker.C
fmt.Printf("Processing request %d at %v\n", req, time.Now().Format("15:04:05.000"))
}
close(done)
}
func main() {
requests := make(chan int, 10)
done := make(chan struct{})
go rateLimiter(requests, 500*time.Millisecond, done)
for i := 1; i <= 5; i++ {
requests <- i
}
close(requests)
<-done // wait for the limiter instead of guessing with time.Sleep
}
Exercise 3: timeout without leaking
func longRunningTask(result chan<- string) {
time.Sleep(3 * time.Second)
result <- "Task completed" // never blocks: the buffer has room
}
func main() {
result := make(chan string, 1) // buffer of 1 prevents the leak
go longRunningTask(result)
select {
case res := <-result:
fmt.Println(res)
case <-time.After(2 * time.Second):
fmt.Println("Timeout: task took too long")
}
}
This prints Timeout: task took too long. The task goroutine still runs to completion; the buffer only guarantees it can finish. To actually stop the work, pass a context and check it.
Exercise 4: a counter guarded two ways
type SafeCounter1 struct {
mu sync.Mutex
count map[string]int
}
func (c *SafeCounter1) Inc(key string) {
c.mu.Lock()
defer c.mu.Unlock()
c.count[key]++
}
// SafeCounter2: one goroutine owns the map; others send it operations
type SafeCounter2 struct {
ops chan func(map[string]int)
}
func NewSafeCounter2() *SafeCounter2 {
c := &SafeCounter2{ops: make(chan func(map[string]int))}
go func() {
count := make(map[string]int)
for op := range c.ops {
op(count)
}
}()
return c
}
func (c *SafeCounter2) Inc(key string) {
c.ops <- func(count map[string]int) { count[key]++ }
}
func (c *SafeCounter2) Value(key string) int {
result := make(chan int)
c.ops <- func(count map[string]int) { result <- count[key] }
return <-result
}
Both versions are correct after 1000 concurrent Inc calls. The mutex version is simpler and faster for this job; the owner-goroutine version shows the pattern that scales to state with complex invariants, at the cost of a channel round trip per operation and a goroutine that lives as long as the counter (close ops to stop it). Plain Go maps are not safe for concurrent writes, and the runtime detects many such cases with fatal error: concurrent map writes.
Concurrency is not parallelism: the GOMAXPROCS=1 check
Goroutines and channels are how you structure a program as independent tasks (concurrency). Whether those tasks run at the same instant on different cores (parallelism) is up to the runtime and GOMAXPROCS. A correctly synchronised program therefore still produces the right result with GOMAXPROCS=1, where only one goroutine runs Go code at a time. If it hangs or prints something different there, it relies on timing, for example on a goroutine happening to finish before main reads a variable instead of on a wg.Wait() or a channel receive. The single-threaded run tends to hide data races rather than expose them, so pair it with go test -race.
Next: Go modules and testing — dependency management and go test.
Series navigation
Go in 2 weeks: #01 · #02 · #03 · #04 · #05 · #06 · #07 · #08 · #09