
本文介绍使用 sync.waitgroup 配合 goroutine 实现对未知数量数据项的并发处理,确保主协程准确等待所有子任务完成,并安全获取结果。
本文介绍使用 sync.waitgroup 配合 goroutine 实现对未知数量数据项的并发处理,确保主协程准确等待所有子任务完成,并安全获取结果。
在 Go 中处理“未知数量”的批量任务(如递归搜索、动态生成的数据流或外部输入源)时,简单地用 go func() 启动协程容易引发竞态或提前退出问题——主函数可能在子任务完成前就结束,导致结果丢失。此时,sync.WaitGroup 是最标准、最可靠的同步机制,它专为“等待一组动态启动的 goroutine 完成”而设计。
✅ 正确模式:WaitGroup + goroutine 分发
核心原则是:先调用 wg.Add(1) 再启动 goroutine,确保计数器在协程启动前已更新;协程内务必通过 defer wg.Done() 标记完成。示例如下:
package main
import (
"fmt"
"sync"
)
func main() {
// 模拟未知数量的待处理项(例如:从文件、API 或 channel 动态读取)
items := []int{1, 2, 3, 4, 5, 6, 7, 8, 9, 10} // 实际中可能是 len 未知的切片或从 chan 接收
var wg sync.WaitGroup
results := make(chan int, len(items)) // 缓冲通道,安全接收结果
// 并发处理每个 item
for _, item := range items {
wg.Add(1)
go func(val int) {
defer wg.Done()
// 模拟耗时处理(如网络请求、文件解析、条件匹配等)
if val%3 == 0 { // 示例:只保留能被 3 整除的项(模拟“找到目标”逻辑)
results <- val
}
}(item)
}
// 启动 goroutine 关闭结果通道(避免阻塞)
go func() {
wg.Wait()
close(results)
}()
// 收集所有结果
found := []int{}
for res := range results {
found = append(found, res)
}
fmt.Printf("Found %d items: %v\n", len(found), found)
}⚠️ 关键注意事项
-
不要在循环内直接传入循环变量地址:如
go process(&item, wg)会导致所有 goroutine 共享同一内存地址,值被覆盖。应显式传值(如go func(val int){...}(item))。 -
wg.Add()必须在go之前调用:否则存在竞态风险(wg.Wait()可能因计数未增加而立即返回)。 -
结果收集需线程安全:若写入共享切片,需加锁;推荐使用带缓冲的 channel,由单独 goroutine 负责关闭,主 goroutine 通过
range安全消费。 -
避免过度并发:对成千上万个 item 直接启动同等数量 goroutine 可能导致调度压力。生产环境建议结合
semaphore或 worker pool 限流(如使用带缓冲 channel 控制并发数)。
? 总结
sync.WaitGroup 是 Go 中协调未知数量 goroutine 的基石工具。它不依赖 channel 的复杂信号传递,语义清晰、开销极小。搭配合理的结果收集方式(如结果 channel),即可构建健壮、可扩展的并行处理流程——无论数据源来自 slice、channel 还是实时流,都能优雅应对“数量未知”这一常见场景。

















