MapReduce中使用Combiner计算人均平均消费的reducer问题求解
错误原因
你得到的结果是预期值两倍的核心问题是reducer的计算逻辑错误,不需要额外加分支区分combiner和mapper输出,原因如下:
- mapper输出的value格式为
(单条记录单价, 单条记录购买数量) - combiner已经对同个姓名下的分片数据做了汇总,输出的value格式为
(该分片总消费金额, 该分片总购买数量),其中总消费金额已经是sum(单价*数量)的结果 - 你当前的reducer仍然按照处理mapper原始输出的逻辑,执行
totalprice += (value[0] * value[1]),相当于把已经算好的总消费金额又乘了一次总数量,直接导致总金额翻倍,平均值偏高。
修改后的reducer代码
def reducer(self, key, values): totalprice, totalqty = 0, 0 for value in values: # 直接累加combiner输出的总金额和总数量即可,不需要二次乘法 totalprice += value[0] totalqty += value[1] average = round(totalprice / totalqty, 2) yield key, average
额外优化建议
你提供的CSV文件中部分字段前后带有空格,比如 MATT、 MORAL,会导致同一个姓名被识别为不同的分组,建议修改mapper中的姓名拼接逻辑:
full_name = words[-3].strip() + ' ' + words[-2].strip()
完整可运行代码
from mrjob.job import MRJob class Job(MRJob): def mapper(self, key, value): words = value.strip().split(',') # 新增strip处理字段前后空格 full_name = words[-3].strip() + ' ' + words[-2].strip() price, qty = float(words[-5]), int(words[-4]) yield full_name, (price, qty) def combiner(self, key, values): totalprice, totalqty = 0,0 for value in values: totalprice += (value[0] * value[1]) totalqty += value[1] yield key, (totalprice, totalqty) def reducer(self, key, values): totalprice, totalqty = 0,0 for value in values: totalprice += value[0] totalqty += value[1] average = round(totalprice/totalqty,2) yield key, average if __name__ == '__main__': Job.run()
内容的提问来源于stack exchange,提问作者JohnDoe34
相关产品推荐
相关产品推荐

