如何使用多步骤mrjob处理JSON文件实现数据统计需求
基于mrjob的JSON评论数据指标统计实现
你需要对原有高频词统计代码做如下修改,核心逻辑改为逐行解析JSON、按需求聚合对应字段即可,不需要原有正则匹配单词、停用词过滤的相关逻辑。
完整修改后代码
import json from mrjob.job import MRJob from mrjob.step import MRStep class MRReviewStats(MRJob): OUTPUT_PROTOCOL = JSONValueProtocol def mapper_parse_json(self, _, line): # 逐行解析JSON对象 try: review = json.loads(line.strip()) # 1. 为品类平均星级统计产出KV:key是品类,value是(星级, 计数1) category = review.get('product_category') stars = int(review.get('stars', 0)) if category and stars: yield ('category', category), (stars, 1) # 2. 为最高评论量语言统计产出KV:key是语言,value是计数1 lang = review.get('language') if lang: yield ('lang', lang), 1 except: # 跳过格式错误的行 pass def reducer_aggregate(self, key, values): type_tag, field_val = key if type_tag == 'category': # 聚合品类的总星级和总评论数 total_stars = 0 total_cnt = 0 for stars, cnt in values: total_stars += stars total_cnt += cnt avg_star = round(total_stars / total_cnt, 2) yield None, ('品类平均星级', field_val, avg_star) elif type_tag == 'lang': # 聚合各语言的总评论数 total_cnt = sum(values) yield None, ('语言评论数', field_val, total_cnt) def reducer_output_result(self, _, items): lang_counts = [] category_results = [] for item in items: if item[0] == '品类平均星级': category_results.append(item) else: lang_counts.append((item[2], item[1])) # 输出所有品类平均星级 for res in category_results: yield res # 输出评论量最高的语言 if lang_counts: max_count, max_lang = max(lang_counts) yield ('评论量最高的语言', max_lang, max_count) def steps(self): return [ MRStep(mapper=self.mapper_parse_json, reducer=self.reducer_aggregate), MRStep(reducer=self.reducer_output_result) ] if __name__ == '__main__': MRReviewStats.run()
代码逻辑说明
- 输入解析:mapper阶段逐行读取JSON格式的评论数据,用
json.loads解析为Python字典,提取需要的字段产出对应KV - 聚合统计:第一个reducer分别计算每个品类的总星级、总评论数,以及每个语言的总评论量
- 结果输出:第二个reducer汇总所有结果,直接输出所有品类的平均星级,同时从语言统计结果中取最大值得到评论量最高的语言
运行方式
将代码保存为review_stats.py,输入文件为每行一个JSON对象的文本文件,运行命令如下:python review_stats.py 你的输入文件路径 > 输出结果文件
内容的提问来源于stack exchange,提问作者PabloGS
相关产品推荐
相关产品推荐

