从一段「漏掉最后一个结果」的 WorkerPool 代码说起,把 WaitGroup 讲透。
一、WaitGroup 是什么
sync.WaitGroup 是 Go 标准库提供的并发同步原语,用来等待一组 goroutine 全部执行完毕。本质就是一个并发安全的计数器:
Add(n int):计数器加 n(可为负)Done():计数器减 1,等价于Add(-1)Wait():阻塞,直到计数器归零
你可以把它理解成「门卫」:Add 是发通行证,Done 是回收通行证,Wait 是「等所有通行证都收回来再放行主流程」。
二、标准使用模板
var wg sync.WaitGroup
for i := 0; i < n; i++ {
wg.Add(1) // ① 启动 goroutine 之前 Add
go func(i int) {
defer wg.Done() // ② goroutine 内 defer Done
// 干活
}(i)
}
wg.Wait() // ③ 外部等待
记住三步:启动前 Add、协程内 Done、外部 Wait。顺序不能乱。
三、五条铁律(违反即 panic 或死锁)
铁律 1:Add 必须在 go 之前调用
// ❌ 错误:在 goroutine 内 Add
go func() {
wg.Add(1) // 主协程的 Wait 可能在此时已经看到 0 并返回
defer wg.Done()
}()
Wait 只判断计数器是否归零。如果在 goroutine 里才 Add,主协程的 Wait 完全可能在 Add 执行之前就看到 0 并解除阻塞——goroutine 还没跑完主流程就走了。这是最常见的并发 bug,而且只在调度时序不巧时才触发,极难复现。
铁律 2:Add 不能和 Wait 并发
// ❌ 错误:一边 Wait 一边 Add
go func() {
for {
wg.Add(1) // Wait 进行中再 Add → panic
...
}
}()
wg.Wait()
一旦开始 Wait,就只允许计数递减(Done),不允许再 Add。违者触发:
sync: WaitGroup is reused before previous Wait has returned
这条铁律是初学者最容易踩的雷。下面会专门展开。
铁律 3:计数不能变负
Done 次数绝不能超过 Add 次数,否则:
sync: negative WaitGroup counter
铁律 4:批量启动用 Delta
// ✅ 推荐
wg.Add(n)
for i := 0; i < n; i++ {
go work()
}
// 也可以,但每次 Add 都要抢锁
for i := 0; i < n; i++ {
wg.Add(1)
go work()
}
Add(n) 一次性加 n,省去循环里逐次抢锁的开销。注意:Add(n) 同样必须在 go 之前调用,n 要预先计算好。
铁律 5:不要拷贝 WaitGroup
WaitGroup 含 noCopy 字段,按值传递会拷贝出一份独立计数器,同步立刻失效。要么用指针,要么作为结构体字段按引用持有:
// ❌ 错误
func bad(wg sync.WaitGroup) { ... }
// ✅ 正确
func good(wg *sync.WaitGroup) { ... }
四、真实案例:为什么这段 WorkerPool 漏掉了「结果4」
原始问题代码
func (p *WorkerPool) CollectResults() (results []string) {
go func() { // ① 启动收集 goroutine
for val := range p.resultChan {
results = append(results, val)
}
}()
p.wg.Wait() // ② 等 worker 全部 Done
close(p.resultChan) // ③ 关通道
return // ④ 立即返回 results
}
两次踩雷的演进
第一次尝试(错误改法)
有人想「让 wg 计数对齐任务数」就能解决,改成这样:
func (p *WorkerPool) worker(id int) {
for val := range p.jobChan {
p.wg.Add(1) // ❌ 在 goroutine 内 Add
// ...
p.resultChan <- val
}
}
func (p *WorkerPool) CollectResults() (results []string) {
go func() {
for val := range p.resultChan {
results = append(results, val)
p.wg.Done() // ❌ 与 Wait 并发的 Add
}
}()
p.wg.Wait()
close(p.resultChan)
return
}
这个改法会直接 panic:
Add放在 worker 的 for 循环里,与CollectResults中的Wait并发执行,触发铁律 2 的 panic。- 退一步说,即便不 panic,「wg 计数 == 任务数」也不是正确性的来源——原 bug 是「worker 完成」和「收集协程完成」是两件事,用一个 wg 混淆了语义。
根因剖析
原 wg 计数的是「worker 是否全部产出完毕」(生产侧),但 CollectResults 真正需要等的是「收集协程把缓冲区读完并退出」(消费侧)。生产 ≠ 消费,两个独立的 happens-before 链,用一个 wg 串起来就会漏:
wg.Wait()返回只表示所有 worker 调了Done(),但 worker 在resultChan <- val之后立刻Done(),值可能还躺在带缓冲的resultChan(容量 3)里没被消费。close(resultChan)之后for range会读完缓冲区剩余值才退出,但return可能在收集协程跑完前就执行了——最后一个 append 丢失,于是「结果4」没了。这是因为 4 是最后写入缓冲区的那一个,主协程 return 时它还没被取出。
五、正确写法:两个 WaitGroup 各管一件事
type WorkerPool struct {
workerCount int
jobChan chan string
resultChan chan string
workerWg sync.WaitGroup // 计 worker 是否全部产出完毕
collectWg sync.WaitGroup // 计收集协程是否读完缓冲区
}
func (p *WorkerPool) start() {
for i := 0; i < p.workerCount; i++ {
p.workerWg.Add(1) // ① 启动前 Add
go p.worker(i)
}
}
func (p *WorkerPool) worker(id int) {
defer p.workerWg.Done() // ② goroutine 内 Done
for val := range p.jobChan {
t := time.Duration(rand.Intn(6)) * time.Second
time.Sleep(t)
p.resultChan <- val
}
}
func (p *WorkerPool) CollectResults() (results []string) {
p.collectWg.Add(1) // ① 启动前 Add
go func() {
defer p.collectWg.Done() // ② goroutine 内 Done
for val := range p.resultChan {
results = append(results, val)
}
}()
p.workerWg.Wait() // 等所有 worker 产出完毕 → 安全 close
close(p.resultChan)
p.collectWg.Wait() // 等收集协程读完 → 安全 return
return
}
两个 wg 各司其职:
| WaitGroup | 计的是什么 | 决定什么 |
|---|---|---|
workerWg | worker 是否全部产出完毕 | 何时 close(resultChan) |
collectWg | 收集协程是否读完缓冲区退出 | 何时安全 return |
记住这句话:「关 diver 等消费者完成」才是关键,不是「wg 计数对齐任务数」。 生产侧的 Wait 决定何时关 chan,消费侧的 Wait 决定何时关完之后还能安全 return。
六、只用一个 WaitGroup 怎么写
如果你坚持只用一个 wg,必须引入一个额外的无缓冲通道顶替第二个 wg 的同步角色:
type WorkerPool struct {
workerCount int
jobChan chan string
resultChan chan string
wg sync.WaitGroup
done chan struct{} // 代替第二个 wg:收集协程完成信号
}
func (p *WorkerPool) start() {
for i := 0; i < p.workerCount; i++ {
p.wg.Add(1)
go p.worker(i)
}
}
func (p *WorkerPool) worker(id int) {
defer p.wg.Done()
for val := range p.jobChan {
t := time.Duration(rand.Intn(6)) * time.Second
fmt.Printf("%d 等待时间: %s\n", id, t)
time.Sleep(t)
fmt.Printf("当前Worker %d 拿到了任务: %v\n", id, val)
p.resultChan <- val
}
}
func (p *WorkerPool) CollectResults() (results []string) {
done := make(chan struct{})
p.done = done
go func() {
for val := range p.resultChan {
results = append(results, val)
}
close(done) // 收集协程读完缓冲区退出 → 发完成信号
}()
p.wg.Wait() // 等 worker 全部产出完
close(p.resultChan) // 安全关 resultChan
<-done // 等收集协程退出 → 安全 return
return
}
done 这个无缓冲通道,功能上就是第二个 WaitGroup 的等价物——它干的事和 collectWg.Wait() 一模一样:把主协程挡在 return 前,直到收集协程把缓冲区读完并发来信号。「一个 chan」和「一个 WaitGroup」互为替代是 Go 里的常见手法,选哪个看可读性。
所以「单 wg 方案」其实是把第二个 wg 换成了 chan 凑数,语义反而更绕。 原架构里 workerWg + collectWg 双 wg 才是量身之选。
为什么不能让 worker 自己 close(resultChan)
有人想「让最后一个退出的 worker 关 resultChan」,这条路必然 panic:
- 有 3 个 worker,A worker 的 range 退出不代表 B、C 也退出了。
- A 贸然
close(resultChan),会和 B、C 还在写的情形冲突 →panic: send on closed channel。 - 要安全 close,只能等所有 worker 都退出——这天然就需要一个专门等 worker 的 wg。
worker 之间是平行的,没有一个可靠的「最后一个」概念。
七、happens-before:为什么 collectWg.Wait() 能保证 append 可见
sync.WaitGroup 的关键不只是一个计数器,而是它建立的 happens-before 关系:
在
Wait返回之前,所有Done调用对内存的写入,对Wait之后的代码可见。
应用到本例:
- 收集协程
append(results, val)→ 后续defer collectWg.Done()。 - 主协程
collectWg.Wait()解除阻塞 → 紧接着return results。
由于 Done → Wait 的 happens-before,收集协程对 results 的 append 写,对主协程 return 时的读一定可见。这就是为什么两个 wg 方案没有 data race 的语义保证。
八、WaitGroup vs Channel:何时用哪个
| 场景 | 推荐 |
|---|---|
| 等 N 个 goroutine 都完成 | WaitGroup |
| 等待一个事件发生 | 无缓冲 chan(done := make(chan struct{})) |
| 在 goroutine 间传递值 | 带缓冲/无缓冲 chan |
| 需要取消/超时 | context + chan |
| 需要收集首批错误 | errgroup.Group |
经验法则:「等 N 个完成」用 wg,「等 1 个事件」用 chan。 本例里 collectWg 等的是 1 个收集协程,所以 chan 也能胜任——这正是单 wg 方案可行的依据。
九、常见误用速查表
| 误用 | 后果 | 正确 |
|---|---|---|
| goroutine 内 Add | Wait 可能在 Add 前返回 | 启动前 Add |
| Add 与 Wait 并发 | panic: reused before previous Wait returned | Add 全部在 Wait 前完成 |
| Done 多于 Add | panic: negative counter | 保证 Add/Done 配对 |
| 按值传 WaitGroup | 同步失效 | 用指针 |
| 多个 goroutine 同时 Wait | 安全,可正常返回 | 无需特殊处理 |
| 用 wg 替代 chan 传值 | 语义错位 | wg 只做同步,值用 chan |
十、一句话总结
Add 在启动前、Done 在协程内、Wait 在外部等。一个 wg 只管一件事:生产或消费,各用各的 wg,不要让一个计数器串起两段本应独立的 happens-before 链。