Go并发:安全关闭通道并避免向已关闭通道发送数据
问题根源分析
1. 计数不一致的原因
jobWrites和nResults是被多个goroutine并发修改的普通int64变量,无任何同步机制(互斥锁、原子操作),存在竞态条件,导致最终计数结果不准。- worker函数中启动的匿名goroutine是异步执行的,
wg.Wait()仅等待worker函数本身结束,不等待这些匿名goroutine完成results写入,因此nResults并未统计全所有结果写入操作。
2. 向已关闭通道发送的原因
main函数中启动发送jobs的goroutine是异步的,循环结束后立刻调用 close(jobs),但此时仍有大量发送jobs的goroutine未执行到 jobs <- j 操作,通道关闭后再发送就会触发panic。
解决方案
1. 安全关闭jobs通道
- 新增
jobSendWgWaitGroup,每启动一个发送jobs的goroutine就调用Add(1),goroutine执行完发送操作后调用Done()。 - 等待所有发送jobs的goroutine完成后,再调用
close(jobs),确保无goroutine在通道关闭后发送数据。
2. 解决计数竞态问题
- 使用
sync/atomic包的原子操作修改jobWrites和nResults,避免并发修改导致的计数错误。
3. 安全关闭results通道
- 新增
resultSendWgWaitGroup,每启动一个向results发送数据的匿名goroutine就调用Add(1),发送完成后调用Done()。 - 等所有worker结束、所有向results发送数据的goroutine完成后,再调用
close(results),接收循环可正常遍历所有数据后退出,不会死锁。
修改后的完整代码
package main import ( "fmt" "sync" "sync/atomic" ) func worker(wg *sync.WaitGroup, resultSendWg *sync.WaitGroup, jobs <-chan int, results chan<- int, nResults *int64) { defer wg.Done() for j := range jobs { atomic.AddInt64(nResults, 1) resultSendWg.Add(1) go func(j int) { defer resultSendWg.Done() if j%2 == 0 { results <- j * 2 } else { results <- j } }(j) } } func main() { workerWg := &sync.WaitGroup{} jobSendWg := &sync.WaitGroup{} resultSendWg := &sync.WaitGroup{} jobs := make(chan int) results := make(chan int) var i int = 1 var jobWrites int64 = 0 for i <= 10000000 { jobSendWg.Add(1) go func(j int) { defer jobSendWg.Done() if j%2 == 0 { atomic.AddInt64((*int64)(&i), 99) j += 99 } atomic.AddInt64(&jobWrites, 1) jobs <- j }(i) i += 1 } var nResults int64 = 0 for w := 1; w < 1000; w++ { workerWg.Add(1) go worker(workerWg, resultSendWg, jobs, results, &nResults) } // 等待所有jobs发送完成后关闭通道 go func() { jobSendWg.Wait() close(jobs) }() // 等待所有worker完成 workerWg.Wait() // 等待所有result发送完成后关闭通道 go func() { resultSendWg.Wait() close(results) }() var sum int32 = 0 var count int64 = 0 for r := range results { count += 1 sum += int32(r) } fmt.Println(sum) fmt.Printf("number of result writes %d\n", atomic.LoadInt64(&nResults)) fmt.Printf("Number of job writes %d\n", atomic.LoadInt64(&jobWrites)) }
关键修改点说明
- 通道关闭时机:通过两个WaitGroup分别等待jobs发送方和results发送方全部完成,再关闭对应通道,彻底避免向已关闭通道发送的panic。
- 并发计数安全:所有跨goroutine的计数修改都使用
atomic.AddInt64和atomic.LoadInt64,消除竞态条件,保证计数准确。 - 避免死锁:results通道在所有发送操作完成后才关闭,接收循环可正常遍历所有数据后自动退出,不会出现死锁或提前关闭通道导致的发送失败。
内容的提问来源于stack exchange,提问作者Mayuresh Anand
相关产品推荐
相关产品推荐

