ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Hadoop MapReduce 过程中 Key 和 Value 分别存储什么值

Hadoop MapReduce 过程中 Key 和 Value 分别存储什么值 摘要本文以 WordCount 经典示例为基础详细解析 Hadoop MapReduce 过程中各个阶段 Key 和 Value 的具体含义与变化过程。通过图文结合的方式清晰展示从输入文件分割到最终输出结果的全流程数据流转。一、示例说明本文以 WordCount词频统计为例通过图解方式直观展示 MapReduce 各阶段 Key/Value 的变化过程。二、MapReduce 处理流程详解1. 输入分割阶段InputFormatInputFormat 将 HDFS 上要处理的文件逐行读入将文件拆分成 splits。由于测试文件较小每个文件为一个 split并将文件按行分割形成 key, value 对如图 4-1 所示。这一步由 MapReduce 框架自动完成其中偏移量即 key 值包括了回车所占的字符数Windows 和 Linux 环境会不同。Key/Value 含义Key行偏移量每行起始字符在文件中的位置Value该行的文本内容示例说明这里是把每个文件按行处理下图有两个文件每个文件有两行。每一行的开头字符所在位置的偏移量第一行的开头偏移量自然是 0hello 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 固定为 1Combine 阶段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 MapperLongWritable, 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 ReducerText, 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, IterableIntWritable values, Context context ) throws IOException, InterruptedException { // 3. Combine 阶段Map 端局部聚合 // 将相同 key单词的 value1累加 // 对应文中图 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 ReducerText, 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, IterableIntWritable 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()的输入参数KeyLongWritable类型表示行偏移量ValueText类型表示该行文本内容图 4-1框架自动将输入文件分割为 行偏移量, 行内容 对Map 处理阶段TokenizerMapper.map()方法Key从行偏移量变为单词Text类型Value从行内容变为固定值 1IntWritable类型图 4-2将每行文本按空格分割为每个单词输出 单词, 1Combine 阶段可选优化IntSumCombiner.reduce()方法Key单词保持不变Value从多个 1 累加为局部词频计数图 4-3在 Map 端对相同单词的计数进行局部累加减少网络传输Reduce 处理阶段IntSumReducer.reduce()方法Key单词保持不变Value从局部词频累加为全局最终词频图 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 变化并将理论知识与实际代码实现相结合。
返回列表