基于Channels的多Goroutine同步问题排查与解决
你的代码出现值丢失或退出时机错误的问题,主要源于几个关键的同步逻辑缺陷:
1. 未同步的内存访问引发竞态条件
全局变量tskCnt由taskCounter goroutine修改,但main goroutine直接读取它的状态(for tskCnt != 0 {})时,没有任何同步机制。根据Go的内存模型,不同goroutine之间的内存访问如果没有通过channel、互斥锁等同步原语,无法保证可见性——main可能一直看到tskCnt的旧值,导致程序提前退出或者无限等待,看起来像是"丢失值"。
2. 无缓冲channel的阻塞顺序导致计数偏差
cntChannel是无缓冲channel,putTask里先发送+1再发送任务到workCh,getTask里先发送-1再处理任务。但无缓冲channel的发送会阻塞直到接收方准备好,这可能导致计数和任务实际状态不一致:比如putTask的cntChannel <-1已经被接收,但workCh <-val还没被Worker接收,此时tskCnt已经+1,但任务还在"传递中";反过来,Worker拿到任务后先发送-1,但任务还没处理完,tskCnt已经-1,也会导致计数不准。
3. 轮询式等待的不可靠性
你用多个for len(...) != 0 {}和for tskCnt !=0 {}来等待任务完成,这种轮询方式本身就不可靠:
len(cntChannel)只能看到当前channel里未被接收的元素数,无法反映任务的真实处理状态;- 对
tskCnt的读取没有同步,无法保证看到最新值; - 即使这些条件都满足,Worker还在无限循环等待新任务,程序退出时会留下僵尸goroutine。
Go标准库的sync.WaitGroup就是为这类"等待一批任务完成"的场景设计的,配合channel关闭机制,可以优雅实现你的需求。下面是修正后的代码:
import ( "fmt" "os" "os/signal" "strconv" "sync" ) const numWorkers = 5 type workerChannel chan uint64 func initWorker(input workerChannel, result chan string, num int, wg *sync.WaitGroup) { // 遍历channel,直到channel关闭自动退出 for inp := range input { defer wg.Done() // 任务完成后递减WaitGroup计数 result <- fmt.Sprintf("Worker %d:%d", num, inp) } } func main() { abort := make(chan os.Signal, 1) signal.Notify(abort, os.Interrupt) // 给结果channel加缓冲,避免Worker因发送结果阻塞 result := make(chan string, numWorkers) // 任务channel加缓冲,减少任务发送时的阻塞 workCh := make(workerChannel, numWorkers) var wg sync.WaitGroup // 启动Worker for i := 0; i < numWorkers; i++ { go initWorker(workCh, result, i, &wg) } // 发送任务并等待完成 go func() { totalTasks := 21 wg.Add(totalTasks) // 提前注册需要等待的任务总数 for i := uint64(0); i < uint64(totalTasks); i++ { fmt.Printf("Put task %d\n", i) workCh <- i } close(workCh) // 所有任务发送完毕,关闭任务channel通知Worker退出 wg.Wait() // 等待所有任务处理完成 close(result) // 关闭结果channel,通知主循环可以退出 }() // 主循环处理结果和退出信号 for { select { case res, ok := <-result: if !ok { // 结果channel关闭,说明所有结果都已处理 fmt.Println("Done") os.Exit(0) } fmt.Println(res) case <-abort: fmt.Println("Aborted.") os.Exit(0) } } }
用
sync.WaitGroup跟踪任务状态:- 发送任务前调用
wg.Add(totalTasks)注册任务总数; - 每个Worker处理完任务后调用
wg.Done()递减计数; - 任务发送goroutine调用
wg.Wait()等待所有任务完成,完全避免了自定义计数的竞态问题。
- 发送任务前调用
关闭channel通知Worker退出:
- 所有任务发送完毕后关闭
workCh,Worker的for inp := range input循环会在channel关闭后自动退出,无需无限循环; - 任务全部完成后关闭
resultchannel,主循环通过res, ok := <-result判断channel是否关闭,从而优雅退出。
- 所有任务发送完毕后关闭
给channel增加缓冲:
- 给
workCh和result加适当缓冲,减少goroutine之间的阻塞,提升程序并发效率,同时避免潜在的死锁场景。
- 给
移除冗余的自定义计数逻辑:
- 去掉了
tskCnt、cntChannel及相关计数函数,用标准库同步原语替代,代码更简洁可靠。
- 去掉了
内容的提问来源于stack exchange,提问作者EagleNN

