Skip to content

Go 并发:goroutine、channel 与可收束的并发结构

面向有 Java 经验的开发者,基于 Go 1.26。

go f()make(chan T) 几分钟就能学会,但能启动并发任务,离可靠地管理它们还有很远。谁拥有 goroutine?谁负责停止它?错误怎样传回?发送方被取消时会不会永久阻塞?缓冲区满了以后,系统选择阻塞、丢弃还是扩容?函数返回之前,自己启动的任务是否已经收束?这些问题必须在代码中找到答案。

Go 并发的重点不是“轻量线程很多”,而是把通信、同步和生命周期协议写清楚。channel 既能传值,也能建立 happens-before 关系,但它不会替你解决所有权、取消、死锁、背压和数据竞争。后面的内容会先说明语言语义,再把这些语义落到可收束的并发结构中。

目录

1. 先建立并发心智模型

一个 goroutine 是由 Go 运行时调度的独立执行单元。它不是一个带 joininterrupt 方法的对象;go 语句也不返回句柄。启动之后,调用方如果需要等待、取消或接收结果,必须另外设计协议。

go
go serve(conn)

一个 channel 是带类型的通信通道:

go
ch := make(chan Result)
ch <- result // 发送
result := <-ch // 接收

阅读后面的代码时,先抓住三条原则:

  1. 启动 goroutine 时,同时决定退出路径。
  2. 共享内存需要同步;channel 是同步手段之一,不是数据竞争豁免证。
  3. 把容量和过载行为当作 API 协议,而不是随手写一个缓冲数字。

“不要通过共享内存来通信;通过通信来共享内存”是一条设计提示,不是禁止使用 sync.Mutex 的教条。计数器、缓存和短临界区通常用锁更直接;流式传递、所有权转移和任务协调则更适合 channel。

2. goroutine 的创建与生命周期

2.1 go 语句

函数调用前加 go

go
go fetch(url)

go func() {
	fmt.Println("background")
}()

调用表达式的函数值和参数在启动方 goroutine 中求值,然后新 goroutine 执行函数体:

go
for _, url := range urls {
	go fetch(url)
}

Go 1.22 起,for 循环声明的迭代变量按每次迭代拥有新变量的语义处理,Go 1.26 代码不再有过去那个经典闭包捕获问题。显式传参仍能让数据边界更清楚:

go
for _, url := range urls {
	go func(u string) {
		fetch(u)
	}(url)
}

不要为了兼容老文章而在现代代码中机械写 url := url,但维护旧 go.mod 语言版本的项目时应确认模块采用的循环语义。

2.2 goroutine 何时结束

函数返回或调用 runtime.Goexit 时,当前 goroutine 结束。没有安全的外部“强杀 goroutine”操作。panic 若未在当前 goroutine 内恢复,会使整个进程崩溃,而不是只杀掉这一条任务。

main 返回时进程直接退出,其他 goroutine 不会被等待:

go
func main() {
	go func() {
		time.Sleep(time.Second)
		fmt.Println("可能永远看不到")
	}()
}

生产代码不能用 time.Sleep 猜任务是否完成,应使用 channel、sync.WaitGroup 或更高层的任务组。

2.3 每个 goroutine 都需要所有者

启动点应能回答:

  • 谁发出取消信号?
  • goroutine 正在阻塞时如何响应取消?
  • 谁等待它结束?
  • 错误交给谁?
  • 它持有的文件、连接、ticker 由谁释放?
  • 允许活到进程结束,还是必须在请求结束前退出?

即便是进程级指标收集器这类常驻 goroutine,也应有明确的进程生命周期和关闭入口。“反正很轻量”不能成为泄漏的理由。

3. 等待 goroutine

3.1 用结果 channel 等待单个任务

go
resultCh := make(chan Result, 1)

go func() {
	resultCh <- compute()
}()

result := <-resultCh

缓冲为 1 的结果 channel 有一个重要性质:即使调用方因取消提前离开,工作 goroutine 完成发送时也不必等待接收者。若结果包含大对象且无人接收,应考虑生命周期和及时释放。

3.2 sync.WaitGroup

多个无返回结果的任务可以使用 WaitGroup

go
var wg sync.WaitGroup

for _, item := range items {
	item := item
	wg.Add(1)
	go func() {
		defer wg.Done()
		process(item)
	}()
}

wg.Wait()

Add 必须在启动 goroutine 前完成,避免 Wait 在计数增加前返回。WaitGroup 使用后不能复制。

Go 1.25 起可用 WaitGroup.Go,Go 1.26 中可以写:

go
var wg sync.WaitGroup
for _, item := range items {
	item := item
	wg.Go(func() {
		process(item)
	})
}
wg.Wait()

Go 会启动函数并在返回时完成计数;其文档要求传入函数不得 panic。若任务需要返回错误、首错取消或限制并发,仅有 WaitGroup 还不够,应增加明确的错误与取消协议。

3.3 WaitGroup 不是取消器

Wait 只能等待,不能通知任务退出。若任务卡在网络读取、channel 发送或无限循环中,Wait 会永远等待。通常把它和 context.Context 或关闭信号结合。

4. Go 内存模型基础

内存模型回答:一个 goroutine 的写入,什么时候保证能被另一个 goroutine 观察到。最实用的结论是:

多个 goroutine 并发访问同一内存,其中至少一个写入时,必须使用 channel、锁、原子操作等同步机制建立顺序。

没有数据竞争的程序具有顺序一致的解释(DRF-SC):结果可视为各 goroutine 操作按某种顺序交错。带数据竞争的程序不是“只会读到旧值”;接口、切片、map、字符串这类多机器字结构甚至可能出现不一致的组合,造成严重错误。

