第三章:Hadoop MapReduce 分布式计算框架


文档摘要

第三章:Hadoop MapReduce 分布式计算框架 第三章:Hadoop MapReduce 分布式计算框架 3.1 MapReduce 框架概述 在大数据时代,数据量呈爆炸式增长,传统的单机计算模式已经无法满足海量数据的处理需求。Hadoop MapReduce 正是为了解决这个问题而诞生的。它借鉴了函数式编程的思想,将复杂的大规模数据处理任务分解成两个主要阶段:Map(映射) 和 Reduce(归约)。这两个阶段可以并行执行,从而实现高效的分布式计算。 核心思想: "分而治之,并行计算"。MapReduce 将输入数据分割成独立的数据块,由 Map 任务并行处理,生成中间结果;然后,Reduce 任务对 Map 阶段产生的中间结果进行合并和汇总,最终得到最终结果。

第三章:Hadoop MapReduce 分布式计算框架

第三章:Hadoop MapReduce 分布式计算框架

3.1 MapReduce 框架概述

在大数据时代,数据量呈爆炸式增长,传统的单机计算模式已经无法满足海量数据的处理需求。Hadoop MapReduce 正是为了解决这个问题而诞生的。它借鉴了函数式编程的思想,将复杂的大规模数据处理任务分解成两个主要阶段:Map(映射)Reduce(归约)。这两个阶段可以并行执行,从而实现高效的分布式计算。

核心思想: "分而治之,并行计算"。MapReduce 将输入数据分割成独立的数据块,由 Map 任务并行处理,生成中间结果;然后,Reduce 任务对 Map 阶段产生的中间结果进行合并和汇总,最终得到最终结果。

主要特点:

  • 简单易用: 开发者只需关注 Map 和 Reduce 函数的编写,框架负责任务调度、数据分发、容错处理等底层细节。

  • 高容错性: 框架能够自动检测和处理节点故障,保证任务的可靠执行。

  • 可扩展性: 可以通过增加集群节点来线性扩展计算能力,处理更大规模的数据。

  • 适用于大规模数据处理: 特别适合处理 PB 甚至 EB 级别的数据。

适用场景:

MapReduce 适用于离线批处理场景,例如:

  • 日志分析: 分析海量日志数据,提取关键指标、统计信息等。

  • 数据挖掘: 进行数据预处理、特征工程、模型训练等。

  • 搜索引擎索引构建: 构建大规模的倒排索引。

  • ETL(抽取、转换、加载): 从各种数据源抽取数据,进行清洗、转换,加载到目标数据仓库。

不适用场景:

MapReduce 不适合对延迟要求较高的在线处理场景,例如:

  • 实时计算: 对数据流进行实时处理和分析。

  • 交互式查询: 需要快速响应用户查询请求。

  • 迭代计算: 算法需要多次迭代才能收敛(例如:图算法、机器学习算法),MapReduce 的迭代效率相对较低。

3.2 MapReduce 框架架构

MapReduce 框架主要由以下几个核心组件构成:

  • 客户端(Client): 用户编写 MapReduce 程序并提交作业的入口。客户端负责作业的配置、提交、监控和结果获取。

  • 作业跟踪器(JobTracker)/资源管理器(ResourceManager): (在 Hadoop YARN 架构中,JobTracker 的功能被 ResourceManager 和 ApplicationMaster 取代,但为了兼容性和理解 MapReduce 的基本概念,这里先以 JobTracker 为例进行讲解,稍后会介绍 YARN 架构下的对应组件)。JobTracker 是 MapReduce 集群的中央协调器,负责接收客户端提交的作业,调度任务(Map Task 和 Reduce Task)到不同的 TaskTracker 上执行,监控任务执行状态,并处理任务失败等情况。

  • 任务跟踪器(TaskTracker)/节点管理器(NodeManager): (同样,在 YARN 架构中,TaskTracker 的功能被 NodeManager 取代)。TaskTracker 运行在集群的每个节点上,负责执行 JobTracker 分配的任务。它会启动和监控 Map Task 和 Reduce Task 的执行,并将任务的执行进度和状态汇报给 JobTracker。

  • HDFS(Hadoop Distributed File System): 分布式文件系统,用于存储输入数据、中间结果和最终输出结果。MapReduce 作业的数据通常存储在 HDFS 上。

