大数据开发升级之路 | WordCount入门程序详解
源代码文件
WordCountDriver.java(Driver驱动类)
driver类负责配置job并提交,比如指定Mapper/Reducer,设置kv类型,输入输出的地址等,同时程序也从这里进入。
public class WordCountDriver {
public static void main(String[] args) throws IOException, ClassNotFoundException, InterruptedException {
// 1 获取配置信息以及获取job对象
Configuration conf = new Configuration();
Job job = Job.getInstance(conf);
// 2 关联本Driver程序的jar
job.setJarByClass(WordCountDriver.class);
// 3 关联Mapper和Reducer的jar
job.setMapperClass(WordCountMapper.class);
job.setReducerClass(WordCountReducer.class);
// 4 设置Mapper输出的kv类型
job.setMapOutputKeyClass(Text.class);
job.setMapOutputValueClass(IntWritable.class);
// 5 设置最终输出kv类型
job.setOutputKeyClass(Text.class);
job.setOutputValueClass(IntWritable.class);
// 6 设置输入和输出路径
FileInputFormat.setInputPaths(job, new Path("E:\\2.txt"));
FileOutputFormat.setOutputPath(job, new Path("E:\\3.txt"));
// 7 提交job
boolean result = job.waitForCompletion(true);
System.exit(result ? 0 : 1);
}
}
WordCountMapper.java(Mapper映射)
Mapper映射类负责将读取的数据映射为<k, v>的键值对形式,WordCountMapper代码为逐行读取文本,将一行内容拆分成若干<字母, 1> 对。
public class WordCountMapper extends Mapper<LongWritable, Text, Text, IntWritable>{
Text k = new Text();
IntWritable v = new IntWritable(1);
@Override
protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException {
// 1 获取一行
String line = value.toString();
// 2 切割
String[] words = line.split("");
// 3 输出
for (String word : words) {
k.set(word);
context.write(k, v);
}
}
}
WordCountReducer.java(Reducer类)
Reducer类是将经过shuffle后的键值对<k, List<v>>进行归并计算,,在WordCountReducer代码中为对相同字母的所有 value 求和,得到最终次数。
public class WordCountReducer extends Reducer<Text, IntWritable, Text, IntWritable>{
int sum;
IntWritable v = new IntWritable();
@Override
protected void reduce(Text key, Iterable<IntWritable> values,Context context) throws IOException, InterruptedException {
// 1 累加求和
sum = 0;
for (IntWritable count : values) {
sum += count.get();
}
// 2 输出
v.set(sum);
context.write(key,v);
}
}
WordCount执行
MapReduce分为Map和Reduce两个阶段,中间还有框架自动完成的shuffle阶段,一般情况下当读取数据之后:
-> 程序会根据文件大小来进行切片,分配MapTask,然后进行Map的映射阶段
-> 随后进入Shuffle阶段,将Map阶段所生成的键值对<k, v>按照k来分区拉取,边拉取边进行归并排序(k)
-> 最后为Reduce阶段,分配ReduceTask进行数据合并计算操作

执行流程剖析
(1)InputFormat & RecordReader
protected void map(LongWritable key, Text value, Context context) throws IOException,InterruptedException
刚开始的时候我有一个疑问,代码没有任何的文件操作逻辑,代码自动就会捕获这三个参数并进行下一步的运算操作,后来知道因为MapReduce 读数据不是直接读文件,而是通过 InputFormat 封装的。
🔹 一、InputFormat(输入格式)
InputFormat的作用是:
-
切分数据 → 把 HDFS 文件分成若干逻辑片段(InputSplit),每个 Split 对应一个 MapTask。
-
提供 RecordReader → 告诉 MapTask 如何把 Split 转换成一对对
<key, value>记录。
👉 所以代码中的 key/value 是 InputFormat 定义的,不是自己传的。
🔹 二、RecordReader(记录读取器)
InputFormat 把 Split 交给 RecordReader,后者负责 从数据流中迭代读取一条条记录,交给 Mapper。
-
key = 当前记录在文件中的位置(通常是偏移量)。
-
value = 当前记录的内容(比如一行文本)。
👉 每读一条,就调用一次 map(key, value, context)。
🔹 三、TextInputFormat(默认实现)
WordCount 里没有修改 InputFormat,所以默认 TextInputFormat。
-
TextInputFormat切分规则:按行切分(每行是一条记录)。
-
RecordReader = LineRecordReader。记录读取器 = 行记录读取器 。
LineRecordReader 的逻辑:
-
key = 该行的 起始字节偏移量(LongWritable)。
-
value = 该行的内容(Text)。
假设现在有文件1.txt
a b a c
b a c
c a b
那么经过TextInputFormat的转换之后:
(0, "a b a c")
(8, "b a c")
(14, "c a b")
(2)Mapper 阶段
// 1 获取一行
String line = value.toString();
// 2 切割
String[] words = line.split("");
// 3 输出
for (String word : words) {
k.set(word);
context.write(k, v);
}
Mapper 将每行字符串拆分为字母,并像context写入 <字母, 1>,这个context随后会传递给Reducer:
("a", 1), ("b", 1), ("a", 1), ("c", 1)
("b", 1), ("a", 1), ("c", 1)
("c", 1), ("a", 1), ("b", 1)
(3)Shuffle 阶段(框架自动完成)
-
分区:相同 key 保证进入同一个 Reducer。
-
排序:在 Reducer 内部,key 会被排序。
-
合并:同一个 key 的所有 value 被聚合到一起。
("a", [1,1,1,1])
("b", [1,1,1])
("c", [1,1,1])
(4)Reducer 阶段
Reducer 遍历同一个 key 的所有 value,做求和得到result:
("a", 4)
("b", 3)
("c", 3)
(5)输出结果
结果写入 HDFS 的输出目录,我在这里定义的为3.txt,这里其实有个问题,我给的这个3.txt系统会建立一个名为3.txt的文件夹,我原来以为是直接写入文件。。。
文件夹内有四个文件,其中_SUCCESS文件没有内容,只是提示你程序运行成功了,真正的数据在part-r-00000文件中,上面的两个crc数据是用来做数据校验保证数据完整性的,通常不必理会。

更多推荐


所有评论(0)