1.2 Hadoop核心概念


文档摘要

1.2 Hadoop核心概念 1.2 Hadoop核心概念 1.2.1 Hadoop分布式文件系统 (HDFS) Hadoop分布式文件系统(HDFS)是Hadoop的核心组件之一,它是一个设计用于在廉价硬件上运行的分布式文件系统。HDFS具有高容错性、高吞吐量和高扩展性等特点,使其成为存储和处理大规模数据集的理想选择。 HDFS架构 HDFS采用主/从(Master/Slave)架构,主要由以下几个核心组件构成: NameNode(NN): NameNode是HDFS的“大脑”,负责管理文件系统的命名空间和元数据。元数据包括文件和目录的结构、文件的块信息(块ID、块的位置等)以及访问权限等。

1.2 Hadoop核心概念

1.2 Hadoop核心概念

1.2.1 Hadoop分布式文件系统 (HDFS)

Hadoop分布式文件系统(HDFS)是Hadoop的核心组件之一,它是一个设计用于在廉价硬件上运行的分布式文件系统。HDFS具有高容错性、高吞吐量和高扩展性等特点,使其成为存储和处理大规模数据集的理想选择。

1. HDFS架构

HDFS采用主/从(Master/Slave)架构,主要由以下几个核心组件构成:

  • NameNode(NN): NameNode是HDFS的“大脑”,负责管理文件系统的命名空间和元数据。元数据包括文件和目录的结构、文件的块信息(块ID、块的位置等)以及访问权限等。NameNode维护着文件系统的目录树和所有文件的元数据信息,并将这些信息存储在内存中以实现快速访问。此外,NameNode还负责管理DataNode,接收DataNode的心跳报告,并向DataNode发送指令。

  • DataNode(DN): DataNode是HDFS的工作节点,负责存储实际的数据块。数据被分割成多个块(默认大小为128MB),并以多副本的形式存储在不同的DataNode上,以提高数据的可靠性和容错性。DataNode定期向NameNode发送心跳报告,汇报自身的状态和存储的块信息。NameNode接收到心跳报告后,可以了解DataNode的健康状况,并根据需要进行数据块的重新分配和复制。

  • Secondary NameNode(SNN): Secondary NameNode并非NameNode的备份,而是一个辅助NameNode的节点。它的主要作用是定期合并NameNode的编辑日志(EditLog)和文件系统镜像(FsImage),并将合并后的FsImage推送给NameNode,以减少NameNode重启时加载编辑日志的时间,并帮助NameNode进行冷备份。在新的Hadoop版本中,Secondary NameNode的角色逐渐被更强大的HA机制所取代,例如使用备NameNode(Standby NameNode)实现高可用性。

  • Client: Client是用户与HDFS交互的接口。用户可以通过Client向NameNode发送请求,例如读取、写入、创建或删除文件。Client会与NameNode和DataNode进行通信,完成数据的读写操作。

图 1.2.1 HDFS架构示意图

2. HDFS数据存储和复制

HDFS将文件分割成大小固定的数据块(Block),默认大小为128MB。每个数据块会存储多个副本,默认情况下是3个副本,分布在不同的DataNode上。这种数据复制机制保证了数据的高可靠性和容错性。即使部分DataNode发生故障,数据仍然可以从其他副本中恢复。

数据复制策略通常遵循以下原则:

  • 第一个副本: 通常放置在写入数据的Client所在的DataNode上,如果Client不在集群内部,则随机选择一个DataNode。

  • 第二个副本: 放置在与第一个副本不同的机架(Rack)上的DataNode上。机架感知策略可以提高容错性,防止整个机架故障导致数据丢失。

  • 第三个副本: 放置在与第二个副本相同机架但不同的DataNode上。

后续副本则会尽量均匀地分布在集群中,以平衡存储负载和提高数据读取的并行性。

3. HDFS读写流程

写流程 (Write)

  1. Client向NameNode发送写请求: Client通过FileSystem接口向NameNode发起文件写入请求。

  2. NameNode检查权限和空间: NameNode检查Client是否有权限写入,以及文件系统是否有足够的空间。如果满足条件,NameNode会返回允许写入的响应。

  3. Client请求DataNode列表: Client向NameNode请求用于存储数据块的DataNode列表。NameNode根据副本策略选择一组DataNode,并返回给Client。

  4. Client将数据写入DataNode管道: Client将数据分成多个数据块,并以管道的方式将数据块写入第一个DataNode,第一个DataNode再将数据块转发给第二个DataNode,以此类推,形成一个数据流管道。每个DataNode在接收到数据后,会将数据写入本地磁盘,并继续转发给下一个DataNode。

  5. DataNode发送确认信息: 当数据块成功写入所有DataNode后,DataNode会向管道中的前一个DataNode发送确认信息,最终确认信息会传递回Client。

  6. Client通知NameNode写入完成: Client在所有数据块写入完成后,通知NameNode数据写入完成。NameNode更新元数据信息。

