【发布时间】:2016-08-21 23:55:54
【问题描述】:
我正在尝试在 Golang 中实现类似 mapreduce 的方法。我的设计如下:
映射工作人员将项目从映射器输入通道中拉出并输出到映射器输出通道
映射器输出通道然后由单个 goroutine 读取。此例程维护先前看到的键值对的映射。如果映射器输出中的下一项具有匹配键,它会将具有匹配键的新旧值都发送到 reduce-input 通道。
- reduce-input 管道将两个值缩减为一个键值对,并将结果提交到同一个 map-output 通道。
这导致了 mapper 输出和 reduce 输入之间的循环依赖,我现在不知道如何表示 mapper 输出已完成(并关闭通道)。
打破这种循环依赖或知道何时关闭具有这种循环行为的通道的最佳方法是什么?
下面的代码有一个死锁,map输出通道和reduce输入通道互相等待。
type MapFn func(input int) (int, int)
type ReduceFn func(a int, b int) int
type kvPair struct {
k int
v int
}
type reducePair struct {
k int
v1 int
v2 int
}
func MapReduce(mapFn MapFn, reduceFn ReduceFn, input []int, nMappers int, nReducers int) (map[int]int, error) {
inputMapChan := make(chan int, len(input))
outputMapChan := make(chan *kvPair, len(input))
reduceInputChan := make(chan *reducePair)
outputMapMap := make(map[int]int)
go func() {
for v := range input {
inputMapChan <- v
}
close(inputMapChan)
}()
for i := 0; i < nMappers; i++ {
go func() {
for v := range inputMapChan {
k, v := mapFn(v)
outputMapChan <- &kvPair{k, v}
}
}()
}
for i := 0; i < nReducers; i++ {
go func() {
for v := range reduceInputChan {
reduceValue := reduceFn(v.v1, v.v2)
outputMapChan <- &kvPair{v.k, reduceValue}
}
}()
}
for v := range outputMapChan {
key := v.k
value := v.v
other, ok := outputMapMap[key]
if ok {
delete(outputMapMap, key)
reduceInputChan <- &reducePair{key, value, other}
} else {
outputMapMap[key] = value
}
}
return outputMapMap, nil
}
【问题讨论】:
-
即使没有
Reducers,outputMapChan也没有关闭,导致永远等待。我认为mapreduce应该包含两个独立的阶段,map和reduce,而不是将它们循环在一起。 -
困难来自
reduce阶段的递归性质。由于减少计算可能需要进一步的folding,并且这可能发生任意次数,我需要一些机制来允许reduce输出的输出流入reduce输入以进一步减少。 -
reduce阶段不具有递归性质。映射的中间结果应该拆分为集合,每个集合都可以由reducer单独处理。这就是为什么mapreduce在分配系统中工作。 -
您可能在描述传统的 marpeduce 架构,但我特别好奇在 Go 中如何处理循环通道依赖关系。 @Amd 接受的答案很好地证明了这一点。不过,谢谢您的建议,因为您的 mapreduce cmets 通常是正确的:)
标签: go concurrency mapreduce channel