4.1 happens-before

若操作 A happens-before 操作 B,B 对相关内存读取就能观察到按规则可见的 A 写入。常用同步规则包括:

  • 启动 goroutine 的 go 语句,先于新 goroutine 开始执行;
  • channel 的一次发送,先于对应接收完成;
  • 关闭 channel,先于因该关闭而返回零值的接收;
  • 对无缓冲 channel,一次接收先于对应发送完成;
  • 容量为 C 的 channel,第 k 次接收先于第 k+C 次发送完成;
  • mutex 的一次 Unlock 先于后续相应 Lock 返回;
  • sync.Once 中函数完成,先于任意 Do 调用返回。

例子:

go
var message string
done := make(chan struct{})

go func() {
	message = "ready"
	close(done)
}()

<-done
fmt.Println(message) // 保证看到 "ready"

写入 message 先于 close(done),关闭又先于接收完成,于是形成完整的 happens-before 链。

4.2 goroutine 结束不是同步

go
var message string

go func() {
	message = "ready"
}()

fmt.Println(message) // 数据竞争;不能靠“它应该跑完了”推理

goroutine 消失本身不保证其写入对其他 goroutine 可见。必须等待一个同步事件。

4.3 原子不变量

即使单个机器字读写在某平台上看似原子,也不代表无竞争,更不代表多个字段的一致性。需要原子计数可使用 sync/atomic;需要维护复合不变量通常使用 mutex;需要转移数据所有权可使用 channel。

所有并发测试都应考虑 race detector:

bash
go test -race ./...

它只能发现实际运行路径触发的竞争,不能证明程序在所有路径上无竞争,因此测试覆盖和设计审查仍然重要。

5. channel 类型与创建

channel 类型写成 chan T

go
var jobs chan Job
jobs = make(chan Job)

buffered := make(chan Job, 32)

chan T 是双向 channel,<-chan T 只接收,chan<- T 只发送。channel 值可以复制,副本指向同一个运行时通道。

5.1 零值是 nil

go
var ch chan int
fmt.Println(ch == nil) // true

对 nil channel 发送或接收会永久阻塞;关闭 nil channel会 panic:

go
// ch <- 1  // 永久阻塞
// <-ch     // 永久阻塞
// close(ch) // panic

nil channel 在 select 中很有用,因为对应 case 永远不会就绪,可以动态关闭某个分支;在普通顺序代码里,它通常表示忘记初始化。

5.2 channel 可比较

channel 可以和 nil 比较,也可以比较是否指向同一通道:

go
a := make(chan int)
b := a
fmt.Println(a == b) // true

这个能力适合身份判断,不应用来推导队列内容或状态。

6. 无缓冲 channel

go
ch := make(chan string)

无缓冲 channel 的发送必须和接收会合。发送方与接收方都到达通信点时,值才交付,双方才继续:

go
go func() {
	ch <- "done"
}()

msg := <-ch

它同时完成:

  • 值传递;
  • 发送与接收之间的同步;
  • 对双方执行进度的协调。

不能简单说“发送一定阻塞到接收方把后续逻辑执行完”。它只等待对应通信完成;接收方拿到值后做什么,与发送方可以再次并发。

6.1 用空结构体表达纯信号

go
done := make(chan struct{})
go func() {
	work()
	close(done)
}()
<-done

struct{} 不携带业务数据,表达“这里只需要事件”。一次性完成通知通常用关闭比发送更灵活,因为关闭能广播给所有等待者。

6.2 无缓冲不等于更正确

无缓冲使交接点明确,但也让两个组件的执行节奏紧耦合。若消费者暂时退出,生产者可能泄漏;若环形依赖中所有人先发送,就会死锁。选择无缓冲应源于协议需要同步交接,而不是默认认为“缓冲会隐藏问题”。

7. 缓冲 channel 与背压

go
jobs := make(chan Job, 100)

发送在缓冲未满时可立即完成,接收在缓冲非空时可立即完成。满时发送阻塞,空时接收阻塞。

go
fmt.Println(len(jobs), cap(jobs))

len 只是一瞬间的排队数量,不能用于“先判断再发送”的并发正确性:

go
// if len(ch) < cap(ch) { ch <- v } // 条件可能立刻过期

非阻塞尝试应使用 select

go
select {
case jobs <- job:
	// accepted
default:
	// overloaded
}

7.1 缓冲容量是系统策略

缓冲并不创造吞吐,只吸收生产和消费之间的短期波动。容量过大可能:

  • 延迟过载暴露;
  • 增加内存占用;
  • 延长排队时间,任务到手时已经过期;
  • 让关闭过程等待大量陈旧任务;
  • 掩盖消费者持续低于生产者的事实。

容量过小则可能让生产者频繁阻塞,降低并行度。应根据到达率、处理速率、允许等待时间和内存预算决定,并通过队列长度、拒绝数和端到端延迟监控验证。

Little's Law 可作容量估算起点:稳定系统中平均在途数量约等于到达率乘平均停留时间。但真实流量有突发,仍需压测。

7.2 channel 作为信号量

容量 channel 可以限制并发:

go
limit := make(chan struct{}, 8)

for _, item := range items {
	item := item
	go func() {
		limit <- struct{}{} // acquire
		defer func() { <-limit }() // release
		process(item)
	}()
}

这会为每个 item 先创建 goroutine,即使它们随后阻塞在限流器上。任务很多或无界时,应使用固定 worker pool,避免堆积大量等待 goroutine。

