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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.29 07:42:43