MapReduce Java程序键值重复问题求助(Hadoop3.2.3+Java8)
MapReduce计数错误问题排查与解决
问题描述
我是MapReduce与Hadoop的新手,使用Hadoop 3.2.3和Java 8开发程序,需求是根据行内符号拆分数据,例如将"q1,a,q0,"转换为('a',"q1,a,q0,")的键值对。我的数据集共10条数据,其中5条对应键'a'、5条对应键'b',但实际运行后,'a'对应5条数据,'b'却对应10条,与预期不符。
数据集
A,q0,a,q1;A,q0,b,q0;A,q1,a,q1;A,q1,b,q2;A,q2,a,q1;A,q2,b,q0;B,s0,a,s0;B,s0,b,s1;B,s1,a,s1;B,s1,b,s0
Mapper类
import java.io.IOException; import org.apache.hadoop.io.ByteWritable; import org.apache.hadoop.io.LongWritable; import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.Mapper; public class MyMapper extends Mapper<LongWritable, Text, ByteWritable ,Text>{ private ByteWritable key1 = new ByteWritable(); private int count =0 ; private Text wordObject = new Text(); @Override public void map(LongWritable key, Text value, Context context)throws IOException, InterruptedException { String ftext = value.toString(); for (String line: ftext.split(";")) { wordObject = new Text(); if (line.split(",")[2].equals("b")) { key1.set((byte) 'b'); wordObject.set(line) ; context.write(key1,wordObject); continue ; } key1.set((byte) 'a'); wordObject.set(line) ; context.write(key1,wordObject); } } }
Reducer类(原错误代码)
import java.io.IOException; import org.apache.hadoop.io.ByteWritable; import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.Reducer; public class MyReducer extends Reducer<ByteWritable, Text, ByteWritable ,Text>{ private Integer count=0 ; @Override public void reduce(ByteWritable key, Iterable<Text> values, Context context) throws IOException, InterruptedException { for(Text val : values ) { count++ ; } Text symb = new Text(count.toString()) ; context.write(key , symb); } }
Driver类
import org.apache.hadoop.fs.Path; import org.apache.hadoop.io.ByteWritable; import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.lib.input.FileInputFormat; import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat; import org.apache.hadoop.mapreduce.Job; import org.apache.hadoop.conf.Configured; import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.util.Tool; import org.apache.hadoop.util.ToolRunner; public class MyDriver extends Configured implements Tool { public int run(String[] args) throws Exception { if (args.length != 2) { System.out.printf("Usage: %s [generic options] <inputdir> <outputdir>\n", getClass().getSimpleName()); return -1; } @SuppressWarnings("deprecation") Job job = new Job(getConf()); job.setJarByClass(MyDriver.class); job.setJobName("separation "); FileInputFormat.setInputPaths(job, new Path(args[0])); FileOutputFormat.setOutputPath(job, new Path(args[1])); job.setMapperClass(MyMapper.class); job.setReducerClass(MyReducer.class); job.setMapOutputKeyClass(ByteWritable.class); job.setMapOutputValueClass(Text.class); job.setOutputKeyClass(ByteWritable.class); job.setOutputValueClass(Text.class); boolean success = job.waitForCompletion(true); return success ? 0 : 1; } public static void main(String[] args) throws Exception { int exitCode = ToolRunner.run(new Configuration(), new MyDriver(), args); System.exit(exitCode); } }
问题原因
错误出在Reducer类的count变量上。该变量是类成员变量,而Hadoop会复用Reducer实例来处理不同的键。当处理键'a'时,count累加至5;处理键'b'时,count不会重置,而是从5继续累加5次,最终得到10,导致结果不符合预期。
解决方案
将count变量移至reduce方法内部,每次处理新的键时重新初始化为0,确保每个键的计数独立:
修改后的Reducer类
import java.io.IOException; import org.apache.hadoop.io.ByteWritable; import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.Reducer; public class MyReducer extends Reducer<ByteWritable, Text, ByteWritable ,Text>{ @Override public void reduce(ByteWritable key, Iterable<Text> values, Context context) throws IOException, InterruptedException { // 将count移至方法内,每次处理键时重置为0 Integer count = 0 ; for(Text val : values ) { count++ ; } Text symb = new Text(count.toString()) ; context.write(key , symb); } }
内容的提问来源于stack exchange,提问作者r.walid
相关产品推荐
相关产品推荐

