You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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通道

  • 新增jobSendWg WaitGroup,每启动一个发送jobs的goroutine就调用Add(1),goroutine执行完发送操作后调用Done()。
  • 等待所有发送jobs的goroutine完成后,再调用close(jobs),确保无goroutine在通道关闭后发送数据。

2. 解决计数竞态问题

  • 使用sync/atomic包的原子操作修改jobWrites和nResults,避免并发修改导致的计数错误。

3. 安全关闭results通道

  • 新增resultSendWg WaitGroup,每启动一个向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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.11 20:50:35