架构图 (graph TD):

YARN 架构下的 MapReduce:

在 Hadoop YARN (Yet Another Resource Negotiator) 架构中,MapReduce 框架的组件有所变化,ResourceManager 和 NodeManager 取代了 JobTracker 和 TaskTracker,并且引入了 ApplicationMaster 的概念。

  • ResourceManager (RM): 负责集群资源的统一管理和调度。它接收来自 ApplicationMaster 的资源请求,并分配资源给应用程序。

  • NodeManager (NM): 运行在集群的每个节点上,负责节点资源的监控和管理,并执行 ResourceManager 分配的任务。

  • ApplicationMaster (AM): 每个 MapReduce 作业都有一个 ApplicationMaster 实例。它负责作业的生命周期管理,包括作业分解、任务调度、任务监控、容错处理等。ApplicationMaster 向 ResourceManager 申请资源,并与 NodeManager 协作执行任务。

YARN 架构下的 MapReduce 架构图 (graph TD):

YARN 架构更加通用和灵活,可以支持多种计算框架(例如 Spark, Flink 等)运行在 Hadoop 集群上,提高了资源利用率。

3.3 MapReduce 工作流程详解

MapReduce 作业的执行流程可以分为以下几个阶段:

  1. 作业提交(Job Submission):

    • 客户端将 MapReduce 程序(包括 Jar 包、配置文件等)提交给 ResourceManager (或 JobTracker)。

    • ResourceManager 接收到作业提交请求后,将作业信息(作业 Jar 包、输入路径、输出路径等)放入调度队列中。

  2. 作业初始化(Job Initialization):

    • ResourceManager 从调度队列中选择一个作业进行初始化。

    • ResourceManager 为该作业启动一个 ApplicationMaster (或 JobTracker)。

    • ApplicationMaster 负责作业的整体调度和管理。

    • ApplicationMaster 从 HDFS 中读取作业的输入数据信息(InputFormat),计算输入分片(InputSplit)。

  3. Map 阶段(Map Phase):

    • ApplicationMaster 根据输入分片信息,向 ResourceManager 申请资源,启动 Map Task。

    • Map Task 从 HDFS 读取对应的输入分片数据。

    • Map 函数 对输入数据进行处理,生成中间结果,以键值对 (key-value pair) 的形式输出。

    • 中间结果通常先写入本地磁盘,并进行分区(Partition)和排序(Sort)。

    • 分区(Partition): 根据 Reduce Task 的数量,将 Map 阶段的输出结果划分到不同的区。默认的分区器是 HashPartitioner,根据 key 的哈希值进行分区。

    • 排序(Sort): 对每个分区内的中间结果按照 key 进行排序,方便 Reduce 阶段的合并和归约。

  4. Shuffle 阶段(Shuffle Phase):

    • Shuffle 阶段是 MapReduce 的核心阶段,负责将 Map 阶段的输出结果传输到 Reduce 阶段。

    • Copy 阶段: Reduce Task 启动后,会主动从运行 Map Task 的节点上拉取属于自己分区的数据。

    • Merge 阶段: Reduce Task 将从不同 Map Task 节点拉取到的数据进行合并(Merge),并再次进行排序(Sort),确保相同 key 的数据聚集在一起。

  5. Reduce 阶段(Reduce Phase):

    • Reduce 函数 对 Shuffle 阶段合并排序后的数据进行处理,执行归约操作,生成最终结果。

    • Reduce 阶段的输出结果通常写入 HDFS。

  6. 作业完成(Job Completion):

    • 当所有 Reduce Task 执行完成后,ApplicationMaster 通知 ResourceManager 作业执行成功。

    • ResourceManager 向客户端返回作业执行结果。

数据流图 (graph TD):

