Hadoop MapReduce温度统计程序报错:类型不匹配问题解决
解决Hadoop MapReduce中"wrong value class: Text is not FloatWritable"错误
错误原因分析
出现这个错误的核心原因有两个:
Combiner与Reducer的类型不匹配
你把TemperatureReducer同时设置为Combiner和Reducer,但TemperatureReducer的输出是<Text, Text>类型,而Reducer的输入需要和Mapper输出一致,即<Text, FloatWritable>。Combiner的输出必须能直接作为Reducer的输入,直接用最终Reducer充当Combiner会导致类型冲突。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
相关产品推荐
相关产品推荐

