Go 并发:goroutine、channel 与可收束的并发结构
面向有 Java 经验的开发者,基于 Go 1.26。
go f() 和 make(chan T) 几分钟就能学会,但能启动并发任务,离可靠地管理它们还有很远。谁拥有 goroutine?谁负责停止它?错误怎样传回?发送方被取消时会不会永久阻塞?缓冲区满了以后,系统选择阻塞、丢弃还是扩容?函数返回之前,自己启动的任务是否已经收束?这些问题必须在代码中找到答案。
Go 并发的重点不是“轻量线程很多”,而是把通信、同步和生命周期协议写清楚。channel 既能传值,也能建立 happens-before 关系,但它不会替你解决所有权、取消、死锁、背压和数据竞争。后面的内容会先说明语言语义,再把这些语义落到可收束的并发结构中。
目录
1. 先建立并发心智模型
一个 goroutine 是由 Go 运行时调度的独立执行单元。它不是一个带 join、interrupt 方法的对象;go 语句也不返回句柄。启动之后,调用方如果需要等待、取消或接收结果,必须另外设计协议。
go serve(conn)一个 channel 是带类型的通信通道:
ch := make(chan Result)
ch <- result // 发送
result := <-ch // 接收阅读后面的代码时,先抓住三条原则:
- 启动 goroutine 时,同时决定退出路径。
- 共享内存需要同步;channel 是同步手段之一,不是数据竞争豁免证。
- 把容量和过载行为当作 API 协议,而不是随手写一个缓冲数字。
“不要通过共享内存来通信;通过通信来共享内存”是一条设计提示,不是禁止使用 sync.Mutex 的教条。计数器、缓存和短临界区通常用锁更直接;流式传递、所有权转移和任务协调则更适合 channel。
2. goroutine 的创建与生命周期
2.1 go 语句
函数调用前加 go:
go fetch(url)
go func() {
fmt.Println("background")
}()调用表达式的函数值和参数在启动方 goroutine 中求值,然后新 goroutine 执行函数体:
for _, url := range urls {
go fetch(url)
}Go 1.22 起,for 循环声明的迭代变量按每次迭代拥有新变量的语义处理,Go 1.26 代码不再有过去那个经典闭包捕获问题。显式传参仍能让数据边界更清楚:
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 不会被等待:
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 等待单个任务
resultCh := make(chan Result, 1)
go func() {
resultCh <- compute()
}()
result := <-resultCh缓冲为 1 的结果 channel 有一个重要性质:即使调用方因取消提前离开,工作 goroutine 完成发送时也不必等待接收者。若结果包含大对象且无人接收,应考虑生命周期和及时释放。
3.2 sync.WaitGroup
多个无返回结果的任务可以使用 WaitGroup:
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 中可以写:
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调用返回。
例子:
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 结束不是同步
var message string
go func() {
message = "ready"
}()
fmt.Println(message) // 数据竞争;不能靠“它应该跑完了”推理goroutine 消失本身不保证其写入对其他 goroutine 可见。必须等待一个同步事件。
4.3 原子不变量
即使单个机器字读写在某平台上看似原子,也不代表无竞争,更不代表多个字段的一致性。需要原子计数可使用 sync/atomic;需要维护复合不变量通常使用 mutex;需要转移数据所有权可使用 channel。
所有并发测试都应考虑 race detector:
go test -race ./...它只能发现实际运行路径触发的竞争,不能证明程序在所有路径上无竞争,因此测试覆盖和设计审查仍然重要。
5. channel 类型与创建
channel 类型写成 chan T:
var jobs chan Job
jobs = make(chan Job)
buffered := make(chan Job, 32)chan T 是双向 channel,<-chan T 只接收,chan<- T 只发送。channel 值可以复制,副本指向同一个运行时通道。
5.1 零值是 nil
var ch chan int
fmt.Println(ch == nil) // true对 nil channel 发送或接收会永久阻塞;关闭 nil channel会 panic:
// ch <- 1 // 永久阻塞
// <-ch // 永久阻塞
// close(ch) // panicnil channel 在 select 中很有用,因为对应 case 永远不会就绪,可以动态关闭某个分支;在普通顺序代码里,它通常表示忘记初始化。
5.2 channel 可比较
channel 可以和 nil 比较,也可以比较是否指向同一通道:
a := make(chan int)
b := a
fmt.Println(a == b) // true这个能力适合身份判断,不应用来推导队列内容或状态。
6. 无缓冲 channel
ch := make(chan string)无缓冲 channel 的发送必须和接收会合。发送方与接收方都到达通信点时,值才交付,双方才继续:
go func() {
ch <- "done"
}()
msg := <-ch它同时完成:
- 值传递;
- 发送与接收之间的同步;
- 对双方执行进度的协调。
不能简单说“发送一定阻塞到接收方把后续逻辑执行完”。它只等待对应通信完成;接收方拿到值后做什么,与发送方可以再次并发。
6.1 用空结构体表达纯信号
done := make(chan struct{})
go func() {
work()
close(done)
}()
<-donestruct{} 不携带业务数据,表达“这里只需要事件”。一次性完成通知通常用关闭比发送更灵活,因为关闭能广播给所有等待者。
6.2 无缓冲不等于更正确
无缓冲使交接点明确,但也让两个组件的执行节奏紧耦合。若消费者暂时退出,生产者可能泄漏;若环形依赖中所有人先发送,就会死锁。选择无缓冲应源于协议需要同步交接,而不是默认认为“缓冲会隐藏问题”。
7. 缓冲 channel 与背压
jobs := make(chan Job, 100)发送在缓冲未满时可立即完成,接收在缓冲非空时可立即完成。满时发送阻塞,空时接收阻塞。
fmt.Println(len(jobs), cap(jobs))len 只是一瞬间的排队数量,不能用于“先判断再发送”的并发正确性:
// if len(ch) < cap(ch) { ch <- v } // 条件可能立刻过期非阻塞尝试应使用 select:
select {
case jobs <- job:
// accepted
default:
// overloaded
}7.1 缓冲容量是系统策略
缓冲并不创造吞吐,只吸收生产和消费之间的短期波动。容量过大可能:
- 延迟过载暴露;
- 增加内存占用;
- 延长排队时间,任务到手时已经过期;
- 让关闭过程等待大量陈旧任务;
- 掩盖消费者持续低于生产者的事实。
容量过小则可能让生产者频繁阻塞,降低并行度。应根据到达率、处理速率、允许等待时间和内存预算决定,并通过队列长度、拒绝数和端到端延迟监控验证。
Little's Law 可作容量估算起点:稳定系统中平均在途数量约等于到达率乘平均停留时间。但真实流量有突发,仍需压测。
7.2 channel 作为信号量
容量 channel 可以限制并发:
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”,也不是给接收者发送一个特殊值。
close(ch)接收会先取完缓冲区已有值,之后立即得到元素零值和 ok == false:
v, ok := <-ch
if !ok {
// 已关闭且已无值
}range 会持续接收,直到 channel 关闭并排空:
for v := range ch {
consume(v)
}8.1 关闭规则
- 向已关闭 channel 发送:panic;
- 再次关闭已关闭 channel:panic;
- 关闭 nil channel:panic;
- 从已关闭且排空的 channel 接收:立即返回零值、false;
- 从仅接收 channel 不能调用
close; - channel 不必总是关闭,只有接收方需要结束信号时才需要。
8.2 由发送方关闭
一般由确定“不会再有发送”的一方关闭。接收方擅自关闭会和仍在运行的发送方竞争,导致 panic。
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
}
}
}多个发送者时,任何单个发送者都不知道其他人是否结束。应让一个协调者等待全部发送者,再关闭:
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 与所有权
函数签名应表达方向:
func generate(ctx context.Context) <-chan int
func consume(ctx context.Context, in <-chan int) error
func publish(out chan<- Event, e Event)双向 channel 可隐式收窄为单向,不能反向拓宽:
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 操作中等待:
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 非阻塞发送和接收
select {
case v := <-ch:
handle(v)
default:
// 当前没有值
}select {
case ch <- v:
// 发送成功
default:
// 当前不可发送
}“当前不可用”不是“永远不可用”。这种模式适合指标、尽力通知或明确允许丢弃的负载;不能用来悄悄丢失订单、账务事件等可靠消息。
10.2 nil channel 动态禁用 case
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 接收立即返回。循环中若不处理,它会持续被选中并形成忙循环:
case v, ok := <-ch:
if !ok {
ch = nil // 禁用这个 case
continue
}11. 超时、定时器与心跳
一次性超时可以写:
select {
case result := <-resultCh:
return result, nil
case <-time.After(500 * time.Millisecond):
return Result{}, context.DeadlineExceeded
}在高频循环里通常创建并复用 time.Timer,便于停止、重置和控制资源:
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:
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 携带取消、截止时间和请求范围值。惯例是作为第一个参数显式传入:
func Load(ctx context.Context, id UserID) (User, error)不要把 Context 存进结构体,不要传 nil;没有合适上下文时使用 context.Background() 或 context.TODO()。
12.1 派生和取消
ctx, cancel := context.WithTimeout(parent, 2*time.Second)
defer cancel()
user, err := repo.Load(ctx, id)即使操作提前完成也要调用 cancel,以释放关联资源。子上下文取消不会取消父上下文;父取消会传播到所有后代。
goroutine 应在可能长时间阻塞的位置监听:
select {
case out <- value:
case <-ctx.Done():
return ctx.Err()
}只在循环顶部偶尔检查一次不够。如果随后可能永久阻塞在 channel 发送,取消仍无法生效。
12.2 Done、Err 和取消原因
ctx.Done() 返回一个接收 channel,取消时关闭。调用者不能关闭它。ctx.Err() 返回 context.Canceled 或 context.DeadlineExceeded。
需要保留业务取消原因时:
ctx, cancel := context.WithCancelCause(parent)
cancel(fmt.Errorf("dependency failed: %w", err))
fmt.Println(context.Cause(ctx))派生超时也有 WithTimeoutCause、WithDeadlineCause。用 Cause 能保留“为什么整组任务停止”,但 API 仍应规定最终返回哪个错误。
12.3 Context value 的边界
Context value 适合请求范围、跨 API 边界传递的数据,如 trace ID 或认证主体。不要用它传可选函数参数、数据库连接或业务必填字段。键应使用包内定义的不可导出类型,避免冲突:
type traceIDKey struct{}
ctx = context.WithValue(ctx, traceIDKey{}, traceID)13. pipeline
pipeline 把数据分阶段处理,每个阶段接收上游 channel 并返回下游 channel:
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
}使用:
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 提前消费是泄漏高发点
如果调用方只取一个值就返回:
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 读取,把任务分摊:
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:
outputs := make([]<-chan Result, workerCount)
for i := range outputs {
outputs[i] = worker(ctx, jobs)
}适合任务彼此独立、顺序不重要的场景。若输出必须保持输入顺序,需要带序号并重排,或选择不同并发模型。
14.2 fan-in
把多个输出合并成一个:
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 对无界任务流更稳妥:
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、令牌桶或专用限流器。
队列满时必须选择策略:
- 阻塞生产者:可靠且形成自然背压,但可能把延迟传播到上游;
- 立即返回错误:适合请求/响应服务,让调用方重试或降级;
- 丢弃最新:适合可损失指标或刷新信号;
- 丢弃最旧:适合只关心最新状态,但需要明确同步实现;
- 合并重复任务:适合缓存刷新、配置重载;
- 溢出到持久队列:需要可靠性,但复杂度和延迟更高;
- 动态扩容:只能应对短期波动,无界扩容最终会耗尽内存。
不要用 default 静默丢消息。丢弃策略至少要有指标、日志采样或调用方可见的返回值。
背压必须端到端。中间某层用无界切片持续吸收任务,会切断 channel 的阻塞反馈,系统最终以内存崩溃的方式“限流”。
17. 错误传播与结构化并发
Go 语言没有强制结构化并发语法,但工程代码可以遵循它的思想:
- 子任务属于一个清楚的词法/请求作用域;
- 父任务退出会取消子任务;
- 任一关键子任务失败可取消同组任务;
- 父函数返回前等待自己启动的子任务结束;
- 错误被收集或明确忽略,不留孤儿 goroutine。
17.1 首错取消的基本结构
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 阻塞发送
func first() int {
ch := make(chan int)
go func() {
ch <- slowCompute()
}()
return fallback() // 无人接收 ch,goroutine 可能永久阻塞
}可根据语义使用容量 1、context 取消,或确保始终接收。
18.2 阻塞接收
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 上发送,却没有并发接收者:
ch := make(chan int)
ch <- 1
fmt.Println(<-ch)所有 goroutine 都阻塞且运行时能确定无进展时,会报告:
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 排查顺序
- 获取 goroutine dump,看每条阻塞在 send、receive、select、mutex 还是 syscall;
- 沿资源所有权找“本应唤醒它的人”;
- 检查取消、关闭和错误路径;
- 检查锁与 channel 混用顺序;
- 用最小复现和超时测试固定问题;
- 再考虑调度器或运行时问题。
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 常用工具
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/pprofgoroutine 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 对照
| Go | Java 中较接近的概念 | 关键差异 |
|---|---|---|
| goroutine | virtual thread / thread task | Go 由运行时调度,go 不返回句柄 |
| channel | BlockingQueue + rendezvous | channel 还是语言级同步事件,可关闭、可 select |
| 无缓冲 channel | SynchronousQueue | Go 直接支持发送/接收语法和 select |
select | 多队列等待的组合 | Java 标准阻塞队列没有完全对应语法 |
context.Context | cancellation token + deadline + request scope | 通过显式参数沿调用链传递 |
WaitGroup | CountDownLatch 的部分用途 | 可复用阶段有严格时序规则,Go 1.26 有 Go |
sync.Mutex | synchronized / Lock | Go 锁不与对象监视器绑定,不能复制 |
| race detector | 并发动态分析 | Go 工具链直接提供 -race |
Java 开发者常想保留 Future 句柄并调用 cancel(true) 中断线程。Go 更强调协作式取消:goroutine 和它调用的下游必须主动观察 context。没有观察点,就没有可靠取消。
Java 的 BlockingQueue 通常不会用“关闭队列”表示生产完成,常见 poison pill;Go channel 关闭是广播式的完成信号,不需要为每个消费者塞哨兵值。
23. 常见误区
- 启动 goroutine 就获得并行加速。 调度、依赖、CPU 核数和下游瓶颈决定收益。
main会等待其他 goroutine。main返回即进程结束。- 用
Sleep等完成。 应使用 channel、WaitGroup 或任务组。 - goroutine 结束会同步内存。 必须有 channel、锁等 happens-before 边。
- channel 发送后可以继续改切片。 发送不复制底层数组,可能竞争。
- channel 必须关闭。 只有接收方需要“不会再有值”的信号时才关闭。
- 接收方应该关闭 channel。 通常由确定不再发送的一方关闭。
- 多个发送者各自 defer close。 会产生重复关闭或发送到已关闭 channel。
- 用
len(ch)判断能否发送。 并发下观察立即过期,用 select。 - 缓冲越大吞吐越高。 大缓冲可能只是延长排队并隐藏过载。
default能避免所有阻塞。 它也可能造成忙循环和静默丢数据。- 关闭 channel 后不能接收。 会先排空,之后持续返回零值、false。
- nil channel 等同关闭 channel。 nil 永久阻塞,关闭 channel 立即可接收。
- Context 会强行终止 goroutine。 取消是协作式信号。
- 只在循环顶部检查取消就够了。 每个潜在长阻塞点都应可取消。
- WaitGroup 能传播错误。 它只计数和等待。
- 每个任务一个 goroutine 再用 semaphore 就是池。 大量任务仍会创建大量等待 goroutine。
- select 按 case 顺序优先。 多个就绪 case 采用伪随机选择。
- 接口或 channel 赋值天然原子。 并发读写变量仍需同步。
- 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 | 排空后零值、false | panic |
并发代码审查清单
- 每条 goroutine 的退出条件在哪里?
- 父作用域返回前是否等待它?
- 所有阻塞发送、接收、I/O 是否响应取消?
- channel 由谁创建、发送、关闭?
- 多发送者关闭是否由协调者完成?
- 队列是否有界,满载策略是什么?
- 发送的数据在交付后谁拥有?
- 错误是否可能因无人接收而阻塞?
- 是否在持锁时进行 channel 或外部 I/O?
- 测试是否运行
-race并覆盖取消、超时和关闭?
可运行示例
并发示例不能只验证“能跑”,还要确认它能结束、不会泄漏,并且输出可验证。这里的程序不使用睡眠猜测时序,而是通过关闭、取消和 WaitGroup 建立明确的生命周期。
示例一:由发送方拥有 channel 的关闭权
先确定关闭权。 channel 通常应由能够确认“不会再有发送”的一方关闭。这个生产者创建输出 channel、发送完整序列并负责关闭,消费者只负责接收。
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)
}运行:
go run ./examples/ch12/channel-ownership预期输出:
sum: 15所有权体现在三处。
- 返回
<-chan int限制调用者只能接收,API 直接表达所有权。 - 生产 goroutine 用
defer close(output)保证正常退出时发出结束信号。 - 消费者
range到关闭后自然结束,不需要额外布尔标记。
进一步验证。 把限制从 5 改成 10 并预测总和;再故意删除 close(output),程序会在收完数据后死锁。这正说明关闭不是附加动作,而是协议的一部分。
示例二:可取消且无泄漏的 pipeline
下游提前离开会发生什么。 下游只取前三个结果便离开时,无限生产者可能永远阻塞在发送上。这个两阶段 pipeline 让每一次发送和接收都能响应 context 取消。
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")
}运行:
go run ./examples/ch12/pipeline-cancel预期输出:
1
4
9
pipeline stopped退出路径由这些动作组成。
generate和square都只关闭自己创建的输出 channel。- channel 操作与
ctx.Done()放在同一个select中,取消后不会卡在无人接收的发送。 cancel()广播退出,wg.Wait()则证明两个 goroutine 确实已经归还。
破坏一次退出协议。 把消费数量改成 1 或 10;再临时删掉 generate 发送处的取消分支,运行后观察 wg.Wait() 为什么无法返回。不要用 time.Sleep 掩盖生命周期问题。
示例三:有界 worker pool 与确定性汇总
并发量必须有上限。 为每个任务创建一个 goroutine,会让并发量随输入一起增长。worker pool 用固定 worker 数限制并行度,并在全部 worker 退出后安全关闭结果通道。
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)
}
}运行:
go run ./examples/ch12/worker-pool预期输出:
job 0 -> 4
job 1 -> 9
job 2 -> 16
job 3 -> 25这段代码有三个所有者。
- 调度 goroutine 独占
jobs的发送和关闭权。 - 多个 worker 共享只读输入和只写输出,数量决定最大并行任务数。
- 单独的收尾 goroutine 等待全部 worker 后关闭
results。 - 并发完成顺序不保证,展示前按任务 ID 排序,使输出与测试稳定。
改变并行度再观察。 把 worker 数分别改成 1 和 8,结果仍应相同;给 Result 增加 WorkerID,观察调度分布,但不要断言某个任务一定由某个 worker 执行。
25. 练习
- 写一个并发程序,让 goroutine 修改变量但不做同步,用
-race观察报告;分别用 channel 和 mutex 修复。 - 实现一个无缓冲请求/响应 channel,证明发送完成只代表交接完成,不代表接收方业务处理完成。
- 写一个容量可配置的队列,分别实现阻塞、拒绝和丢弃最新三种过载策略,并记录指标。
- 实现三个阶段的 pipeline;让消费者只取第一个结果就取消,证明所有 goroutine 能退出。
- 实现泛型
merge[T],合并多个 channel,并测试零输入、一个输入、提前取消和所有输入关闭。 - 实现固定大小 worker pool,支持 context、首错取消、有界任务队列和优雅关闭。
- 写一个多发送者 channel,先让每个发送者都 close 复现 panic,再改成 WaitGroup 协调关闭。
- 在 select 循环中处理两个输入 channel;一个关闭后设为 nil,避免零值忙循环。
- 为 HTTP 服务添加后台刷新 goroutine,设计
Close(ctx),保证停止 ticker、取消请求并等待完成。 - 用 block profile 找出一个 channel 死锁,用 goroutine dump 解释每条 goroutine 在等谁。
- 比较“每任务一 goroutine + 信号量”和固定 worker pool 在十万任务下的 goroutine 峰值与内存。
- 设计一个必须保序的并发处理器:任务携带序号,worker 并行处理,聚合器按序输出,并处理某任务失败。
26. 官方资料
- Go 语言规范:Go statements
- Go 语言规范:Channel types
- Go 语言规范:Send statements
- Go 语言规范:Receive operator
- Go 语言规范:Select statements
- Go 内存模型
- Go 官方博客:Go Concurrency Patterns: Pipelines and cancellation
- Go 官方博客:Share Memory By Communicating
- Go 官方博客:Context
context包sync包sync/atomic包runtime/pprof包- Data Race Detector
- Diagnostics