详细流程步骤:

  1. InputFormat: InputFormat 组件负责将输入数据分割成 InputSplit,并提供 RecordReader 从 InputSplit 中读取数据,生成 <key, value> 对作为 Map 函数的输入。常见的 InputFormat 包括 TextInputFormat (处理文本文件), SequenceFileInputFormat (处理 SequenceFile 文件) 等。

  2. InputSplit: 逻辑分片,将输入数据划分为多个逻辑块,每个 Map Task 处理一个 InputSplit。InputSplit 的大小通常与 HDFS 的 block 大小一致,以提高数据本地性。

  3. Map Task: 每个 Map Task 负责处理一个 InputSplit 的数据。执行用户自定义的 Map 函数,将输入的 <key, value> 对转换为新的 <key, value> 对作为中间结果。

  4. Partitioner: Partitioner 组件负责将 Map 阶段的输出结果划分到不同的区,决定哪些 key 由哪个 Reduce Task 处理。默认的 HashPartitioner 根据 key 的哈希值对 Reduce Task 的数量取模来确定分区号。

  5. Sort & Spill: Map Task 将中间结果写入环形缓冲区,当缓冲区达到一定阈值时,溢写(Spill)到本地磁盘。在溢写过程中,会对数据进行排序。如果配置了 Combiner,则在溢写前会先执行 Combiner 函数,减少网络传输量。

  6. Combiner (可选): Combiner 是一个可选组件,本质上是一个本地的 Reduce 函数。它在 Map 阶段的输出结果写入磁盘之前,先对中间结果进行本地聚合,减少需要传输到 Reduce 阶段的数据量,提高效率。Combiner 的使用需要满足一定的条件,即 Combiner 的输出结果不会影响最终的 Reduce 结果。

  7. Shuffle (Copy & Merge & Sort): Reduce Task 启动后,会并发地从多个 Map Task 节点拉取属于自己分区的数据(Copy)。拉取到的数据会先存储在 Reduce Task 的内存缓冲区和磁盘上。当数据量达到一定阈值时,会进行合并(Merge)和排序(Sort)。

  8. Reduce Task: 每个 Reduce Task 负责处理一个或多个分区的数据。执行用户自定义的 Reduce 函数,对输入数据进行归约操作,生成最终结果。Reduce 函数的输入是经过 Shuffle 阶段合并排序后的 <key, Iterable> 对。

  9. OutputFormat: OutputFormat 组件负责将 Reduce 阶段的输出结果写入到指定的存储介质(例如 HDFS, 本地文件系统, 数据库等)。常见的 OutputFormat 包括 TextOutputFormat (输出文本文件), SequenceFileOutputFormat (输出 SequenceFile 文件) 等。

3.4 MapReduce 代码实践:WordCount

WordCount 是 MapReduce 的经典入门案例,用于统计文本文件中每个单词出现的次数。以下是一个使用 Java 编写的 WordCount MapReduce 程序的示例代码。

1. Mapper 类 (WordCountMapper.java):

import java.io.IOException; import org.apache.hadoop.io.IntWritable; import org.apache.hadoop.io.LongWritable; import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.Mapper; import java.util.StringTokenizer; public class WordCountMapper extends Mapper<LongWritable, Text, Text, IntWritable> { private final static IntWritable one = new IntWritable(1); private Text word = new Text(); @Override public void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { String line = value.toString(); StringTokenizer tokenizer = new StringTokenizer(line); while (tokenizer.hasMoreTokens()) { word.set(tokenizer.nextToken()); context.write(word, one); } } }

