Hadoop MapReduce WordPairsCount 结果不一致问题求助
Hadoop MapReduce词对统计结果不一致问题
运行Hadoop MapReduce词对统计程序时出现结果矛盾问题:
- 输入为包含100个维基百科单行文档(.lineperdoc格式)的目录
- 程序逻辑:Map阶段提取文本行内的连续词对,Reduce阶段对相同词对计数求和,仅输出求和值达到指定阈值(如sum>=500)的词对
- 异常表现:设置不同阈值(>=100、>=500、>=1000)时,输出结果存在逻辑矛盾;修改Mapper文本分割逻辑后,问题仍未解决
原程序代码
import java.io.IOException; import java.util.TreeMap; import javax.naming.Context; import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.conf.Configured; import org.apache.hadoop.fs.Path; import org.apache.hadoop.io.IntWritable; import org.apache.hadoop.io.LongWritable; 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.FileInputFormat; import org.apache.hadoop.mapreduce.lib.input.TextInputFormat; import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat; import org.apache.hadoop.mapreduce.lib.output.TextOutputFormat; import org.apache.hadoop.util.Tool; import org.apache.hadoop.util.ToolRunner; public class HadoopWordPairsb2 extends Configured implements Tool { public static class Map extends Mapper<LongWritable, Text, Text, IntWritable> { private final static IntWritable one = new IntWritable(1); private Text pair = new Text(); private Text lastWord = new Text(); @Override public void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { String[] splitLine = value.toString().split(" "); for (String w : splitLine) { if (lastWord.getLength() > 0) { // only consider words if (w.matches("[a-zA-Z]+")) { pair.set(lastWord + ":" + w); context.write(pair, one); } } lastWord.set(w); } } } public static class Reduce extends Reducer<Text, IntWritable, Text, IntWritable> { @Override public void reduce(Text key, Iterable<IntWritable> values, Context context) throws IOException, InterruptedException { Integer sum = 0; for (IntWritable value : values) sum += value.get(); // output only words that occured more than 500 times if (sum >= 500) { context.write(key, new IntWritable(sum)); } } } @Override public int run(String[] args) throws Exception { Job job = Job.getInstance(new Configuration(), "HadoopWordPairsb2"); job.setJarByClass(HadoopWordPairsb2.class); job.setOutputKeyClass(Text.class); job.setOutputValueClass(IntWritable.class); job.setMapperClass(Map.class); job.setCombinerClass(Reduce.class); job.setReducerClass(Reduce.class); job.setInputFormatClass(TextInputFormat.class); job.setOutputFormatClass(TextOutputFormat.class); FileInputFormat.setInputPaths(job, args[0]); FileOutputFormat.setOutputPath(job, new Path(args[1])); job.waitForCompletion(true); return 0; } public static void main(String[] args) throws Exception { int ret = ToolRunner.run(new Configuration(), new HadoopWordPairsb2(), args); System.exit(ret); } }
修改后的Mapper代码
@Override public void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { String line = value.toString().toLowerCase(); String[] tokens = line.split("[^a-z0-9_.]+"); for (String token : tokens) { if (lastWord.getLength() > 0) { pair.set(lastWord + ":" + token); context.write(pair, one); } lastWord.set(token); } }
结果截图
- 初始逻辑下sum>=100与sum>=500的输出对比:

- 修改Mapper后sum>=1000与sum>=500的输出对比:

问题原因及解决方法
核心问题:Combiner的误用
你将Reduce类同时设置为Combiner和Reducer,但Reduce逻辑中包含阈值过滤(仅输出sum>=阈值的词对),这会导致统计数据丢失:
- Combiner在Map端局部聚合时,会直接丢弃未达阈值的词对计数,这些数据无法传递到Reducer进行全局求和
- 不同阈值下,Combiner过滤的数据集不同,最终全局求和结果自然出现逻辑矛盾
举个例子:某词对在Map1局部计数300、Map2局部计数300,全局总和600。若阈值设为500,Combiner会在Map端过滤掉两个300(均未达阈值),Reducer收不到任何数据,最终不会输出该词对,但实际全局总和符合阈值要求,导致结果缺失。
解决步骤
拆分Combiner与Reducer逻辑:
- Combiner仅负责局部求和,不做阈值过滤
- 阈值过滤仅在Reducer中执行
代码修改:
- 创建独立的Combiner类,仅实现求和逻辑:
public static class Combine extends Reducer<Text, IntWritable, Text, IntWritable> { @Override public void reduce(Text key, Iterable<IntWritable> values, Context context) throws IOException, InterruptedException { int sum = 0; for (IntWritable value : values) { sum += value.get(); } context.write(key, new IntWritable(sum)); } } - 在run方法中替换CombinerClass:
job.setCombinerClass(Combine.class); - 保留原Reducer类的阈值过滤逻辑不变
- 创建独立的Combiner类,仅实现求和逻辑:
额外优化点:
- 原Mapper中
lastWord会跨文本行保留状态,导致错误的跨行词对(需求是单行内连续词对),需在map方法开头重置:@Override public void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { lastWord.clear(); // 处理每行前重置lastWord,避免跨行词对 // 原逻辑... } - 文本分割可能产生空字符串,需添加空值判断:
for (String token : tokens) { if (token.isEmpty()) continue; // 原逻辑... }
- 原Mapper中
内容的提问来源于stack exchange,提问作者ztsv-av
相关产品推荐
相关产品推荐

