MapReduce案例之wordcount

1.案例的主要流程

在这里插入图片描述

step1.数据格式的准备

在这里插入图片描述
结果如下图所示:
在这里插入图片描述

2.step2 Mapper代码的编写

package mapReduce;

import org.apache.hadoop.io.LongWritable;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.mapreduce.Mapper;

import java.io.IOException;

/*
四个泛型的解释:
    KEYIN:K1的类型
    VALUEIN:V1的类型

    KEYOUT:K2的类型
    VALUEOUT:V2的类型
 */

public class wordcount_Mapper extends Mapper<LongWritable, Text,Text,LongWritable>{
    //map方法就是将K1和V1转为K2和V2
    /*
        参数:
        key     :K1   行偏移量
        value   :V1   每一行的文本数据
        context :    表示上下文对象
     */
    /*
        如何将K1V1转化为K2V2
        K1            V1
         0      hello,world,hadoop
         18    hive,sqoop,flume,hello
         40    kitty,tom,jerry,world
         59    hadoop
     -------------------------------
         K2            V2
         hello          1
         world          1
         hdfs           1
         hadoop         1
         hello          1

     */
    @Override
    protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException {
        Text text =new Text();
        LongWritable longWritable =new LongWritable();
        //1:将一行的文本数据进行拆分
        String [] split =value.toString().split(",");
        //2.便利数组,组装K2和V2
        for (String word : split) {
            //3.将K2和V2写入上下文
            text.set(word);
            longWritable.set(1);
            context.write(text,longWritable);
        }
    }
}

3.step3 Reduce代码的编写

package mapReduce;

import org.apache.hadoop.io.LongWritable;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.mapreduce.Reducer;

import java.io.IOException;

/*
四个泛型的解释:
    KEYIN:K2
    VALUEIN:V2的类型

    KEYOUT:K3的类型
    VALUEOUT:V3的类型
 */
public class wordcount_Reduce extends Reducer<Text, LongWritable,Text,LongWritable>{
    //reduce方法就是将新的K2和V2转为K3和V3,将K3和V3写入上下文
     /*
        参数:
        key     :新的K2
        value   :集合 新V2
        context :    表示上下文对象
     */
    /*
        如何将K2V2转化为K3V3
        新的  K2           V2
          hello           <1,1,1>
          world           <1,1>
          hadoop          <1>
     -------------------------------
         K3           V3
         hello          3
         world          2
         hadoop         1
     */
    @Override
    protected void reduce(Text key, Iterable<LongWritable> values, Context context) throws IOException, InterruptedException {
         long count =0 ;
        //1.遍历集合,将集合中的数字相加,得到v3
        for (LongWritable value : values) {
            count += value.get();
        }
        //2.将k3和v3 写入上下文中
        context.write(key,new LongWritable(count));
    }
}

4. step4 主类代码的编写

package mapReduce;

import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.conf.Configured;
import org.apache.hadoop.fs.Path;
import org.apache.hadoop.io.LongWritable;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.mapreduce.Job;
import org.apache.hadoop.mapreduce.lib.input.TextInputFormat;
import org.apache.hadoop.mapreduce.lib.output.TextOutputFormat;
import org.apache.hadoop.util.Tool;
import org.apache.hadoop.util.ToolRunner;

public class job_Main extends Configured implements Tool {

    //该方法用于指定一个job任务
    @Override
    public int run(String[] strings) throws Exception {
        //创建job任务对象
        Job job =Job.getInstance(super.getConf(),"wordcount");
        //配置job任务对象
        //第一步  指定文件的读取方式和读取路径
        job.setInputFormatClass(TextInputFormat.class);
        TextInputFormat.addInputPath(job,new Path("hdfs://hadoop1:8020/wordcount"));//直接写目录名就可,会自动读取路径下的文件

        //第二步  指定map阶段的处理方式和数据类型
        job.setMapperClass(wordcount_Mapper.class);
        //设置map阶段K2的类型
        job.setOutputKeyClass(Text.class);
        //设置map阶段V2的类型
        job.setOutputValueClass(LongWritable.class);

        //第3、4、5、6步 采用默认的方式

        //第七步 指定Reduce阶段的处理方式和数据类型
        job.setReducerClass(wordcount_Reduce.class);
        //设置reduce阶段K3的类型
        job.setOutputKeyClass(Text.class);
        //设置map阶段V3的类型
        job.setOutputValueClass(LongWritable.class);

        //第八步  设置输出类型
        job.setOutputFormatClass(TextOutputFormat.class);
        //设置输出路径
        TextOutputFormat.setOutputPath(job,new Path("hdfs://hadoop1:8020/wordcount_out"));

        //等待任务结束
        boolean b1 =job.waitForCompletion(true);

        return b1 ? 0:1;
    }

    public static void main(String[] args) throws Exception {
        final Configuration configuration = new Configuration();
        //启动job任务
        int run = ToolRunner.run(configuration, new job_Main(), args);
        System.exit(run);
    }
}

更多推荐