代码详解:

  • public class WordCountMapper extends Mapper<LongWritable, Text, Text, IntWritable>: 定义 Mapper 类,继承自 org.apache.hadoop.mapreduce.Mapper 类。

    • <LongWritable, Text, Text, IntWritable> 指定 Mapper 的输入和输出键值对类型。

      • LongWritable: 输入 key 类型,表示行偏移量 (通常用不到,可以忽略)。

      • Text: 输入 value 类型,表示一行文本数据。

      • Text: 输出 key 类型,表示单词。

      • IntWritable: 输出 value 类型,表示单词计数 (固定为 1)。

  • private final static IntWritable one = new IntWritable(1);: 定义静态常量 one,表示单词计数 1。IntWritable 是 Hadoop 提供的用于整数类型的 Writable 接口实现。

  • private Text word = new Text();: 定义 word 变量,用于存储单词。Text 是 Hadoop 提供的用于文本类型的 Writable 接口实现。

  • @Override public void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException: 重写 map() 方法,实现 Map 阶段的逻辑。

    • LongWritable key: 输入 key (行偏移量)。

    • Text value: 输入 value (一行文本数据)。

    • Context context: 上下文对象,用于输出中间结果。

    • String line = value.toString();: 将输入 value (Text 类型) 转换为 String 类型。

    • StringTokenizer tokenizer = new StringTokenizer(line);: 使用 StringTokenizer 将一行文本数据分割成单词。

    • while (tokenizer.hasMoreTokens()) { ... }: 循环遍历每个单词。

    • word.set(tokenizer.nextToken());: 将当前单词设置到 word 变量中。

    • context.write(word, one);: 将单词和计数 1 作为键值对输出。context.write() 方法将中间结果写入到框架,框架会负责后续的处理 (分区、排序、Shuffle 等)。

2. Reducer 类 (WordCountReducer.java):

import java.io.IOException; import org.apache.hadoop.io.IntWritable; import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.Reducer; import java.util.Iterator; public class WordCountReducer extends Reducer<Text, IntWritable, Text, IntWritable> { private IntWritable result = new IntWritable(); @Override public void reduce(Text key, Iterable<IntWritable> values, Context context) throws IOException, InterruptedException { int sum = 0; Iterator<IntWritable> iterator = values.iterator(); while (iterator.hasNext()) { sum += iterator.next().get(); } result.set(sum); context.write(key, result); } }

代码详解:

  • public class WordCountReducer extends Reducer<Text, IntWritable, Text, IntWritable>: 定义 Reducer 类,继承自 org.apache.hadoop.mapreduce.Reducer 类。

    • <Text, IntWritable, Text, IntWritable> 指定 Reducer 的输入和输出键值对类型。

      • Text: 输入 key 类型,表示单词 (与 Mapper 输出 key 类型一致)。

      • IntWritable: 输入 value 类型,表示单词计数 (与 Mapper 输出 value 类型一致)。

      • Text: 输出 key 类型,表示单词。

      • IntWritable: 输出 value 类型,表示单词总计数。

  • private IntWritable result = new IntWritable();: 定义 result 变量,用于存储单词的总计数。

  • @Override public void reduce(Text key, Iterable<IntWritable> values, Context context) throws IOException, InterruptedException: 重写 reduce() 方法,实现 Reduce 阶段的逻辑。

    • Text key: 输入 key (单词)。

    • Iterable<IntWritable> values: 输入 value 集合,表示相同单词的所有计数 (例如: [1, 1, 1, ...]),框架会将相同 key 的 values 聚合在一起传递给 Reduce 函数。

    • Context context: 上下文对象,用于输出最终结果。

    • int sum = 0;: 初始化单词总计数 sum 为 0。

    • Iterator<IntWritable> iterator = values.iterator();: 获取 value 集合的迭代器。

    • while (iterator.hasNext()) { sum += iterator.next().get(); }: 遍历 value 集合,累加每个计数到 sum 中。

    • result.set(sum);: 将总计数 sum 设置到 result 变量中。

    • context.write(key, result);: 将单词和总计数作为键值对输出。

3. Driver 类 (WordCountDriver.java):

import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.fs.Path; import org.apache.hadoop.io.IntWritable; import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.Job; import org.apache.hadoop.mapreduce.lib.input.FileInputFormat; import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat; public class WordCountDriver { public static void main(String[] args) throws Exception { if (args.length != 2) { System.err.println("Usage: WordCount <input path> <output path>"); System.exit(-1); } Configuration conf = new Configuration(); Job job = Job.getInstance(conf, "Word Count"); job.setJarByClass(WordCountDriver.class); job.setMapperClass(WordCountMapper.class); job.setReducerClass(WordCountReducer.class); 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); } }

