MapReduce按年月统计评分记录数的计数与时间戳转换问题
MapReduce按年月统计评分记录数问题修正
两个核心疑问的明确结论
- 时间戳格式转换必须放在Mapper阶段完成:MapReduce的Shuffle流程会严格按照Mapper输出的Key做排序、分桶、分组,Reducer接收到的已经是同Key的聚合数据集。如果不在Mapper端把秒级Unix时间戳转换成
年-月格式,Shuffle会把每一个不同秒值的时间戳当成独立Key分组,完全达不到按年月聚合的效果,等数据到Reducer再做格式转换没有任何聚合意义。 - 记录总数统计逻辑需要调整:你当前Reducer里的
count += int(value)是在累加评分值,算的是评分总和,不是记录条数。统计记录条数只需要对同分组下的有效数据行做计数即可,和评分具体数值无关。
修正后可直接运行的代码
Mapper端代码
import sys from datetime import datetime for line in sys.stdin: fields = line.strip().split("\t") # 跳过格式异常的脏行 if len(fields) < 4: continue # 秒级时间戳转年-月格式作为聚合Key ts = int(fields[3]) agg_key = datetime.fromtimestamp(ts).strftime("%Y-%m") # Value固定输出1,代表1条有效评分记录,方便后续累加计数,也支持Combiner优化 print(agg_key, 1, sep="\t")
Reducer端代码
import sys (last_key, record_count) = (None, 0) for line in sys.stdin: line = line.strip() if not line: continue current_key, _ = line.split("\t", 1) # 切换到新Key时,输出上一个年月的累计记录数 if last_key is not None and last_key != current_key: print(last_key, record_count, sep="\t") record_count = 0 last_key = current_key # 每读一条有效记录,计数+1 record_count += 1 # 兜底输出最后一个年月的统计结果 if last_key is not None: print(last_key, record_count, sep="\t")
小提示:Mapper端固定输出1做计数是这类统计任务的标准写法,后续如果集群资源紧张可以直接加Combiner做Map端预聚合,大幅减少Shuffle阶段的数据传输量,跑数速度会快很多。
内容的提问来源于stack exchange,提问作者NicoleL
相关产品推荐
相关产品推荐