内存模型还规定:容量为 C 的 channel,第 k 次接收先于第 k+C 次发送完成,这正是用缓冲 channel 建模计数信号量的同步基础。

8. 关闭 channel

close(ch) 表示“今后不会再发送值”,不是“销毁 channel”,也不是给接收者发送一个特殊值。

go
close(ch)

接收会先取完缓冲区已有值,之后立即得到元素零值和 ok == false

go
v, ok := <-ch
if !ok {
	// 已关闭且已无值
}

range 会持续接收,直到 channel 关闭并排空:

go
for v := range ch {
	consume(v)
}

8.1 关闭规则

  • 向已关闭 channel 发送:panic;
  • 再次关闭已关闭 channel:panic;
  • 关闭 nil channel:panic;
  • 从已关闭且排空的 channel 接收:立即返回零值、false;
  • 从仅接收 channel 不能调用 close
  • channel 不必总是关闭,只有接收方需要结束信号时才需要。

8.2 由发送方关闭

一般由确定“不会再有发送”的一方关闭。接收方擅自关闭会和仍在运行的发送方竞争,导致 panic。

go
func produce(ctx context.Context, out chan<- Item) {
	defer close(out)
	for {
		item, ok := next()
		if !ok {
			return
		}
		select {
		case out <- item:
		case <-ctx.Done():
			return
		}
	}
}

多个发送者时,任何单个发送者都不知道其他人是否结束。应让一个协调者等待全部发送者,再关闭:

go
var wg sync.WaitGroup
for _, source := range sources {
	source := source
	wg.Go(func() {
		for v := range source {
			out <- v
		}
	})
}

go func() {
	wg.Wait()
	close(out)
}()

8.3 不要用 recover 实现“安全关闭”

包装 close 并恢复 panic 不能解决关闭与发送间的数据竞争,也掩盖所有权不清。正确方案是单一关闭者、锁保护状态,或重新设计协议。

9. 单向 channel 与所有权

函数签名应表达方向:

go
func generate(ctx context.Context) <-chan int
func consume(ctx context.Context, in <-chan int) error
func publish(out chan<- Event, e Event)

双向 channel 可隐式收窄为单向,不能反向拓宽:

go
ch := make(chan int)
var recv <-chan int = ch
var send chan<- int = ch

方向是编译期约束,不会创建新通道。它能防止消费者误发、生产者误收,也让“谁能关闭”更清楚。只有持有发送能力的一方可能关闭,但拥有 chan<- T 并不等于一定拥有关闭权,API 文档仍要约定。

返回只读 channel 的函数通常由内部 goroutine 创建并关闭它。调用者负责持续接收或取消,否则内部发送可能阻塞。隐藏 goroutine 的 API 必须同时暴露停止方式。

10. select

select 在多个 channel 操作中等待:

go
select {
case v := <-values:
	use(v)
case err := <-errorsCh:
	return err
case <-ctx.Done():
	return ctx.Err()
}

若多个 case 同时就绪,会用伪随机的均匀选择挑一个;不能依赖源码顺序表达优先级。若没有 case 就绪:

  • default:立即执行 default;
  • default:阻塞;
  • select {}:永久阻塞。

进入 select 时,各 case 的 channel 操作数和发送右侧表达式会按规范求值一次。不要在其中放有意外副作用的复杂表达式。

10.1 非阻塞发送和接收

go
select {
case v := <-ch:
	handle(v)
default:
	// 当前没有值
}
go
select {
case ch <- v:
	// 发送成功
default:
	// 当前不可发送
}

“当前不可用”不是“永远不可用”。这种模式适合指标、尽力通知或明确允许丢弃的负载;不能用来悄悄丢失订单、账务事件等可靠消息。

10.2 nil channel 动态禁用 case

go
var out chan<- Item
var next Item

if len(queue) > 0 {
	out = destination
	next = queue[0]
}

select {
case out <- next:
	queue = queue[1:]
case item := <-source:
	queue = append(queue, item)
}

out == nil 时发送 case 不会被选择。这是状态机循环中控制分支的常用技巧。要确保至少有一个分支最终可进展,否则会永久阻塞。

10.3 关闭 channel 在 select 中会一直就绪

已关闭 channel 接收立即返回。循环中若不处理,它会持续被选中并形成忙循环:

go
case v, ok := <-ch:
	if !ok {
		ch = nil // 禁用这个 case
		continue
	}

11. 超时、定时器与心跳

一次性超时可以写:

go
select {
case result := <-resultCh:
	return result, nil
case <-time.After(500 * time.Millisecond):
	return Result{}, context.DeadlineExceeded
}

在高频循环里通常创建并复用 time.Timer,便于停止、重置和控制资源:

go
timer := time.NewTimer(timeout)
defer timer.Stop()

select {
case v := <-ch:
	return v, nil
case <-timer.C:
	return zero, ErrTimeout
}

跨调用链的超时优先用 context.WithTimeout,让下游数据库、HTTP 和 RPC 共享同一截止时间,而不是每层各造一个彼此无关的 timer。

周期任务使用 time.NewTicker 并在不再需要时 Stop

go
ticker := time.NewTicker(time.Second)
defer ticker.Stop()

for {
	select {
	case <-ticker.C:
		heartbeat()
	case <-ctx.Done():
		return
	}
}

ticker 不是精确实时调度器;接收方太慢时 tick 可能被调整或丢弃。对账、结算等任务应基于实际时间和持久状态计算,不应假定“每个 tick 都收到”。

12. context 取消树

