如何在Python的mrjob中按多参数条件统计相同条目的数量
问题核心错误点
- 键设计错误:你需要统计分类+版本两个维度的数量,Mapper输出的键不能只传分类,必须把
(分类, 版本)组成复合键,否则同分类不同版本的数据会被混合统计 - 函数名不匹配:
steps中注册的函数名和实际定义的函数名完全不一致,运行会直接抛出找不到函数的错误 - 缺少Reducer步骤:只有Mapper和Combiner无法完成全局聚合,必须补充Reducer逻辑
- Combiner求和逻辑错误:你传入的value是
(1, 版本号)元组,不能直接调用sum求和
修正后完整代码
from mrjob.job import MRJob from mrjob.step import MRStep class MRFrequencyCount(MRJob): def steps(self): return [ MRStep( mapper=self.mapper_extract, combiner=self.combine_counts, reducer=self.reduce_counts ) ] def mapper_extract(self, _, line): # 跳过表头行 if line.startswith(' product_name'): return # 默认按任意空白符分割,如果你实际文件是*分隔就替换为split('*') fields = line.split() # 按字段顺序取分类和版本,strip处理可能的引号 category = fields[2].strip() version = fields[4].strip('"') # 复合键:(分类, 版本),value为计数1 yield (category, version), 1 def combine_counts(self, key, counts): # key是(分类, 版本),counts是当前分片内该key对应的所有1,直接求和 yield key, sum(counts) def reduce_counts(self, key, counts): category, version = key total = sum(counts) # 输出你要求的格式 yield None, f"<{category}, {{{total}, {version}}}>" if __name__ == '__main__': MRFrequencyCount.run()
逻辑说明
- Mapper阶段:跳过表头后,逐行拆分字段,把
(分类, 版本)作为复合键,value输出1 - Combiner阶段:在每个分片本地先聚合同一个key的计数,减少网络传输数据量
- Reducer阶段:合并所有分片的计数结果,按要求格式化输出
内容的提问来源于stack exchange,提问作者user17488887
相关产品推荐
相关产品推荐