代码详解:

  • public class WordCountDriver { ... }: 定义 Driver 类,作为 MapReduce 程序的入口。

  • public static void main(String[] args) throws Exception { ... }: 主函数,程序入口。

  • if (args.length != 2) { ... }: 检查命令行参数,需要输入路径和输出路径两个参数。

  • Configuration conf = new Configuration();: 创建 Hadoop 配置对象。

  • Job job = Job.getInstance(conf, "Word Count");: 创建 Job 对象,指定作业名称为 "Word Count"。

  • job.setJarByClass(WordCountDriver.class);: 设置 Job 的 Jar 包,框架会根据 Driver 类所在的 Jar 包来查找 Mapper 和 Reducer 类。

  • job.setMapperClass(WordCountMapper.class);: 设置 Mapper 类。

  • job.setReducerClass(WordCountReducer.class);: 设置 Reducer 类。

  • job.setOutputKeyClass(Text.class);: 设置输出 key 类型 (与 Reducer 输出 key 类型一致)。

  • job.setOutputValueClass(IntWritable.class);: 设置输出 value 类型 (与 Reducer 输出 value 类型一致)。

  • FileInputFormat.addInputPath(job, new Path(args[0]));: 设置输入路径,从命令行参数获取。

  • FileOutputFormat.setOutputPath(job, new Path(args[1]));: 设置输出路径,从命令行参数获取。注意:输出路径不能已存在,否则作业会失败。

  • System.exit(job.waitForCompletion(true) ? 0 : 1);: 提交 Job 并等待作业完成。

    • job.waitForCompletion(true): 提交作业并等待完成,true 表示打印作业进度信息。

    • System.exit(...): 根据作业执行结果退出程序,0 表示成功,1 表示失败。

编译和运行 WordCount 程序:

  1. 环境准备: 确保已安装 Hadoop 环境,并配置好环境变量。

  2. 打包: 将 WordCountMapper.java, WordCountReducer.java, WordCountDriver.java 编译成 class 文件,并打包成 Jar 文件 (例如: wordcount.jar)。

  3. 上传数据: 将输入文本文件 (例如: input.txt) 上传到 HDFS 的输入路径 (例如: /input)。

  4. 运行命令: 在 Hadoop 集群的节点上,使用以下命令运行 WordCount 程序:

    hadoop jar wordcount.jar WordCountDriver /input /output
    • hadoop jar wordcount.jar: 运行 Jar 包。

    • WordCountDriver: Driver 类的完整类名。

    • /input: HDFS 输入路径。

    • /output: HDFS 输出路径。

  5. 查看结果: 作业执行完成后,可以在 HDFS 的输出路径 (例如: /output) 下查看结果文件 (通常是 part-r-00000 文件),其中包含了单词计数结果。

示例输入 (input.txt):

Hello Hadoop Hello World Hadoop MapReduce

示例输出 (part-r-00000):

Hadoop 2 Hello 2 MapReduce 1 World 1

3.5 总结

Hadoop MapReduce 框架作为大数据处理的基石,为大规模数据并行计算提供了强大而易用的解决方案。本章详细介绍了 MapReduce 的架构、工作原理、核心概念,并通过 WordCount 案例进行了代码实践。理解 MapReduce 的运行机制对于深入学习 Hadoop 生态系统和进行大数据应用开发至关重要。尽管现在出现了更先进的计算框架 (如 Spark, Flink),但 MapReduce 的思想和原理仍然是大数据领域的重要基础。掌握 MapReduce,能够帮助我们更好地理解分布式计算的本质,并为学习和应用其他大数据技术打下坚实的基础。


作者与出处
原作者: 灏天文库
来源:灏天文库
整理: 灏天文库整理
由灏天文库平台收录,内容或由平台用户上传,仅供学习交流
发布者: 作者: 灏天文库 转发
评论区 (0)
U