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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.30 17:32:51