MapReduce任务统计异常:五星好评电影计数始终为1
问题描述
我正在开发MapReduce任务,目标是筛选出拥有500条及以上五星好评的动作电影。目前已完成两个前置MapReduce任务:
- 从电影列表中筛选动作电影ID
- 筛选带有单条五星好评的电影ID
新任务的两个Mapper分别接收这两类ID作为输入:
- MapperA输出
<MovieID, "1">,标记该ID属于动作电影 - MapperB输出
<MovieID, "2">,标记该ID对应一条五星好评
我需要在Reducer中完成两个逻辑:统计每个电影ID的五星好评总数(判断是否≥500),同时验证该电影是否属于动作电影。但现在遇到的问题是:统计好评数的HashMap里,每个电影ID的计数始终只有1,无法累计。我怀疑Reducer未收到同一MovieID对应的多个"2"输入,相关代码如下:
public class JoinRatings extends Configured implements Tool { public static class TokenizerMapperA extends Mapper<Object, Text, Text, Text> { private Text node1; private Text node2 = new Text("1"); @Override protected void map(Object key, Text value, Context context) throws IOException, InterruptedException { //Write movie ID and int writeable 1 node1 = new Text(value.toString()); context.write(node1, node2); } } public static class TokenizerMapperB extends Mapper<Object, Text, Text, Text> { private Text node1; private Text node2 = new Text("2"); @Override protected void map(Object key, Text value, Context context) throws IOException, InterruptedException { node1 = new Text(value.toString()); context.write(node1, node2); } } public static class CountReducer extends Reducer<Text, Text, Text, NullWritable>{ private Text node1; private Set<String> distinctNodes; Map<String, Integer> map; private final static IntWritable one = new IntWritable(1); private final static IntWritable two = new IntWritable(2); @Override protected void setup(Context context) { distinctNodes = new HashSet<String>(); map = new HashMap<String,Integer>(); } @Override protected void reduce(Text key, Iterable<Text> values, Context context) { //Add all action movie IDs to unique set distinctNodes String keyString = key.toString().strip(); String valueString = "NULL"; int counter = 0 for (Text text : values) { String value = text.toString().strip(); // try { // Text testText = new Text(value); // context.write(testText, NullWritable.get()); // } catch (IOException e) { // e.printStackTrace(); // } catch (InterruptedException e) { // e.printStackTrace(); // } if (value.equals("1")) { distinctNodes.add(keyString); } else if (value.equals("2")) { if (map.containsKey(keyString)) { map.put(keyString, map.get(keyString) + 1); } else { map.put(keyString, 1); } } } } @Override protected void cleanup(Context context) throws IOException, InterruptedException { String numOfDistinctNodes = Integer.toString(distinctNodes.size()); context.write(new Text(numOfDistinctNodes), NullWritable.get()); String sizeOfMap = Integer.toString(map.size()); context.write(new Text(sizeOfMap), NullWritable.get()); for(Map.Entry<String, Integer> mapElement : map.entrySet()) { String test = Integer.toString(mapElement.getValue()); context.write(new Text(test), NullWritable.get()); } for (String s: distinctNodes) { if (map.containsKey(s)) { if (map.get(s) >= 500) { node1 = new Text(s); context.write(node1, NullWritable.get()); } } } //Order movieID by ascending //Write all movie IDs with 500+ 5-star reviews (to context) } }
问题根源及修复方案
核心逻辑错误
你误用了全局HashMap统计计数——MapReduce的reduce方法是按Key分组执行的:同一个MovieID的所有标记("1"和"2")会被一次性传入当前的reduce调用,根本不需要用全局HashMap跨Key累计。这种写法不仅逻辑错误,还会导致内存占用过高,甚至OOM。
修正后的Reducer代码
public static class CountReducer extends Reducer<Text, Text, Text, NullWritable> { @Override protected void reduce(Text key, Iterable<Text> values, Context context) throws IOException, InterruptedException { boolean isActionMovie = false; int fiveStarCount = 0; String movieId = key.toString().strip(); // 针对当前MovieID统计数据 for (Text text : values) { String value = text.toString().strip(); if ("1".equals(value)) { isActionMovie = true; } else if ("2".equals(value)) { fiveStarCount++; } } // 满足条件直接输出 if (isActionMovie && fiveStarCount >= 500) { context.write(key, NullWritable.get()); } } // 调试用的cleanup代码可按需保留或删除 @Override protected void cleanup(Context context) throws IOException, InterruptedException { // ... } }
额外检查点
- 验证MapperB的输入数据:确认每个五星好评对应的MovieID是否正确输出,可临时在MapperB中输出所有Key-Value,查看同一MovieID是否有多条"2"记录。
- 修复原代码的语法错误:
reduce方法缺少throws IOException, InterruptedException声明,distinct nodes是拼写错误(应为distinctNodes)。
内容的提问来源于stack exchange,提问作者wirly
相关产品推荐
相关产品推荐

