MapReduce中首个Reducer输出<sum,count>键值对的类型选择问题
问题分析与解决方案
你遇到的核心问题是Hadoop无法直接使用HashMap作为Map/Reducer的输出值类型——因为HashMap没有实现Hadoop的Writable接口,Hadoop无法完成它的序列化/反序列化操作,这直接导致了MapOutputCollector初始化失败和后续的NullPointerException。
要实现「Combiner输出<sum, count>,Reducer计算平均值」的需求,正确的做法是自定义一个实现Writable接口的类来封装sum和count,而不是用HashMap。下面是具体的修正步骤和完整代码:
1. 自定义Writable类存储Sum和Count
首先创建一个SumCountWritable类,用来封装每门课程的总分和成绩数量,实现Hadoop的序列化逻辑:
import java.io.DataInput; import java.io.DataOutput; import java.io.IOException; import org.apache.hadoop.io.Writable; public class SumCountWritable implements Writable { private float sum; private long count; // 无参构造函数(必须,Hadoop反射需要) public SumCountWritable() {} public SumCountWritable(float sum, long count) { this.sum = sum; this.count = count; } // 序列化方法 @Override public void write(DataOutput out) throws IOException { out.writeFloat(sum); out.writeLong(count); } // 反序列化方法 @Override public void readFields(DataInput in) throws IOException { this.sum = in.readFloat(); this.count = in.readLong(); } // getter和setter方法 public float getSum() { return sum; } public void setSum(float sum) { this.sum = sum; } public long getCount() { return count; } public void setCount(long count) { this.count = count; } // 可选:toString方便调试 @Override public String toString() { return sum + "," + count; } }
2. 修改Mapper类
Mapper不再输出HashMap,而是直接输出封装好的SumCountWritable——每个成绩对应一个sum为成绩值、count为1的实例:
public static class MapForAverage extends Mapper<LongWritable, Text, LongWritable, SumCountWritable> { @Override public void map(LongWritable key, Text value, Context con) throws IOException, InterruptedException { String[] word = value.toString().split(", "); float grade = Float.parseFloat(word[1]); int course = Integer.parseInt(word[0]); // 输出课程ID和对应的(成绩, 1) con.write(new LongWritable(course), new SumCountWritable(grade, 1)); } }
3. 修改Combiner类
Combiner的输入是Iterable<SumCountWritable>,我们需要遍历所有值,累加sum和count,然后输出新的SumCountWritable:
public static class ReduceForAverage extends Reducer<LongWritable, SumCountWritable, LongWritable, SumCountWritable> { @Override public void reduce(LongWritable course, Iterable<SumCountWritable> values, Context con) throws IOException, InterruptedException { float totalSum = 0; long totalCount = 0; for (SumCountWritable scw : values) { totalSum += scw.getSum(); totalCount += scw.getCount(); } // 输出课程ID和累加后的(sum, count) con.write(course, new SumCountWritable(totalSum, totalCount)); } }
4. 修改Final Reducer类
Final Reducer同样接收Iterable<SumCountWritable>,累加所有Combiner输出的sum和count,计算平均值后输出:
public static class ReduceForFinal extends Reducer<LongWritable, SumCountWritable, LongWritable, FloatWritable> { private FloatWritable result = new FloatWritable(); @Override public void reduce(LongWritable course, Iterable<SumCountWritable> values, Context con) throws IOException, InterruptedException { float totalSum = 0; long totalCount = 0; for (SumCountWritable scw : values) { totalSum += scw.getSum(); totalCount += scw.getCount(); } // 计算平均值 float average = totalSum / totalCount; result.set(average); con.write(course, result); } }
5. 修改Job配置
更新Job的输出类型配置,替换原来的Object.class为自定义的SumCountWritable.class:
public static void main(String[] args) throws IllegalArgumentException, IOException, ClassNotFoundException, InterruptedException { Configuration conf = new Configuration(); Job job = Job.getInstance(conf, "avg grading"); job.setJarByClass(AvgGrading.class); job.setMapperClass(MapForAverage.class); job.setCombinerClass(ReduceForAverage.class); job.setNumReduceTasks(2); job.setReducerClass(ReduceForFinal.class); // 修改Map输出类型 job.setMapOutputKeyClass(LongWritable.class); job.setMapOutputValueClass(SumCountWritable.class); job.setOutputKeyClass(LongWritable.class); job.setOutputValueClass(FloatWritable.class); FileInputFormat.addInputPath(job, new Path(args[0])); FileOutputFormat.setOutputPath(job, new Path(args[1])); System.exit(job.waitForCompletion(true) ? 0 : 1); }
关键说明
- Hadoop的所有在节点间传输的数据类型必须实现
Writable接口,这是Hadoop序列化框架的要求,普通JDK集合类(比如HashMap)不满足这个条件,所以不能直接使用。 - 你的原代码中Reducer的输入参数写错了(比如
Map<Float,Long> values),正确的Reducer输入值应该是Iterable<Value类型>,因为Reducer会接收同一个Key对应的所有Value的集合。
内容的提问来源于stack exchange,提问作者joy_jlee
相关产品推荐
相关产品推荐

