You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.24 15:54:06