context.Context 携带取消、截止时间和请求范围值。惯例是作为第一个参数显式传入:

go
func Load(ctx context.Context, id UserID) (User, error)

不要把 Context 存进结构体,不要传 nil;没有合适上下文时使用 context.Background()context.TODO()

12.1 派生和取消

go
ctx, cancel := context.WithTimeout(parent, 2*time.Second)
defer cancel()

user, err := repo.Load(ctx, id)

即使操作提前完成也要调用 cancel,以释放关联资源。子上下文取消不会取消父上下文;父取消会传播到所有后代。

goroutine 应在可能长时间阻塞的位置监听:

go
select {
case out <- value:
case <-ctx.Done():
	return ctx.Err()
}

只在循环顶部偶尔检查一次不够。如果随后可能永久阻塞在 channel 发送,取消仍无法生效。

12.2 DoneErr 和取消原因

ctx.Done() 返回一个接收 channel,取消时关闭。调用者不能关闭它。ctx.Err() 返回 context.Canceledcontext.DeadlineExceeded

需要保留业务取消原因时:

go
ctx, cancel := context.WithCancelCause(parent)
cancel(fmt.Errorf("dependency failed: %w", err))

fmt.Println(context.Cause(ctx))

派生超时也有 WithTimeoutCauseWithDeadlineCause。用 Cause 能保留“为什么整组任务停止”,但 API 仍应规定最终返回哪个错误。

12.3 Context value 的边界

Context value 适合请求范围、跨 API 边界传递的数据,如 trace ID 或认证主体。不要用它传可选函数参数、数据库连接或业务必填字段。键应使用包内定义的不可导出类型,避免冲突:

go
type traceIDKey struct{}
ctx = context.WithValue(ctx, traceIDKey{}, traceID)

13. pipeline

pipeline 把数据分阶段处理,每个阶段接收上游 channel 并返回下游 channel:

go
func generate(ctx context.Context, values ...int) <-chan int {
	out := make(chan int)
	go func() {
		defer close(out)
		for _, v := range values {
			select {
			case out <- v:
			case <-ctx.Done():
				return
			}
		}
	}()
	return out
}

func square(ctx context.Context, in <-chan int) <-chan int {
	out := make(chan int)
	go func() {
		defer close(out)
		for v := range in {
			select {
			case out <- v * v:
			case <-ctx.Done():
				return
			}
		}
	}()
	return out
}

使用:

go
ctx, cancel := context.WithCancel(context.Background())
defer cancel()

for v := range square(ctx, generate(ctx, 2, 3, 4)) {
	fmt.Println(v)
}

可靠阶段的约定:

  • 只接收上游 channel;
  • 自己创建并关闭下游;
  • 上游关闭后完成;
  • 下游提前退出时,能通过 context 停止;
  • 发送和耗时操作都响应取消;
  • 错误有独立且不会阻塞的传播路径。

13.1 提前消费是泄漏高发点

如果调用方只取一个值就返回:

go
v := <-square(ctx, generate(ctx, values...))
return v

上游可能永远阻塞在后续发送。必须在返回前 cancel(),而每一阶段也必须监听该取消。仅由最末端监听 context 不能解开上游阻塞。

13.2 数据所有权

通过 channel 发送 []byte、map 或指针只是发送描述符/指针,不会深复制。发送后若生产者继续修改,而消费者同时读取,就会竞争。常见策略:

  • 发送后转移所有权,生产者不再访问;
  • 发送不可变值;
  • 发送前克隆;
  • 用锁保护共享对象。

channel 建立发送前写入对接收方可见,但不会自动禁止发送后的再次写入。

14. fan-out 与 fan-in

14.1 fan-out

多个 worker 从同一个输入 channel 读取,把任务分摊:

go
func worker(ctx context.Context, jobs <-chan Job) <-chan Result {
	out := make(chan Result)
	go func() {
		defer close(out)
		for job := range jobs {
			result := handle(job)
			select {
			case out <- result:
			case <-ctx.Done():
				return
			}
		}
	}()
	return out
}

创建多个 worker:

go
outputs := make([]<-chan Result, workerCount)
for i := range outputs {
	outputs[i] = worker(ctx, jobs)
}

适合任务彼此独立、顺序不重要的场景。若输出必须保持输入顺序,需要带序号并重排,或选择不同并发模型。

14.2 fan-in

把多个输出合并成一个:

go
func merge[T any](ctx context.Context, inputs ...<-chan T) <-chan T {
	out := make(chan T)
	var wg sync.WaitGroup

	forward := func(in <-chan T) {
		defer wg.Done()
		for v := range in {
			select {
			case out <- v:
			case <-ctx.Done():
				return
			}
		}
	}

	wg.Add(len(inputs))
	for _, in := range inputs {
		go forward(in)
	}

	go func() {
		wg.Wait()
		close(out)
	}()

	return out
}

必须在所有转发 goroutine 完成后关闭 out。若在启动后立即关闭,发送会 panic;若从不关闭,range 消费方无法结束。

fan-in 的输出顺序取决于调度,不应作为业务顺序。公平性也不是严格保证;如果某输入需要配额或优先级,应显式调度。

15. worker pool

固定 worker pool 对无界任务流更稳妥:

go
type Job func(context.Context) error

