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

如何使用多步骤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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.27 00:57:06