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

Hadoop MapReduce温度统计程序报错:类型不匹配问题解决

解决Hadoop MapReduce中"wrong value class: Text is not FloatWritable"错误

错误原因分析

出现这个错误的核心原因有两个:

  1. Combiner与Reducer的类型不匹配
    你把TemperatureReducer同时设置为Combiner和Reducer,但TemperatureReducer的输出是<Text, Text>类型,而Reducer的输入需要和Mapper输出一致,即<Text, FloatWritable>。Combiner的输出必须能直接作为Reducer的输入,直接用最终Reducer充当Combiner会导致类型冲突。

  2. Job输出类型配置错误
    程序最终输出的Value类型是Text(Reducer中context.write(key, new Text(...))),但你在Job配置里设置了job.setOutputValueClass(FloatWritable.class),这和实际输出类型不匹配。

解决方案

方案1:移除Combiner(快速解决)

如果不需要Combiner做性能优化,直接删除job.setCombinerClass(TemperatureReducer.class)这一行,同时修正Job的输出Value类型。

方案2:自定义Combiner类(推荐,优化性能)

如果要保留Combiner,需要编写专门的Combiner类,让它的输出类型和Mapper一致,传递sum、count、min、max这些中间统计值,而非直接输出最终文本结果。

修改后的完整代码(方案1:移除Combiner)

import java.io.IOException;
import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.fs.Path;
import org.apache.hadoop.io.FloatWritable;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.mapreduce.Job;
import org.apache.hadoop.mapreduce.Mapper;
import org.apache.hadoop.mapreduce.Reducer;
import org.apache.hadoop.mapreduce.lib.input.TextInputFormat;
import org.apache.hadoop.mapreduce.lib.output.TextOutputFormat;

public class TemperatureStats {

    public static class TemperatureMapper extends Mapper<Object, Text, Text, FloatWritable> {
        private Text year = new Text();
        private FloatWritable temperature = new FloatWritable();

        public void map(Object key, Text value, Context context) throws IOException, InterruptedException {
            String[] tokens = value.toString().split(",");
            String date = tokens[0];
            String[] dateParts = date.split("-");
            String yearStr = dateParts[0];
            year.set(yearStr);
            float temp = Float.parseFloat(tokens[1]);
            temperature.set(temp);
            context.write(year, temperature);
        }
    }

    public static class TemperatureReducer extends Reducer<Text, FloatWritable, Text, Text> {
        public void reduce(Text key, Iterable<FloatWritable> values, Context context)
                throws IOException, InterruptedException {
            float sum = 0;
            int count = 0;
            float minTemp = Float.MAX_VALUE;
            float maxTemp = Float.MIN_VALUE;

            // 计算每年的温度总和、数量、最小值、最大值
            for (FloatWritable val : values) {
                float temp = val.get();
                sum += temp;
                count++;
                minTemp = Math.min(minTemp, temp);
                maxTemp = Math.max(maxTemp, temp);
            }

            // 计算平均温度
            float averageTemp = sum / count;

            // 输出年份和统计结果
            context.write(key, new Text("Average: " + averageTemp + ", Min: " + minTemp + ", Max: " + maxTemp));
        }
    }

    public static void main(String[] args) throws Exception {
        Configuration conf = new Configuration();
        Job job = Job.getInstance(conf, "temperature stats");
        job.setJarByClass(TemperatureStats.class);
        job.setMapperClass(TemperatureMapper.class);
        // 移除Combiner设置
        job.setReducerClass(TemperatureReducer.class);
        job.setOutputKeyClass(Text.class);
        // 修正输出Value类型为Text
        job.setOutputValueClass(Text.class);
        job.setInputFormatClass(TextInputFormat.class);
        job.setOutputFormatClass(TextOutputFormat.class);
        TextInputFormat.addInputPath(job, new Path(args[0]));
        TextOutputFormat.setOutputPath(job, new Path(args[1]));
        System.exit(job.waitForCompletion(true) ? 0 : 1);
    }
}