func runPool(ctx context.Context, workers int, jobs <-chan Job) error {
	if workers <= 0 {
		return fmt.Errorf("workers must be positive")
	}

	ctx, cancel := context.WithCancelCause(ctx)
	defer cancel(nil)

	errCh := make(chan error, 1)
	var wg sync.WaitGroup

	for range workers {
		wg.Go(func() {
			for {
				select {
				case <-ctx.Done():
					return
				case job, ok := <-jobs:
					if !ok {
						return
					}
					if err := job(ctx); err != nil {
						select {
						case errCh <- err:
							cancel(err)
						default:
						}
						return
					}
				}
			}
		})
	}

	done := make(chan struct{})
	go func() {
		wg.Wait()
		close(done)
	}()

	select {
	case err := <-errCh:
		<-done
		return err
	case <-done:
		return context.Cause(ctx)
	case <-ctx.Done():
		<-done
		return context.Cause(ctx)
	}
}

这个示例展示结构,但实际 API 还要决定:遇到首错是否丢弃队列、外部取消返回哪个错误、job 是否允许 panic、生产者在 pool 退出后如何停止发送。

worker 数量不是越大越快:

  • CPU 密集任务从接近 GOMAXPROCS 开始测;
  • I/O 密集任务可更高,但受连接池、对端容量和内存限制;
  • 下游数据库只有 20 个连接时,开 1000 个 worker 只会把等待转移到连接池。

池的入口队列必须有界,并定义满载策略。

16. 并发限制、速率与背压策略

三个概念不要混淆:

  • 并发限制:同一时刻最多多少任务在执行;
  • 速率限制:单位时间最多启动或完成多少任务;
  • 队列容量:最多允许多少任务等待。

channel 缓冲天然能实现有界队列和计数信号量,但速率限制通常需要 ticker、令牌桶或专用限流器。

队列满时必须选择策略:

  1. 阻塞生产者:可靠且形成自然背压,但可能把延迟传播到上游;
  2. 立即返回错误:适合请求/响应服务,让调用方重试或降级;
  3. 丢弃最新:适合可损失指标或刷新信号;
  4. 丢弃最旧:适合只关心最新状态,但需要明确同步实现;
  5. 合并重复任务:适合缓存刷新、配置重载;
  6. 溢出到持久队列:需要可靠性,但复杂度和延迟更高;
  7. 动态扩容:只能应对短期波动,无界扩容最终会耗尽内存。

不要用 default 静默丢消息。丢弃策略至少要有指标、日志采样或调用方可见的返回值。

背压必须端到端。中间某层用无界切片持续吸收任务,会切断 channel 的阻塞反馈,系统最终以内存崩溃的方式“限流”。

17. 错误传播与结构化并发

Go 语言没有强制结构化并发语法,但工程代码可以遵循它的思想:

  • 子任务属于一个清楚的词法/请求作用域;
  • 父任务退出会取消子任务;
  • 任一关键子任务失败可取消同组任务;
  • 父函数返回前等待自己启动的子任务结束;
  • 错误被收集或明确忽略,不留孤儿 goroutine。

17.1 首错取消的基本结构

go
func runAll(ctx context.Context, tasks []func(context.Context) error) error {
	ctx, cancel := context.WithCancelCause(ctx)
	defer cancel(nil)

	errCh := make(chan error, 1)
	var wg sync.WaitGroup

	for _, task := range tasks {
		task := task
		wg.Go(func() {
			if err := task(ctx); err != nil {
				select {
				case errCh <- err:
					cancel(err)
				default:
				}
			}
		})
	}

	done := make(chan struct{})
	go func() {
		wg.Wait()
		close(done)
	}()

	select {
	case err := <-errCh:
		<-done
		return err
	case <-done:
		return context.Cause(ctx)
	case <-ctx.Done():
		<-done
		return context.Cause(ctx)
	}
}

关键不是复制这段代码,而是看清四个角色:取消源、首错缓冲、完成等待、返回前收束。

官方维护的 golang.org/x/sync/errgroup 封装了常见的任务组、首错和 context 取消,还支持并发上限。项目允许依赖时通常比重复手写更可靠。但仍需保证任务函数真正响应 context,否则组等待也会卡住。

17.2 错误 channel 的容量

多个任务都可能失败时,无缓冲错误 channel 很容易泄漏:调用方取到首错后返回,其他发送者永远阻塞。常见做法:

  • 首错模式用容量 1 + 非阻塞发送;
  • 全量收集由专门 goroutine 持续接收,并在所有发送者结束后关闭;
  • 直接由任务组库管理。

17.3 panic 的处理

不要把 panic 当普通错误通道。若框架边界必须隔离第三方任务,可以在启动它的同一 goroutine 中 recover,转成带堆栈的失败并清理资源。恢复点应少而明确;随处 recover 会让进程继续运行在未知不变量下。

18. goroutine 泄漏

goroutine 泄漏指任务已经没有业务价值,却因阻塞或循环永远存活。它会保留栈上引用、timer、连接和其他资源,数量不断增长时最终拖垮进程。

18.1 阻塞发送

go
func first() int {
	ch := make(chan int)
	go func() {
		ch <- slowCompute()
	}()

	return fallback() // 无人接收 ch,goroutine 可能永久阻塞
}

可根据语义使用容量 1、context 取消,或确保始终接收。

18.2 阻塞接收

go
go func() {
	for v := range ch {
		process(v)
	}
}()

若发送方永不关闭,且没有取消 case,这个 goroutine 永不结束。长期服务可以这样设计,但 channel 和服务必须拥有匹配生命周期。

18.3 忘记下游提前退出

pipeline 消费者返回后,上游发送者仍在发送。每一级都必须响应共同取消,而不只是关闭最末端。

18.4 永不返回的外部调用

即使 goroutine 监听 context,它调用的函数若不接收 context、没有连接超时,也无法中断。选择支持上下文的 API,为网络客户端配置 timeout,并在取消时关闭可中断资源。

