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

基于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)
		}
	}
}

关键改动说明
  1. 用sync.WaitGroup跟踪任务状态:

    • 发送任务前调用wg.Add(totalTasks)注册任务总数;
    • 每个Worker处理完任务后调用wg.Done()递减计数;
    • 任务发送goroutine调用wg.Wait()等待所有任务完成,完全避免了自定义计数的竞态问题。
  2. 关闭channel通知Worker退出:

    • 所有任务发送完毕后关闭workCh,Worker的for inp := range input循环会在channel关闭后自动退出,无需无限循环;
    • 任务全部完成后关闭result channel,主循环通过res, ok := <-result判断channel是否关闭,从而优雅退出。
  3. 给channel增加缓冲:

    • 给workCh和result加适当缓冲,减少goroutine之间的阻塞,提升程序并发效率,同时避免潜在的死锁场景。
  4. 移除冗余的自定义计数逻辑:

    • 去掉了tskCnt、cntChannel及相关计数函数,用标准库同步原语替代,代码更简洁可靠。

内容的提问来源于stack exchange,提问作者EagleNN

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.27 09:27:23