Go Worker Pool并发词频统计程序输出结果不一致的排查与修复方案
你的程序出现结果不一致的情况,主要源于两个关键问题:countWord函数的逻辑错误,以及**main函数与computeTotal goroutine之间的同步缺失**。让我们逐个分析并修复:
1. 修复countWord函数的逻辑bug
先看你的countWord实现,这里有两个明显的错误:
- 当单词第一次出现时,你没有将其加入
tempMap,导致后续遇到同一个单词时仍然认为它不存在; - 当单词存在时,你在
tempMap[word]++后返回了tempMap[word]+1,这会导致统计的计数比实际多1。
修复后的countWord函数可以简化为(直接在worker中统计更高效,也能避免逻辑错误):
// 可以直接去掉countWord函数,在worker里直接统计 func worker(wg *sync.WaitGroup) { defer wg.Done() var tempMap = make(map[string]int) for w := range words { tempMap[w]++ // 不管单词是否存在,直接累加(不存在时默认初始为0,加1后为1) } // 处理完所有单词后批量发送结果 for word, count := range tempMap { resultC <- Result{word, count} } }
如果你想保留countWord函数,修正后的版本应该是:
func countWord(word string, tempMap map[string]int) Result { _, ok := tempMap[word] if ok { tempMap[word]++ } else { tempMap[word] = 1 // 首次出现时初始化计数 } return Result{word, tempMap[word]} // 返回当前正确的计数 }
2. 修复main与computeTotal的同步问题
你的main函数在workerPool()执行完成后立即打印total,但此时computeTotal goroutine可能还在处理resultC通道中的数据(尤其是添加fmt.Println(i)拖慢了computeTotal的执行速度时)。workerPool关闭resultC后,computeTotal需要把通道里所有缓冲的结果处理完才能完成统计,但main没有等待这个过程就输出了total,导致结果不完整。
我们可以用sync.WaitGroup来实现同步:
修改computeTotal函数
添加WaitGroup参数,确保函数结束时通知等待组:
func computeTotal(wg *sync.WaitGroup) { defer wg.Done() // 函数结束时标记完成 i := 0 for e := range resultC { total[e.word] += e.count i += 1 fmt.Println(i) } }
修改main函数
初始化等待组,启动computeTotal时增加计数,并在打印结果前等待其完成:
func main() { startTime := time.Now() var computeWg sync.WaitGroup computeWg.Add(1) go readText() go computeTotal(&computeWg) workerPool() // 阻塞等待所有worker完成并关闭resultC computeWg.Wait() // 等待computeTotal处理完所有结果 fmt.Println(total) endTime := time.Now() timeTaken := endTime.Sub(startTime) fmt.Println("total words: ", len(total)) fmt.Println("Time taken for reading the book", timeTaken) }
修复后的完整代码
整合所有修改后的完整代码如下:
package main import ( "bufio" "fmt" "log" "os" "strings" "sync" "time" ) const FILENAME = "cat.txt" const BUFFERSIZE = 3000 const NUMOFWORKER = 5 var words = make(chan string, BUFFERSIZE) // 任务通道 var resultC = make(chan Result, BUFFERSIZE) var total = map[string]int{} type Result struct { word string count int } func readText() { file, err := os.Open(FILENAME) if err != nil { log.Fatal(err) } defer file.Close() scanner := bufio.NewScanner(file) scanner.Split(bufio.ScanWords) for scanner.Scan() { word := strings.ToLower(scanner.Text()) words <- strings.Trim(word, ".,:;") } close(words) } func workerPool() { var wg sync.WaitGroup for i := 1; i <= NUMOFWORKER; i++ { wg.Add(1) go worker(&wg) } wg.Wait() fmt.Println("all goroutines finished") close(resultC) } func worker(wg *sync.WaitGroup) { defer wg.Done() var tempMap = make(map[string]int) for w := range words { tempMap[w]++ } for word, count := range tempMap { resultC <- Result{word, count} } } func computeTotal(wg *sync.WaitGroup) { defer wg.Done() i := 0 for e := range resultC { total[e.word] += e.count i += 1 fmt.Println(i) } } func main() { startTime := time.Now() var computeWg sync.WaitGroup computeWg.Add(1) go readText() go computeTotal(&computeWg) workerPool() computeWg.Wait() fmt.Println(total) endTime := time.Now() timeTaken := endTime.Sub(startTime) fmt.Println("total words: ", len(total)) fmt.Println("Time taken for reading the book", timeTaken) }
为什么原来的代码有时看似正常?
当你移除fmt.Println(i)时,computeTotal的执行速度很快,可能刚好在main函数打印total之前处理完所有结果,这是一种竞态条件下的巧合,并不是可靠的执行顺序。添加fmt.Println(i)后,computeTotal的执行被延迟,main先打印了未完成的total,就暴露了同步问题。
内容的提问来源于stack exchange,提问作者Teriyaki Boy

