使用Worker Pool构建通用管道库时的close与send竞态问题排查
你的问题核心在于WaitGroup的使用方式错误,以及对Worker Pool生命周期的同步逻辑理解有误。我们先看修正后的代码,再一步步分析问题根源:
修正后的代码
func (p *pipe) Process(in chan interface{}) (out chan interface{}) { var wg sync.WaitGroup out = make(chan interface{}, 100) workerCount := 100 // 为每个worker预注册WaitGroup计数 wg.Add(workerCount) for i := 0; i < workerCount; i++ { go func() { defer wg.Done() // worker退出时标记完成 // 循环读取输入通道,直到输入通道被关闭 for j := range in { res := doSomethingWith(j) out <- res } }() } // 单独goroutine等待所有worker完成,再关闭输出通道 go func() { wg.Wait() close(out) }() return out }
原代码的核心错误分析
WaitGroup的Add与Done顺序颠倒
你在处理每个任务的匿名函数中先defer wg.Done(),再调用wg.Add(1)。这会导致致命问题:如果doSomethingWith(j)执行极快,Done()可能在Add()之前就被触发,直接造成WaitGroup misuse: Done called too many times的panic。错误地用WaitGroup跟踪任务而非Worker生命周期
原代码试图用WaitGroup跟踪每个任务的完成,但这完全不符合Worker Pool的逻辑:当部分任务完成后,WaitGroup计数器可能瞬间归零,导致wg.Wait()提前执行并关闭out通道;但此时输入通道in可能还在发送数据,剩余Worker仍在运行,向已关闭的out发送数据就会触发send on closed channel的panic。数据竞态的根源
由于close(out)和out <- res之间没有正确的同步机制,Race Detector检测到两者可能并发执行——原代码的WaitGroup逻辑无法保证所有Worker都停止向out发送数据后再关闭通道。
修正后的逻辑解释
跟踪Worker的生命周期
我们用WaitGroup跟踪每个Worker的完整生命周期:启动前调用wg.Add(workerCount)预注册所有Worker,每个Worker退出时(即输入通道in被关闭,for range循环结束)调用wg.Done()。正确的通道关闭时机
单独启动一个goroutine等待所有Worker完成(wg.Wait()),此时意味着所有输入任务都已处理完毕,再关闭out通道,从根本上避免了向关闭通道发送数据的问题。外部调用的注意事项
使用这个修正后的Process函数时,外部必须在所有数据发送到in通道后关闭in——这是Worker能够退出for range循环的前提,否则Worker会一直阻塞等待输入,wg.Wait()永远不会完成,out通道也不会被关闭。
内容的提问来源于stack exchange,提问作者Nicky Logan

