摘要:本文以 WordCount 经典示例为基础,详细解析 Hadoop MapReduce 过程中各个阶段 Key 和 Value 的具体含义与变化过程。通过图文结合的方式,清晰展示从输入文件分割到最终输出结果的全流程数据流转。
一、示例说明
本文以 WordCount(词频统计)为例,通过图解方式直观展示 MapReduce 各阶段 Key/Value 的变化过程。
二、MapReduce 处理流程详解
1. 输入分割阶段(InputFormat)
InputFormat 将 HDFS 上要处理的文件逐行读入,将文件拆分成 splits。由于测试文件较小,每个文件为一个 split,并将文件按行分割形成 <key, value> 对,如图 4-1 所示。
这一步由 MapReduce 框架自动完成,其中偏移量(即 key 值)包括了回车所占的字符数(Windows 和 Linux 环境会不同)。
Key/Value 含义:
- Key:行偏移量(每行起始字符在文件中的位置)
- Value:该行的文本内容
示例说明:这里是把每个文件按行处理,下图有两个文件,每个文件有两行。每一行的开头字符所在位置的偏移量,第一行的开头偏移量自然是 0,"hello world" 共 10 个字符,加上中间的空格 11 个字符,回车再算一个,第二行的开头偏移量是 12。
图 4-1 分割过程
2. Map 处理阶段
将分割好的 <key, value> 对交给用户定义的 map 方法进行处理,生成新的 <key, value> 对,如图 4-2 所示。
这里是用户自定义的 map 处理程序,每一行的字符按空格分割,分割的每一个元素都记为 1,也就是 map 节点的所有 value 都是 1。
Key/Value 含义:
- Key:单词(分割后的每个元素)
- Value:计数 1(每个单词出现一次记为 1)
图 4-2 执行 map 方法
3. Map 端排序与 Combine 阶段
得到 map 方法输出的 <key, value> 对后,Mapper 会将它们按照 key 值进行排序,并执行 Combine 过程,将 key 相同的 value 值累加,得到 Mapper 的最终输出结果,如图 4-3 所示。
Key/Value 含义:
- Key:单词(保持不变)
- Value:局部累加后的词频计数
图 4-3 Map 端排序及 Combine 过程
4. Reduce 处理阶段
Reducer 先对从 Mapper 接收的数据进行排序,再交由用户自定义的 reduce 方法进行处理,得到新的 <key, value> 对,并作为 WordCount 的输出结果,如图 4-4 所示。
Key/Value 含义:
- Key:单词(最终统计的单词)
- Value:全局累加后的最终词频
图 4-4 Reduce 端排序及输出结果
三、总结
通过 WordCount 示例可以清晰地看到,在 MapReduce 过程中:
- 输入阶段:Key 为行偏移量,Value 为行内容
- Map 阶段:Key 转换为单词,Value 固定为 1
- Combine 阶段:Key 保持不变,Value 进行局部累加
- Reduce 阶段:Key 保持不变,Value 进行全局累加得到最终结果
这种 Key/Value 的设计模式是 MapReduce 编程模型的核心,理解各阶段 Key/Value 的含义对于编写高效的 MapReduce 程序至关重要。
五、实战代码示例
下面是一个完整的 Hadoop MapReduce WordCount 程序(Java 版本),代码中包含了详细的注释,明确指出每个阶段对应的代码位置,并与文中图解的关键步骤相对应。
import java.io.IOException; import java.util.StringTokenizer; import org.apache.hadoop.conf.Configuration; 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.output.FileOutputFormat; /** WordCount 示例程序 对应文中图解的各阶段 Key/Value 变化过程 */ public class WordCount { /** Mapper 类 对应文中 "2. Map 处理阶段" 图解 */ public static class TokenizerMapper extends Mapper<LongWritable, Text, Text, IntWritable> { private final static IntWritable one = new IntWritable(1); private Text word = new Text(); /** map 方法 - Map 阶段核心处理逻辑 @param key: 行偏移量(InputFormat 阶段生成的 key) @param value: 该行的文本内容(InputFormat 阶段生成的 value) @param context: MapReduce 上下文 */ public void map(LongWritable key, Text value, Context context ) throws IOException, InterruptedException { // 1. InputFormat 阶段(框架自动完成): // - key: 行偏移量(如文中示例的 0, 12 等) // - value: 该行文本内容(如 "hello world") // 对应文中图 4-1 的分割过程 // 2. Map 处理阶段(用户自定义逻辑): // 将每行文本按空格分割成单词,每个单词输出 <单词, 1> StringTokenizer itr = new StringTokenizer(value.toString()); while (itr.hasMoreTokens()) { word.set(itr.nextToken()); // 输出 <单词, 1>,对应文中图 4-2 的 Map 输出 // 此时 key 变为单词,value 固定为 1 context.write(word, one); } } } /** Combiner 类(可选优化) 对应文中 "3. Map 端排序与 Combine 阶段" 图解 注意:Combiner 本质是本地 Reducer,在 Map 端执行局部聚合 */ public static class IntSumCombiner extends Reducer<Text, IntWritable, Text, IntWritable> { private IntWritable result = new IntWritable(); /** reduce 方法(Combiner 使用) @param key: 单词(Map 输出的 key) @param values: 该单词对应的所有 1 的集合 @param context: MapReduce 上下文 */ public void reduce(Text key, Iterable<IntWritable> values, Context context ) throws IOException, InterruptedException { // 3. Combine 阶段(Map 端局部聚合): // 将相同 key(单词)的 value(1)累加 // 对应文中图 4-3 的 Combine 过程 int sum = 0; for (IntWritable val : values) { sum += val.get(); } result.set(sum); // 输出 <单词, 局部累加值> context.write(key, result); } } /** Reducer 类 对应文中 "4. Reduce 处理阶段" 图解 */ public static class IntSumReducer extends Reducer<Text, IntWritable, Text, IntWritable> { private IntWritable result = new IntWritable(); /** reduce 方法 - Reduce 阶段核心处理逻辑 @param key: 单词(经过 Shuffle 排序后的 key) @param values: 该单词对应的所有计数值(可能来自多个 Mapper) @param context: MapReduce 上下文 */ public void reduce(Text key, Iterable<IntWritable> values, Context context ) throws IOException, InterruptedException { // 4. Reduce 处理阶段: // a) Shuffle & Sort:框架自动对 Mapper 输出按键排序 // b) Reduce:对相同 key 的所有 value 进行全局累加 // 对应文中图 4-4 的 Reduce 输出 int sum = 0; for (IntWritable val : values) { sum += val.get(); } result.set(sum); // 输出最终结果 <单词, 总词频> context.write(key, result); } } /** 主函数 - 作业配置和提交 */ public static void main(String[] args) throws Exception { Configuration conf = new Configuration(); Job job = Job.getInstance(conf, "word count"); // 设置 Jar 包 job.setJarByClass(WordCount.class); // 设置 Mapper job.setMapperClass(TokenizerMapper.class); // 设置 Combiner(可选,但推荐使用以减少网络传输) job.setCombinerClass(IntSumCombiner.class); // 设置 Reducer job.setReducerClass(IntSumReducer.class); // 设置输出 key/value 类型 job.setOutputKeyClass(Text.class); job.setOutputValueClass(IntWritable.class); // 设置输入输出路径 FileInputFormat.addInputPath(job, new Path(args[0])); FileOutputFormat.setOutputPath(job, new Path(args[1])); // 提交作业并等待完成 System.exit(job.waitForCompletion(true) ? 0 : 1); } }代码关键点说明:
| 处理阶段 | 对应代码 | Key/Value 变化 | 对应图解 | 说明 |
|---|---|---|---|---|
| InputFormat 阶段 (框架自动完成) | FileInputFormat和Mapper.map()的输入参数 |
| 图 4-1 | 框架自动将输入文件分割为 <行偏移量, 行内容> 对 |
| Map 处理阶段 | TokenizerMapper.map()方法 |
| 图 4-2 | 将每行文本按空格分割,为每个单词输出 <单词, 1> |
| Combine 阶段 (可选优化) | IntSumCombiner.reduce()方法 |
| 图 4-3 | 在 Map 端对相同单词的计数进行局部累加,减少网络传输 |
| Reduce 处理阶段 | IntSumReducer.reduce()方法 |
| 图 4-4 | 对来自所有 Mapper 的相同单词计数进行全局累加,输出最终结果 |
运行说明:
- 将代码保存为
WordCount.java - 编译:
javac -cp $(hadoop classpath) WordCount.java - 打包:
jar -cvf wordcount.jar *.class - 运行:
hadoop jar wordcount.jar WordCount /input/path /output/path - 查看结果:
hdfs dfs -cat /output/path/part-r-00000
通过这个完整的代码示例,您可以更直观地理解文中图解的各阶段 Key/Value 变化,并将理论知识与实际代码实现相结合。