18.5 泄漏测试

可以观察测试前后的 goroutine 数量或使用泄漏检测库,但数量法容易受运行时后台任务影响。更根本的是:

  • API 测试取消后能否在期限内返回;
  • 关闭服务后 Wait 是否完成;
  • 用 goroutine profile 查阻塞栈;
  • 在压测中监控 goroutine 数是否随请求单调增长。

19. 死锁、活锁与饥饿

19.1 常见死锁

同一 goroutine 在无缓冲 channel 上发送,却没有并发接收者:

go
ch := make(chan int)
ch <- 1
fmt.Println(<-ch)

所有 goroutine 都阻塞且运行时能确定无进展时,会报告:

text
fatal error: all goroutines are asleep - deadlock!

运行时不一定能识别所有业务死锁。例如还有网络轮询 goroutine 或 cgo 调用时,进程可能只是挂住。

其他高发原因:

  • 环形 channel 依赖;
  • 等待一个永远不会关闭的 channel;
  • 持锁发送 channel,而接收方需要同一把锁;
  • 多把锁顺序不一致;
  • WaitGroup.Add / Done 数量不匹配;
  • 向 nil channel 发送或接收;
  • 所有 select case 都被 nil channel 禁用;
  • 错误路径忘记释放信号量。

19.2 活锁

goroutine 都在运行、重试、让步,却因相互礼让无法完成。例如双方检测冲突后同时退避,再同时重试。加入随机退避可能缓解,但应首先消除对称协议。

19.3 饥饿

某个 goroutine 长期得不到锁、CPU 或队列机会。Go 调度器和 select 提供一定公平性努力,但不是业务级严格公平保证。若租户、优先级或顺序有 SLA,应实现明确队列和配额。

19.4 排查顺序

  1. 获取 goroutine dump,看每条阻塞在 send、receive、select、mutex 还是 syscall;
  2. 沿资源所有权找“本应唤醒它的人”;
  3. 检查取消、关闭和错误路径;
  4. 检查锁与 channel 混用顺序;
  5. 用最小复现和超时测试固定问题;
  6. 再考虑调度器或运行时问题。

20. channel 还是锁

场景更自然的工具
保护一个 map 或小型复合状态sync.Mutex / RWMutex
高频计数、标志sync/atomic
一次初始化sync.Once
等待一组任务sync.WaitGroup
传递任务、结果或事件流channel
转移可变对象所有权channel
限制并发worker pool / 信号量
请求取消和截止时间context.Context
状态必须串行执行复杂命令单所有者 goroutine + channel,或锁

“一个 goroutine 独占 map,其他人通过 channel 请求”能避免锁,但每次访问都会调度和分配消息,读写 API 也变复杂。简单缓存用 mutex 常常更清楚、更快。

反过来,拿锁保护一条异步任务流,自己实现条件变量、关闭和排队,也可能不如 channel 直观。选择能直接表达不变量的最小工具。

不要在持锁状态下做可能长期阻塞的 channel 操作或外部 I/O,除非协议经过严格证明。它会扩大临界区并产生锁-channel 环形等待。

21. 性能与调试

21.1 不要按 goroutine 数量炫技

goroutine 初始栈小且可增长,但不是零成本。每条 goroutine 都有调度、栈、引用保留和观测成本。十万条正在等待同一个下游的 goroutine,通常说明背压失效。

21.2 常用工具

bash
go test -race ./...
go test -run TestName -count=100
go test -bench . -benchmem
go test -blockprofile=block.out ./...
go tool pprof block.out
go test -mutexprofile=mutex.out ./...
go tool pprof mutex.out

还可以使用:

  • runtime/pprof goroutine profile 查看阻塞栈;
  • net/http/pprof 在线采样服务;
  • go tool trace 查看调度、网络阻塞、同步和延迟;
  • runtime/metrics 监控 goroutine、调度和 GC 指标;
  • 日志中的 trace ID 与取消原因连接跨 goroutine 调用。

21.3 基准测试要测系统而非语法

比较 channel 和 mutex 时,应使用真实竞争度、数据大小和临界区。单核无竞争 microbenchmark 不能代表生产环境;反过来,极端争用也可能夸大差异。

关注最终指标:

  • 吞吐;
  • p50/p95/p99 延迟;
  • goroutine 数;
  • 队列深度;
  • 拒绝与超时;
  • 分配与 GC;
  • 取消到完全退出的时间。

21.4 调度不是确定的

测试不能依赖 Sleep、某条 goroutine “通常先跑”或 select 的选择顺序。使用同步事件表达前置条件,用超时只作为测试失败保险,而不是正常控制流程。

22. 与 Java 对照

GoJava 中较接近的概念关键差异
goroutinevirtual thread / thread taskGo 由运行时调度,go 不返回句柄
channelBlockingQueue + rendezvouschannel 还是语言级同步事件,可关闭、可 select
无缓冲 channelSynchronousQueueGo 直接支持发送/接收语法和 select
select多队列等待的组合Java 标准阻塞队列没有完全对应语法
context.Contextcancellation token + deadline + request scope通过显式参数沿调用链传递
WaitGroupCountDownLatch 的部分用途可复用阶段有严格时序规则,Go 1.26 有 Go
sync.Mutexsynchronized / LockGo 锁不与对象监视器绑定,不能复制
race detector并发动态分析Go 工具链直接提供 -race

Java 开发者常想保留 Future 句柄并调用 cancel(true) 中断线程。Go 更强调协作式取消:goroutine 和它调用的下游必须主动观察 context。没有观察点,就没有可靠取消。

