【问题标题】:Why does this goroutine not call wg.Done()?为什么这个 goroutine 不调用 wg.Done()?
【发布时间】:2020-03-03 16:45:02
【问题描述】:

假设任何时候 registerChan 上最多有两个元素(工作地址)。然后由于某种原因,以下代码在最后两个 goroutine 中没有调用 wg.Done()。

func schedule(jobName string, mapFiles []string, nReduce int, phase jobPhase, registerChan chan string) {
    var ntasks int
    var nOther int // number of inputs (for reduce) or outputs (for map)
    switch phase {
    case mapPhase:
        ntasks = len(mapFiles)
        nOther = nReduce
    case reducePhase:
        ntasks = nReduce
        nOther = len(mapFiles)
    }

    fmt.Printf("Schedule: %v %v tasks (%d I/Os)\n", ntasks, phase, nOther)

    const rpcname = "Worker.DoTask"
    var wg sync.WaitGroup
    for taskNumber := 0; taskNumber < ntasks; taskNumber++ {
        file := mapFiles[taskNumber%len(mapFiles)]
        taskArgs := DoTaskArgs{jobName, file, phase, taskNumber, nOther}
        wg.Add(1)
        go func(taskArgs DoTaskArgs) {
            workerAddr := <-registerChan
            print("hello\n")
            // _ = call(workerAddr, rpcname, taskArgs, nil)
            registerChan <- workerAddr
            wg.Done()
        }(taskArgs)
    }
    wg.Wait()
    fmt.Printf("Schedule: %v done\n", phase)
}

如果我将wg.Done() 放在registerChan &lt;- workerAddr 之前,它工作得很好,我不知道为什么。我也尝试过推迟 wg.Done() ,但即使我期望它似乎也不起作用。我想我对 goroutines 和 channels 的工作方式有一些误解,这导致了我的困惑。

【问题讨论】:

  • 当你启动这个函数的时候,channel中最多有两个worker地址,但是如果你消费了一个,那么你在channel的其他地方添加了吗?如果是这样,当你想再次添加worker地址时,通道可能已满,这将阻塞。
  • 您应该想出一个重点突出、可运行的示例来显示该行为。您的问题可能与 sync.Workgroups 完全无关,而是与错误的同步和阻塞有关。

标签: go concurrency synchronization channel


【解决方案1】:

因为它在此处停止

workerAddr := <-registerChan

对于缓冲通道:
要让这个workerAddr := &lt;-registerChan 工作:频道registerChan 必须有一个值;否则,代码将在此处停止等待频道


我设法以这种方式运行您的代码(试试this):

package main

import (
    "fmt"
    "sync"
)

func main() {
    registerChan := make(chan int, 1)
    for i := 1; i <= 10; i++ {
        wg.Add(1)
        go fn(i, registerChan)
    }
    registerChan <- 0 // seed
    wg.Wait()
    fmt.Println(<-registerChan)
}

func fn(taskArgs int, registerChan chan int) {
    workerAddr := <-registerChan
    workerAddr += taskArgs
    registerChan <- workerAddr
    wg.Done()
}

var wg sync.WaitGroup

输出:

55

说明:
这段代码使用一个通道和 10 个 goroutine 加上一个主 goroutine 将 1 加到 10。

我希望这会有所帮助。

【讨论】:

    【解决方案2】:

    当你运行这条语句registerChan &lt;- workerAddr时,如果通道容量已满,你不能添加它,它会阻塞。如果你有一个池,比如 10 个 workerAddr,你可以在调用 schedule 之前将它们全部添加到容量为 10 的缓冲通道中。调用后不要添加,以保证如果从通道中取值,之后有空间再次添加。在 goroutine 的开头使用 defer 很好。

    【讨论】:

      猜你喜欢
      • 2017-07-14
      • 2013-06-25
      • 1970-01-01
      • 1970-01-01
      • 2017-08-27
      • 1970-01-01
      • 1970-01-01
      • 2020-05-24
      • 1970-01-01
      相关资源
      最近更新 更多