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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.24 20:18:20