图 1.2.2 HDFS写流程示意图

读流程 (Read)

  1. Client向NameNode发送读请求: Client通过FileSystem接口向NameNode发起文件读取请求。

  2. NameNode返回数据块位置信息: NameNode根据请求的文件名,查找元数据信息,返回包含数据块位置信息的DataNode列表。

  3. Client连接DataNode读取数据: Client选择就近的DataNode(通常是网络拓扑距离最近的DataNode)建立连接,并向DataNode发送读取数据块的请求。

  4. DataNode读取数据并返回给Client: DataNode从本地磁盘读取数据块,并将数据块返回给Client。

  5. Client接收数据并组装文件: Client接收到所有数据块后,将数据块组装成完整的文件。

图 1.2.3 HDFS读流程示意图

4. HDFS代码实践 (Java API)

以下代码示例展示了如何使用Java HDFS API进行基本的文件操作,例如创建目录、上传文件、下载文件和列出目录内容。

环境准备:

  • 确保已安装并配置好Hadoop环境,并且HDFS服务正在运行。

  • 在你的Java项目中引入Hadoop客户端依赖,例如Maven:

<dependency> <groupId>org.apache.hadoop</groupId> <artifactId>hadoop-client</artifactId> <version>你的Hadoop版本</version> </dependency>

代码示例:

import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.fs.FileSystem; import org.apache.hadoop.fs.Path; import org.apache.hadoop.io.IOUtils; import java.io.IOException; import java.io.InputStream; import java.io.OutputStream; public class HDFSExample { public static void main(String[] args) throws IOException { Configuration conf = new Configuration(); // 如果你的Hadoop集群配置了core-site.xml等配置文件,Hadoop会根据配置文件自动连接 // 如果没有配置,你需要手动指定NameNode的地址,例如: // conf.set("fs.defaultFS", "hdfs://localhost:9000"); FileSystem fs = FileSystem.get(conf); // 1. 创建目录 Path dirPath = new Path("/example_dir"); if (!fs.exists(dirPath)) { fs.mkdirs(dirPath); System.out.println("目录创建成功: " + dirPath); } else { System.out.println("目录已存在: " + dirPath); } // 2. 上传文件 Path localFilePath = new Path("local_file.txt"); // 本地文件路径 Path hdfsFilePath = new Path("/example_dir/uploaded_file.txt"); // HDFS目标路径 try (OutputStream os = fs.create(hdfsFilePath)) { InputStream is = HDFSExample.class.getClassLoader().getResourceAsStream("local_file.txt"); // 假设local_file.txt在resources目录下 IOUtils.copy(is, os); System.out.println("文件上传成功: " + localFilePath + " -> " + hdfsFilePath); } catch (Exception e) { System.err.println("文件上传失败: " + e.getMessage()); } // 3. 下载文件 Path downloadHdfsFilePath = new Path("/example_dir/uploaded_file.txt"); Path localDownloadPath = new Path("downloaded_file.txt"); try (InputStream is = fs.open(downloadHdfsFilePath); OutputStream os = fs.create(localDownloadPath)) { IOUtils.copy(is, os); System.out.println("文件下载成功: " + downloadHdfsFilePath + " -> " + localDownloadPath); } catch (Exception e) { System.err.println("文件下载失败: " + e.getMessage()); } // 4. 列出目录内容 Path listDirPath = new Path("/example_dir"); System.out.println("目录内容: " + listDirPath); for (org.apache.hadoop.fs.FileStatus fileStatus : fs.listStatus(listDirPath)) { System.out.println("- " + fileStatus.getPath()); } fs.close(); } }

local_file.txt (resources目录下,用于上传测试)

This is a local file for HDFS upload example.

代码详解:

  • Configuration conf = new Configuration();: 创建Hadoop配置对象,用于加载Hadoop配置文件,例如core-site.xml, hdfs-site.xml等。

  • FileSystem fs = FileSystem.get(conf);: 获取HDFS文件系统实例。FileSystem.get(conf)会根据配置信息连接到HDFS集群。

  • fs.mkdirs(dirPath);: 在HDFS上创建目录。

  • fs.create(hdfsFilePath);: 在HDFS上创建文件,并返回输出流OutputStream,用于写入数据。

  • fs.open(downloadHdfsFilePath);: 打开HDFS上的文件,并返回输入流InputStream,用于读取数据.

  • fs.listStatus(listDirPath);: 列出指定目录下的文件和子目录信息。

  • IOUtils.copy(is, os);: 使用Hadoop的IO工具类IOUtils高效地进行流的复制,实现文件上传和下载。

  • fs.close();: 关闭FileSystem连接,释放资源。

运行代码前:

  1. 确保 local_file.txt 文件放置在 src/main/resources 目录下 (或者根据代码调整路径)。

  2. 如果你的Hadoop集群没有配置默认的 fs.defaultFS,需要在代码中手动设置,例如 conf.set("fs.defaultFS", "hdfs://你的NameNode地址:端口"); 替换 你的NameNode地址:端口 为你实际的NameNode地址和端口。

  3. 编译并运行Java代码。

运行成功后,你将在HDFS的 /example_dir 目录下看到 uploaded_file.txt 文件,并且本地会生成 downloaded_file.txt 文件,内容与 local_file.txt 相同。同时,控制台会输出目录创建、文件上传、文件下载和目录列表的信息。

1.2.2 MapReduce计算模型

MapReduce是Hadoop的核心计算框架,它是一种用于处理大规模数据集的分布式计算模型。MapReduce将复杂的计算任务分解成两个阶段:Map阶段和Reduce阶段,这两个阶段可以并行执行,从而实现高效的数据处理。

1. MapReduce工作原理

MapReduce的核心思想是“分而治之”。它将输入数据分割成多个独立的数据块,分配给不同的计算节点(Mapper)并行处理,生成中间结果,然后将中间结果进行Shuffle和排序,最后由Reducer对排序后的中间结果进行汇总和处理,生成最终结果。

MapReduce的整个工作流程可以概括为以下几个阶段:

  1. Input: 输入数据,通常存储在HDFS上。

  2. Split: 将输入数据分割成多个InputSplit,每个InputSplit对应一个Mapper任务。InputSplit是MapReduce任务的最小输入单元。

  3. Map: Mapper阶段并行处理InputSplit,对每个InputSplit中的数据进行处理,生成键值对形式的中间结果 <key, value>

  4. Shuffle: Shuffle阶段是MapReduce的核心阶段,负责将Mapper输出的中间结果按照Key进行分区、排序和分组,并将相同Key的Value汇集到一起,发送给相应的Reducer。Shuffle阶段包括Partition、Sort和Spill等子阶段。

  5. Reduce: Reducer阶段接收Shuffle阶段处理后的数据,对相同Key的Value列表进行汇总和处理,生成最终结果。

  6. Output: 输出结果,通常存储回HDFS。

图 1.2.4 MapReduce工作流程示意图

2. MapReduce核心组件

  • InputFormat: InputFormat负责处理输入数据,包括数据分割成InputSplit和将InputSplit解析成Mapper可以处理的键值对 <key, value>。常用的InputFormat包括TextInputFormat(处理文本文件)、KeyValueTextInputFormat(处理键值对文件)等。

  • Mapper: Mapper是Map阶段的核心组件,负责对输入的键值对进行处理,并输出中间结果键值对 <intermediate_key, intermediate_value>。用户需要自定义Mapper类,实现 map() 方法,定义具体的Map逻辑。

  • Partitioner: Partitioner负责将Mapper输出的中间结果按照Key进行分区,决定哪些Key的中间结果发送给哪个Reducer。默认的Partitioner是HashPartitioner,根据Key的哈希值进行分区。用户可以自定义Partitioner实现更复杂的分区策略。

  • Shuffle & Sort: Shuffle和Sort是MapReduce的核心阶段,由框架自动完成,不需要用户干预。Shuffle阶段包括分区(Partition)、排序(Sort)、合并(Combine,可选)、分组(Group)等子阶段。排序保证了相同Key的中间结果会被发送到同一个Reducer,并且Key在Reducer端是有序的。

  • Reducer: Reducer是Reduce阶段的核心组件,负责对Shuffle阶段处理后的数据进行汇总和处理,生成最终结果键值对 <output_key, output_value>。用户需要自定义Reducer类,实现 reduce() 方法,定义具体的Reduce逻辑。

  • OutputFormat: OutputFormat负责处理输出结果,将Reducer输出的键值对写入到输出文件。常用的OutputFormat包括TextOutputFormat(输出文本文件)、SequenceFileOutputFormat(输出SequenceFile文件)等。

3. MapReduce代码实践 (Java API - WordCount)

经典的WordCount程序是学习MapReduce的入门示例。以下代码示例展示了如何使用Java MapReduce API实现WordCount程序,统计文本文件中每个单词出现的次数。

环境准备:

  • 确保已安装并配置好Hadoop环境,并且YARN服务正在运行 (MapReduce任务需要在YARN上运行)。

  • 在你的Java项目中引入Hadoop客户端依赖,例如Maven (与HDFS示例相同)。

代码示例:

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.Mapper; import org.apache.hadoop.mapreduce.Reducer; import org.apache.hadoop.mapreduce.lib.input.FileInputFormat; import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat; import java.io.IOException; import java.util.StringTokenizer; public class WordCount { // Mapper类 public static class TokenizerMapper extends Mapper<Object, Text, Text, IntWritable> { private final static IntWritable one = new IntWritable(1); private final Text word = new Text(); public void map(Object key, Text value, Context context) throws IOException, InterruptedException { StringTokenizer itr = new StringTokenizer(value.toString()); // 将每行文本按空格分割成单词 while (itr.hasMoreTokens()) { word.set(itr.nextToken()); // 设置单词 context.write(word, one); // 输出 <单词, 1> 键值对 } } } // Reducer类 public static class IntSumReducer extends Reducer<Text, IntWritable, Text, IntWritable> { private final IntWritable result = new IntWritable(); public void reduce(Text key, Iterable<IntWritable> values, Context context) throws IOException, InterruptedException { 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"); // 创建Job实例 job.setJarByClass(WordCount.class); // 设置Job的启动类 job.setMapperClass(TokenizerMapper.class); // 设置Mapper类 // job.setCombinerClass(IntSumReducer.class); // 可以设置Combiner,优化性能 (可选) job.setReducerClass(IntSumReducer.class); // 设置Reducer类 job.setOutputKeyClass(Text.class); // 设置输出Key类型 job.setOutputValueClass(IntWritable.class); // 设置输出Value类型 FileInputFormat.addInputPath(job, new Path(args[0])); // 设置输入文件路径 (从命令行参数获取) FileOutputFormat.setOutputPath(job, new Path(args[1])); // 设置输出目录路径 (从命令行参数获取) System.exit(job.waitForCompletion(true) ? 0 : 1); // 提交Job并等待完成 } }

输入文件 (input.txt,假设放在HDFS的 /input 目录下)

Hello Hadoop Hello World Hadoop MapReduce

编译和运行:

  1. 编译: 使用Maven等工具编译Java代码,生成jar包,例如 wordcount.jar

  2. 上传jar包和输入文件到HDFS: 将 wordcount.jar 上传到HDFS的某个目录,例如 /user/hadoop/wordcount.jar。将 input.txt 上传到HDFS的 /input 目录。

  3. 运行MapReduce任务: 使用 Hadoop 命令提交 MapReduce 任务,例如:

hadoop jar /user/hadoop/wordcount.jar WordCount /input /output
  • /user/hadoop/wordcount.jar 是 jar 包在 HDFS 上的路径。

  • WordCount 是主类名。

  • /input 是输入文件目录 (HDFS)。

  • /output 是输出目录 (HDFS),注意输出目录不能事先存在,Hadoop 会自动创建。

代码详解:

  • TokenizerMapper: Mapper 类,继承 Mapper<Object, Text, Text, IntWritable>

    • map(Object key, Text value, Context context) 方法:

      • 输入: key (行偏移量,这里不需要使用), value (一行文本)。

      • 输出: <Text word, IntWritable one>,例如 <"Hello", 1>, <"Hadoop", 1>, <"World", 1>, <"MapReduce", 1>

      • 使用 StringTokenizer 将每行文本按空格分割成单词。

      • 使用 context.write(word, one) 输出键值对。

  • IntSumReducer: Reducer 类,继承 Reducer<Text, IntWritable, Text, IntWritable>

    • reduce(Text key, Iterable<IntWritable> values, Context context) 方法:

      • 输入: key (单词), values (相同单词的计数器列表,例如 [1, 1, 1])。

      • 输出: <Text key, IntWritable result>,例如 <"Hadoop", 2>, <"Hello", 2>, <"MapReduce", 1>, <"World", 1>

      • 遍历 values 列表,累加计数器值。

      • 使用 context.write(key, result) 输出键值对。

  • main 方法:

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

    • Job job = Job.getInstance(conf, "word count");: 创建 Job 实例,并设置 Job 名称为 "word count"。

    • job.setJarByClass(WordCount.class);: 设置 Job 的启动类,Hadoop 可以根据这个类找到 jar 包。

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

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

    • job.setOutputKeyClass(Text.class);: 设置输出 Key 类型。

    • job.setOutputValueClass(IntWritable.class);: 设置输出 Value 类型。

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

    • FileOutputFormat.setOutputPath(job, new Path(args[1]));: 设置输出目录路径,从命令行参数 args[1] 获取。

    • System.exit(job.waitForCompletion(true) ? 0 : 1);: 提交 Job 并等待完成,根据 Job 执行结果决定程序退出状态。

运行结果:

任务运行完成后,在 HDFS 的 /output 目录下会生成输出文件,例如 part-r-00000,内容如下 (顺序可能不同):

Hadoop 2 Hello 2 MapReduce 1 World 1

1.2.3 资源管理框架 YARN (Yet Another Resource Negotiator)

YARN是Hadoop 2.0引入的资源管理框架,它将JobTracker的资源管理和任务调度功能分离出来,分别由ResourceManager和ApplicationMaster负责。YARN使得Hadoop能够支持更多类型的计算框架(例如Spark、Storm等),而不仅仅是MapReduce。

1. YARN架构

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

  • ResourceManager (RM): ResourceManager是YARN集群的资源管理器,负责整个集群的资源管理和调度。RM接收Client提交的应用程序,并为应用程序分配资源(Container)。RM由两个主要组件组成:

    • Scheduler: 调度器,负责资源的调度和分配。YARN支持多种调度器,例如FIFO Scheduler、Capacity Scheduler、Fair Scheduler等。

    • ApplicationsManager: 应用程序管理器,负责管理集群中运行的应用程序,包括应用程序的提交、启动、监控和重启等。

  • NodeManager (NM): NodeManager是YARN集群的工作节点,负责管理单个节点上的资源和Container。NM定期向ResourceManager汇报节点资源使用情况,并接收ResourceManager的命令,启动和管理Container。

  • ApplicationMaster (AM): ApplicationMaster是每个应用程序的管理者,负责应用程序的任务调度和监控。当一个应用程序提交到YARN集群后,ResourceManager会为该应用程序启动一个ApplicationMaster。ApplicationMaster向ResourceManager申请资源,并将任务分配给Container执行。对于MapReduce应用程序,ApplicationMaster就是MRAppMaster。

  • Container: Container是YARN集群中的资源分配单位,它封装了CPU、内存、磁盘、网络等资源。每个任务都运行在Container中。ResourceManager根据应用程序的资源需求,为应用程序分配Container。

图 1.2.5 YARN架构示意图

2. YARN工作流程

  1. Client提交应用程序: Client向ResourceManager提交应用程序,包括应用程序的jar包、配置文件和启动命令等。

  2. ResourceManager分配Container并启动ApplicationMaster: ResourceManager接收到应用程序提交请求后,从集群中选择一个NodeManager,为应用程序分配一个Container,并在该Container中启动ApplicationMaster。

  3. ApplicationMaster向ResourceManager注册: ApplicationMaster启动后,向ResourceManager注册自己,以便ResourceManager知道该应用程序的存在。

  4. ApplicationMaster申请资源: ApplicationMaster根据应用程序的任务需求,向ResourceManager申请资源(Container)。

  5. ResourceManager分配Container给ApplicationMaster: ResourceManager根据集群资源情况和调度策略,为ApplicationMaster分配Container。

  6. ApplicationMaster启动Task: ApplicationMaster在分配到的Container中启动Task(例如MapTask和ReduceTask)。

  7. Task运行并向ApplicationMaster汇报状态: Task在Container中运行,并定期向ApplicationMaster汇报任务状态、进度和资源使用情况。

  8. ApplicationMaster监控任务: ApplicationMaster监控所有Task的运行状态,如果Task运行失败,ApplicationMaster会尝试重新启动Task。

  9. 应用程序完成: 当所有Task都运行完成后,ApplicationMaster向ResourceManager注销自己,并释放资源。

YARN的引入使得Hadoop集群能够更有效地管理和利用资源,支持更多类型的计算框架,提高了Hadoop的灵活性和扩展性。

1.2.4 总结

本章节详细介绍了Hadoop的核心概念,包括HDFS、MapReduce和YARN。

  • HDFS 提供了高可靠、高吞吐量、高扩展性的分布式文件存储系统,是Hadoop生态系统的基石。

  • MapReduce 是一种用于处理大规模数据集的分布式计算模型,通过Map和Reduce两个阶段实现并行计算。

  • YARN 是Hadoop的资源管理框架,负责集群资源的统一管理和调度,支持多种计算框架的运行。

理解这些核心概念是深入学习和应用Hadoop的关键。通过代码实践和图文详解,希望读者能够对Hadoop的核心组件和工作原理有更清晰的认识,为后续学习Hadoop生态系统的其他组件和技术打下坚实的基础。


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