【发布时间】: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 <- workerAddr 之前,它工作得很好,我不知道为什么。我也尝试过推迟 wg.Done() ,但即使我期望它似乎也不起作用。我想我对 goroutines 和 channels 的工作方式有一些误解,这导致了我的困惑。
【问题讨论】:
-
当你启动这个函数的时候,channel中最多有两个worker地址,但是如果你消费了一个,那么你在channel的其他地方添加了吗?如果是这样,当你想再次添加worker地址时,通道可能已满,这将阻塞。
-
您应该想出一个重点突出、可运行的示例来显示该行为。您的问题可能与 sync.Workgroups 完全无关,而是与错误的同步和阻塞有关。
标签: go concurrency synchronization channel