Java 的 BlockingQueue 通常不会用“关闭队列”表示生产完成,常见 poison pill;Go channel 关闭是广播式的完成信号,不需要为每个消费者塞哨兵值。

23. 常见误区

  1. 启动 goroutine 就获得并行加速。 调度、依赖、CPU 核数和下游瓶颈决定收益。
  2. main 会等待其他 goroutine。 main 返回即进程结束。
  3. Sleep 等完成。 应使用 channel、WaitGroup 或任务组。
  4. goroutine 结束会同步内存。 必须有 channel、锁等 happens-before 边。
  5. channel 发送后可以继续改切片。 发送不复制底层数组,可能竞争。
  6. channel 必须关闭。 只有接收方需要“不会再有值”的信号时才关闭。
  7. 接收方应该关闭 channel。 通常由确定不再发送的一方关闭。
  8. 多个发送者各自 defer close。 会产生重复关闭或发送到已关闭 channel。
  9. len(ch) 判断能否发送。 并发下观察立即过期,用 select。
  10. 缓冲越大吞吐越高。 大缓冲可能只是延长排队并隐藏过载。
  11. default 能避免所有阻塞。 它也可能造成忙循环和静默丢数据。
  12. 关闭 channel 后不能接收。 会先排空,之后持续返回零值、false。
  13. nil channel 等同关闭 channel。 nil 永久阻塞,关闭 channel 立即可接收。
  14. Context 会强行终止 goroutine。 取消是协作式信号。
  15. 只在循环顶部检查取消就够了。 每个潜在长阻塞点都应可取消。
  16. WaitGroup 能传播错误。 它只计数和等待。
  17. 每个任务一个 goroutine 再用 semaphore 就是池。 大量任务仍会创建大量等待 goroutine。
  18. select 按 case 顺序优先。 多个就绪 case 采用伪随机选择。
  19. 接口或 channel 赋值天然原子。 并发读写变量仍需同步。
  20. race detector 没报错就证明正确。 它只覆盖实际执行路径。

24. 速查表

goroutine 与同步

需求写法
启动任务go f()
等待一组任务sync.WaitGroup / wg.Go(f)
传单个结果容量 1 的结果 channel
首错取消context.WithCancelCause + 任务组
检测数据竞争go test -race ./...

channel

操作写法/行为
无缓冲创建make(chan T)
有缓冲创建make(chan T, n)
发送ch <- v
接收v := <-ch
检测关闭v, ok := <-ch
遍历到关闭for v := range ch
关闭close(ch)
只接收<-chan T
只发送chan<- T
非阻塞尝试select + default
动态禁用 select case把 channel 设为 nil

特殊状态

状态发送接收关闭
nil channel永久阻塞永久阻塞panic
正常无缓冲等待接收者等待发送者成功
正常缓冲未满/非空可进展可进展成功
已关闭panic排空后零值、falsepanic

并发代码审查清单

  • 每条 goroutine 的退出条件在哪里?
  • 父作用域返回前是否等待它?
  • 所有阻塞发送、接收、I/O 是否响应取消?
  • channel 由谁创建、发送、关闭?
  • 多发送者关闭是否由协调者完成?
  • 队列是否有界,满载策略是什么?
  • 发送的数据在交付后谁拥有?
  • 错误是否可能因无人接收而阻塞?
  • 是否在持锁时进行 channel 或外部 I/O?
  • 测试是否运行 -race 并覆盖取消、超时和关闭?

可运行示例

并发示例不能只验证“能跑”,还要确认它能结束、不会泄漏,并且输出可验证。这里的程序不使用睡眠猜测时序,而是通过关闭、取消和 WaitGroup 建立明确的生命周期。

示例一:由发送方拥有 channel 的关闭权

先确定关闭权。 channel 通常应由能够确认“不会再有发送”的一方关闭。这个生产者创建输出 channel、发送完整序列并负责关闭,消费者只负责接收。

go
package main

import "fmt"

func produce(limit int) <-chan int {
	output := make(chan int)
	go func() {
		// 创建并发送 channel 的一方拥有关闭权。
		// defer 确保生产者无论从哪个正常路径退出,消费者都能结束 range。
		defer close(output)
		for number := 1; number <= limit; number++ {
			output <- number
		}
	}()
	return output
}

func main() {
	sum := 0
	for number := range produce(5) {
		sum += number
	}
	fmt.Println("sum:", sum)
}

运行:

bash
go run ./examples/ch12/channel-ownership

预期输出:

text
sum: 15

所有权体现在三处。

  1. 返回 <-chan int 限制调用者只能接收,API 直接表达所有权。
  2. 生产 goroutine 用 defer close(output) 保证正常退出时发出结束信号。
  3. 消费者 range 到关闭后自然结束,不需要额外布尔标记。

进一步验证。 把限制从 5 改成 10 并预测总和;再故意删除 close(output),程序会在收完数据后死锁。这正说明关闭不是附加动作,而是协议的一部分。

示例二:可取消且无泄漏的 pipeline

下游提前离开会发生什么。 下游只取前三个结果便离开时,无限生产者可能永远阻塞在发送上。这个两阶段 pipeline 让每一次发送和接收都能响应 context 取消。

go
package main

import (
	"context"
	"fmt"
	"sync"
)

func generate(ctx context.Context, wg *sync.WaitGroup) <-chan int {
	output := make(chan int)
	wg.Add(1)
	go func() {
		defer wg.Done()
		defer close(output)
		for number := 1; ; number++ {
			select {
			case output <- number:
			case <-ctx.Done():
				return // 下游提前退出时停止发送,避免 goroutine 永久阻塞。
			}
		}
	}()
	return output
}

