MapReduce案例之wordcount
·
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);
}
}
更多推荐



所有评论(0)