【问题标题】:Problem synchronizing between Goroutines, one blocks the otherGoroutines之间的同步问题,一个阻塞另一个
【发布时间】:2022-11-05 07:08:32
【问题描述】:

我试图处理从 AWS S3 读取的 CSV 文件,对于每一行我想激活 worker 函数来做一些工作并返回结果

理想情况下,我希望将结果作为原始 CSV 排序,但这不是必需的,由于某种原因,当我运行此代码时,我会遇到奇怪的数据竞争和这一行:

for result := range output {
   results = append(results, result)
}

永远阻塞

我尝试使用也不起作用的 WaitGroup,关闭 output 频道也导致我出现“试图将某些内容放入关闭的频道”的错误

func main() {
    resp, err := ReadCSV(bucket, key)
    if err != nil {
        log.Fatal(err)
    }
    defer resp.Body.Close()
    reader := csv.NewReader(resp.Body)

    detector := NewDetector(languages)
    var results []DetectionResult

    numWorkers := 4
    input := make(chan string, numWorkers)
    output := make(chan DetectionResult, numWorkers)

    start := time.Now()

    for w := 1; w < numWorkers+1; w++ {
        go worker(w, detector, input, output)
    }

    go func() {
        for {
            record, err := reader.Read()
            if err == io.EOF {
                close(input)
                break
            }

            if err != nil {
                log.Fatal(err)
            }

            text := record[0]
            input <- text
        }
    }()

    for result := range output {
        results = append(results, result)
    }

    elapsed := time.Since(start)

    log.Printf("Decoded %d lines of text in %s", len(results), elapsed)
}

func worker(id int, detector lingua.LanguageDetector, input chan string, output chan DetectionResult) {
    log.Printf("worker %d started\n", id)
    for t := range input {
        result := DetectText(detector, t)
        output <- result
    }
    log.Printf("worker %d finished\n", id)
}

尝试处理 CSV(最好按顺序),并使用对 worker 的函数调用结果来丰富它

尝试设置WaitGroup,尝试在完成读取(EOF)时关闭输出通道 - 导致错误

【问题讨论】:

  • 通过编写 4 个文本检测器(工作人员),我从您那里读到了一个隐含的声明,即检测比从 CSV 读取慢 4 倍。那正确吗?我可以看到,也许你正试图检测一种给定未知输入的语言,所以要进行大量的猜测和检查?

标签: csv go concurrency channels routines


【解决方案1】:

for 循环将一直读取到output 通道关闭。处理完所有输入后(而不是读完输入),您必须关闭 output 频道。

您可以为此使用等待组:

func worker(detector lingua.LanguageDetector, wg *sync.WaitGroup) func(id int, input chan string, output chan DetectionResult) {
   wg.Add(1)
   return func(id int, input chan string, output chan DetectionResult) {
      defer wg.Done() // Notify wg when processing is finished
      log.Printf("worker %d started
", id)
      for t := range input {
         result := DetectText(detector, t)
         output <- result
      }
      log.Printf("worker %d finished
", id)
   }
}

然后:

go func() {
    wg.Wait()
    close(output)
}()
for result := range output {
        results = append(results, result)
}

【讨论】:

    【解决方案2】:

    我发现您缺少一种向工人发出没有更多工作并且他们应该停止工作的信号。您还需要一种方法让工作人员回复他们确实完成了的信号。当所有这些信号都已发送和接收后,main 应该控制所有工作人员的累积结果。

    我们可以在所有 CSV 记录迭代后通过关闭输入来向工作人员发出信号,并且所有作业都已通过输入发送:

    nWorkers := 4
    
    input := make(chan Tx, nWorkers*2) // buffer so input (the "jobs queue") is always full; see rationale at bottom of answer
    output := make(chan Ty)
    done := make(chan bool)
    
    for i := 1; i < nWorkers+1; i++ {
        go worker(input, output, done)
    }
    
    go func() {
        for {
            record, _ := reader.Read()
            input <- record[0]
        }
        close(input)
    }()
    

    当没有更多作业时,通过输入发送作业的 goroutine 可以安全地关闭输入。即使在关闭之后,工人仍然可以接收仍在输入中的任何工作。

    当输入关闭并最终为空时,工作人员的范围循环退出。然后,工作人员通过在 done 通道上发送信号返回:

    func worker(input <-chan Tx, output chan<- Ty, done <-chan bool) {
        for x := range input { // loop until input is closed
            output <- doWork(x)
        }
        done <- true // finally send done
    }
    

    当我们收到 nWorker-number of done 消息时,我们知道所有工作都已完成并且工作人员不会发送输出,因此关闭输出是安全的:

    go func() {
        log.Println("counting done workers")
        var doneCtr int
        for {
            select {
            case <-done:
                log.Println("got done")
                doneCtr++
            }
    
            if doneCtr == nWorkers {
                close(output) // signal the results appender to stop
                log.Println("closed output")
            }
        }
    }()
    

    关闭输出是向 main 发出的信号,它可以停止尝试接收和累积结果:

    results := make([]result, 0)
    for result := range output {
        results = append(results, result)
    }
    

    最后:所有其他 goroutine 都已终止,main 可以继续累积结果。

    至于按原始顺序获取结果,只需将原始订单与每个作业一起发送,将该订单与结果一起发送回,然后按订单排序:

    type row struct {
        num  int
        text string
    }
    
    type result struct {
        lang language
        row  row
    }
    
    ...
    
    input <- row{rowNum, record[0]}
    rowNum++
    
    ...
    
    output <- result{detect(row.text), row}
    
    ...
    
    results = append(results, result)
    
    ...
    
    sort.Slice(results, func(i, j int) bool { return results[i].row.num < results[j].row.num})
    

    我在The Go Playground 中制作了一个完整的工作模型。

    我的推理可能与缓冲有关,但是,正如我所见,唯一真正令人失望的是发现工人在等待输入时停滞不前。将输入缓冲 2 倍工人数量可确保每个工人在任何时刻平均有两个工作在等待。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2022-11-02
      • 2023-03-28
      • 2017-11-01
      • 1970-01-01
      • 1970-01-01
      • 2015-04-29
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多