func square(ctx context.Context, wg *sync.WaitGroup, input <-chan int) <-chan int {
	output := make(chan int)
	wg.Add(1)
	go func() {
		defer wg.Done()
		defer close(output)
		for {
			select {
			case number, ok := <-input:
				if !ok {
					return
				}
				select {
				case output <- number * number:
				case <-ctx.Done():
					return
				}
			case <-ctx.Done():
				return
			}
		}
	}()
	return output
}

func main() {
	ctx, cancel := context.WithCancel(context.Background())
	var wg sync.WaitGroup
	values := square(ctx, &wg, generate(ctx, &wg))

	for range 3 {
		fmt.Println(<-values)
	}
	cancel() // 消费够了就广播取消,而不是直接遗弃上游。
	wg.Wait()
	fmt.Println("pipeline stopped")
}

运行:

bash
go run ./examples/ch12/pipeline-cancel

预期输出:

text
1
4
9
pipeline stopped

退出路径由这些动作组成。

  • generatesquare 都只关闭自己创建的输出 channel。
  • channel 操作与 ctx.Done() 放在同一个 select 中,取消后不会卡在无人接收的发送。
  • cancel() 广播退出,wg.Wait() 则证明两个 goroutine 确实已经归还。

破坏一次退出协议。 把消费数量改成 1 或 10;再临时删掉 generate 发送处的取消分支,运行后观察 wg.Wait() 为什么无法返回。不要用 time.Sleep 掩盖生命周期问题。

示例三:有界 worker pool 与确定性汇总

并发量必须有上限。 为每个任务创建一个 goroutine,会让并发量随输入一起增长。worker pool 用固定 worker 数限制并行度,并在全部 worker 退出后安全关闭结果通道。

go
package main

import (
	"fmt"
	"sort"
	"sync"
)

type Job struct {
	ID    int
	Value int
}

type Result struct {
	ID      int
	Squared int
}

func worker(jobs <-chan Job, results chan<- Result, wg *sync.WaitGroup) {
	defer wg.Done()
	for job := range jobs {
		results <- Result{ID: job.ID, Squared: job.Value * job.Value}
	}
}

func main() {
	jobs := make(chan Job)
	results := make(chan Result)
	var workers sync.WaitGroup

	for range 3 {
		workers.Add(1)
		go worker(jobs, results, &workers)
	}

	go func() {
		// 调度方拥有 jobs 的发送与关闭权。
		defer close(jobs)
		for id, value := range []int{2, 3, 4, 5} {
			jobs <- Job{ID: id, Value: value}
		}
	}()

	go func() {
		// 只有确认全部 worker 停止后才能关闭 results,避免 send on closed channel。
		workers.Wait()
		close(results)
	}()

	var collected []Result
	for result := range results {
		collected = append(collected, result)
	}

	// worker 完成顺序不确定;排序后再展示,测试与文档输出才稳定。
	sort.Slice(collected, func(i, j int) bool {
		return collected[i].ID < collected[j].ID
	})
	for _, result := range collected {
		fmt.Printf("job %d -> %d\n", result.ID, result.Squared)
	}
}

运行:

bash
go run ./examples/ch12/worker-pool

预期输出:

text
job 0 -> 4
job 1 -> 9
job 2 -> 16
job 3 -> 25

这段代码有三个所有者。

  1. 调度 goroutine 独占 jobs 的发送和关闭权。
  2. 多个 worker 共享只读输入和只写输出,数量决定最大并行任务数。
  3. 单独的收尾 goroutine 等待全部 worker 后关闭 results
  4. 并发完成顺序不保证,展示前按任务 ID 排序,使输出与测试稳定。

改变并行度再观察。 把 worker 数分别改成 1 和 8,结果仍应相同;给 Result 增加 WorkerID,观察调度分布,但不要断言某个任务一定由某个 worker 执行。

25. 练习

  1. 写一个并发程序,让 goroutine 修改变量但不做同步,用 -race 观察报告;分别用 channel 和 mutex 修复。
  2. 实现一个无缓冲请求/响应 channel,证明发送完成只代表交接完成,不代表接收方业务处理完成。
  3. 写一个容量可配置的队列,分别实现阻塞、拒绝和丢弃最新三种过载策略,并记录指标。
  4. 实现三个阶段的 pipeline;让消费者只取第一个结果就取消,证明所有 goroutine 能退出。
  5. 实现泛型 merge[T],合并多个 channel,并测试零输入、一个输入、提前取消和所有输入关闭。
  6. 实现固定大小 worker pool,支持 context、首错取消、有界任务队列和优雅关闭。
  7. 写一个多发送者 channel,先让每个发送者都 close 复现 panic,再改成 WaitGroup 协调关闭。
  8. 在 select 循环中处理两个输入 channel;一个关闭后设为 nil,避免零值忙循环。
  9. 为 HTTP 服务添加后台刷新 goroutine,设计 Close(ctx),保证停止 ticker、取消请求并等待完成。
  10. 用 block profile 找出一个 channel 死锁,用 goroutine dump 解释每条 goroutine 在等谁。
  11. 比较“每任务一 goroutine + 信号量”和固定 worker pool 在十万任务下的 goroutine 峰值与内存。
  12. 设计一个必须保序的并发处理器:任务携带序号,worker 并行处理,聚合器按序输出,并处理某任务失败。

26. 官方资料

以 Go 官方规范与标准库文档为准,示例面向 Go 1.26。