修改后的完整代码(方案2:自定义Combiner)

如果需要Combiner优化,我们可以用Text传递中间统计值(也可以自定义Writable类型,这里用Text做示例):

import java.io.IOException;
import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.fs.Path;
import org.apache.hadoop.io.FloatWritable;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.mapreduce.Job;
import org.apache.hadoop.mapreduce.Mapper;
import org.apache.hadoop.mapreduce.Reducer;
import org.apache.hadoop.mapreduce.lib.input.TextInputFormat;
import org.apache.hadoop.mapreduce.lib.output.TextOutputFormat;

public class TemperatureStats {

    public static class TemperatureMapper extends Mapper<Object, Text, Text, FloatWritable> {
        private Text year = new Text();
        private FloatWritable temperature = new FloatWritable();

        public void map(Object key, Text value, Context context) throws IOException, InterruptedException {
            String[] tokens = value.toString().split(",");
            String date = tokens[0];
            String[] dateParts = date.split("-");
            String yearStr = dateParts[0];
            year.set(yearStr);
            float temp = Float.parseFloat(tokens[1]);
            temperature.set(temp);
            context.write(year, temperature);
        }
    }

    // 自定义Combiner,输出中间统计值:sum,count,min,max
    public static class TemperatureCombiner extends Reducer<Text, FloatWritable, Text, Text> {
        public void reduce(Text key, Iterable<FloatWritable> values, Context context)
                throws IOException, InterruptedException {
            float sum = 0;
            int count = 0;
            float minTemp = Float.MAX_VALUE;
            float maxTemp = Float.MIN_VALUE;

            for (FloatWritable val : values) {
                float temp = val.get();
                sum += temp;
                count++;
                minTemp = Math.min(minTemp, temp);
                maxTemp = Math.max(maxTemp, temp);
            }
            // 用Text传递中间值,格式为"sum,count,min,max"
            context.write(key, new Text(sum + "," + count + "," + minTemp + "," + maxTemp));
        }
    }

    // Reducer接收Combiner的输出,合并最终结果
    public static class TemperatureReducer extends Reducer<Text, Text, Text, Text> {
        public void reduce(Text key, Iterable<Text> values, Context context)
                throws IOException, InterruptedException {
            float totalSum = 0;
            int totalCount = 0;
            float globalMin = Float.MAX_VALUE;
            float globalMax = Float.MIN_VALUE;

            for (Text val : values) {
                String[] parts = val.toString().split(",");
                float sum = Float.parseFloat(parts[0]);
                int count = Integer.parseInt(parts[1]);
                float min = Float.parseFloat(parts[2]);
                float max = Float.parseFloat(parts[3]);

                totalSum += sum;
                totalCount += count;
                globalMin = Math.min(globalMin, min);
                globalMax = Math.max(globalMax, max);
            }

            float averageTemp = totalSum / totalCount;
            context.write(key, new Text("Average: " + averageTemp + ", Min: " + globalMin + ", Max: " + globalMax));
        }
    }

    public static void main(String[] args) throws Exception {
        Configuration conf = new Configuration();
        Job job = Job.getInstance(conf, "temperature stats");
        job.setJarByClass(TemperatureStats.class);
        job.setMapperClass(TemperatureMapper.class);
        job.setCombinerClass(TemperatureCombiner.class);
        job.setReducerClass(TemperatureReducer.class);
        job.setMapOutputKeyClass(Text.class);
        job.setMapOutputValueClass(FloatWritable.class);
        job.setOutputKeyClass(Text.class);
        job.setOutputValueClass(Text.class);
        job.setInputFormatClass(TextInputFormat.class);
        job.setOutputFormatClass(TextOutputFormat.class);
        TextInputFormat.addInputPath(job, new Path(args[0]));
        TextOutputFormat.setOutputPath(job, new Path(args[1]));
        System.exit(job.waitForCompletion(true) ? 0 : 1);
    }
}

内容的提问来源于stack exchange,提问作者Ashok Kumar

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 13:06:04