MRjob:减速机可以执行2次操作吗?

Nic*_*ung 2 python mapreduce mrjob

我试图得出映射器生成的每个键值对的概率.

所以,让我们说mapper产量:

a, (r, 5)
a, (e, 6)
a, (w, 7)
Run Code Online (Sandbox Code Playgroud)

我需要添加5 + 6 + 7 = 18,然后找到概率5/18,6/18,7/18

所以减速器的最终输出看起来像:

a, [[r, 5, 0.278], [e, 6, 0.33], [w, 7, 0.389]]
Run Code Online (Sandbox Code Playgroud)

到目前为止,我只能得到reducer来从值中求和所有整数.如何让它返回并将每个实例除以总和?

谢谢!

小智 5

Pai的解决方案在技术上是正确的,但实际上这会给你带来很大的冲突,因为设置分区可能会很麻烦(参见https://groups.google.com/forum/#!topic/mrjob/aV7bNn0sJ2k).

您可以使用mrjob.step更轻松地完成此任务,然后创建两个reducers,例如在此示例中:https://github.com/Yelp/mrjob/blob/master/mrjob/examples/mr_next_word_stats.py

要按照你所描述的方式去做:

from mrjob.job import MRJob
import re
from mrjob.step import MRStep
from collections import defaultdict

wordRe = re.compile(r"[\w]+")

class MRComplaintFrequencyCount(MRJob):

    def mapper(self, _, line):
        self.increment_counter('group','num_mapper_calls',1)

        #Issue is third column in csv
        issue = line.split(",")[3]

        for word in wordRe.findall(issue):
            #Send all map outputs to same reducer
            yield word.lower(), 1

    def reducer(self, key, values):
        self.increment_counter('group','num_reducer_calls',1)  
        wordCounts = defaultdict(int)
        total = 0         
        for value in values:
            word, count = value
            total+=count
            wordCounts[word]+=count

        for k,v in wordCounts.iteritems():
            # word, frequency, relative frequency 
            yield k, (v, float(v)/total)

    def combiner(self, key, values):
        self.increment_counter('group','num_combiner_calls',1) 
        yield None, (key, sum(values))


if __name__ == '__main__':
    MRComplaintFrequencyCount.run()
Run Code Online (Sandbox Code Playgroud)

这样做的标准字数和聚合主要在组合器中,然后使用"无"作为公共密钥,因此每个字间接地在相同的密钥下发送到reducer.在减速器中,您可以获得总字数并计算相对频率.