第五章:Apache Hadoop — 分布式存储与计算实践指南 核心摘要:本章系统讲解 Apache Hadoop 的三大核心组件——HDFS(分布式文件系统)、MapReduce(并行计算模型)与 YARN(资源调度框架)的代码级实践。涵盖环境配置、HDFS 文件操作、Java API 编程、WordCount 完整 MapReduce 实现、YARN 作业提交流程,以及生产级性能调优策略。内容面向大数据开发工程师与平台运维人员,强调可落地的工程实践与最佳实践。 Hadoop 架构概览 Apache Hadoop 是一个高容错、可扩展的开源框架,专为在商用硬件集群上可靠地存储和处理海量结构化与非结构化数据而设计。
核心摘要:本章系统讲解 Apache Hadoop 的三大核心组件——HDFS(分布式文件系统)、MapReduce(并行计算模型)与 YARN(资源调度框架)的代码级实践。涵盖环境配置、HDFS 文件操作、Java API 编程、WordCount 完整 MapReduce 实现、YARN 作业提交流程,以及生产级性能调优策略。内容面向大数据开发工程师与平台运维人员,强调可落地的工程实践与最佳实践。
Apache Hadoop 是一个高容错、可扩展的开源框架,专为在商用硬件集群上可靠地存储和处理海量结构化与非结构化数据而设计。其架构采用主从(Master-Slave)模式,核心由以下三层组成:
| 组件 | 定位 | 关键能力 |
|---|---|---|
| HDFS(Hadoop Distributed File System) | 分布式存储层 | 提供高吞吐量的数据访问,支持超大文件(TB/PB级)、多副本容错、流式数据读写 |
| MapReduce | 分布式计算层 | 基于“分而治之”思想的编程模型,将计算任务自动切分为 Map(映射)与 Reduce(规约)阶段,实现并行化处理 |
| YARN(Yet Another Resource Negotiator) | 资源管理层 | 解耦资源调度与计算框架,支持 MapReduce、Spark、Flink 等多种计算引擎共存,提升集群资源利用率 |
✅ 技术演进说明:YARN 的引入标志着 Hadoop 从单一 MapReduce 计算平台向通用大数据操作系统转型,是现代 Hadoop 生态的基石。
部署前需确保集群节点(NameNode、DataNode、ResourceManager、NodeManager)已安装 Hadoop(推荐 3.3.x 或 3.4.x 稳定版本),并完成基础网络与 Java 环境配置(JDK 11+)。核心配置文件位于 $HADOOP_HOME/etc/hadoop/ 目录,需在所有节点保持同步。
| 配置文件 | 作用 | 必配项示例 |
|---|---|---|
core-site.xml |
定义 Hadoop 全局参数 | <property><name>fs.defaultFS</name><value>hdfs://namenode:9000</value></property> |
hdfs-site.xml |
配置 HDFS 特定参数 | <property><name>dfs.replication</name><value>3</value></property>(副本数)<property><name>dfs.namenode.name.dir</name><value>/data/hadoop/namenode</value></property> |
mapred-site.xml |
指定 MapReduce 运行框架 | <property><name>mapreduce.framework.name</name><value>yarn</value></property> |
yarn-site.xml |
配置 YARN 资源管理行为 | <property><name>yarn.resourcemanager.hostname</name><value>rm-host</value></property><property><name>yarn.nodemanager.aux-services</name><value>mapreduce_shuffle</value></property> |
⚠️ 验证要点:配置完成后,执行
hdfs namenode -format初始化文件系统,并通过start-dfs.sh与start-yarn.sh启动服务;使用jps命令确认NameNode、DataNode、ResourceManager、NodeManager进程正常运行。
HDFS 提供类 Unix 的命令行接口(CLI)与 Java API 两种主流操作方式,适用于数据准备、调试与集成开发场景。
| 操作类型 | 命令示例 | 说明 |
|---|---|---|
| 上传文件 | hadoop fs -put ./local_data.txt /input/ |
将本地文件上传至 HDFS /input/ 目录 |
| 创建目录 | hadoop fs -mkdir -p /output/result |
递归创建多级目录 |
| 查看内容 | hadoop fs -cat /input/data.txt | head -n 5 |
流式读取并显示前 5 行(适合大文件) |
| 统计信息 | hadoop fs -du -h /input/ |
显示目录下各文件大小(人类可读格式) |
| 权限设置 | hadoop fs -chmod 755 /input/ |
修改 HDFS 目录权限(遵循 Linux 权限模型) |
以下代码展示安全、健壮的 HDFS 文件操作模式,包含异常处理与资源自动释放:
import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.fs.*; import org.apache.hadoop.io.IOUtils; import java.io.*; import java.net.URI; public class HdfsClient { private final FileSystem fs; public HdfsClient(String namenodeUri) throws IOException { Configuration conf = new Configuration(); conf.set("fs.defaultFS", namenodeUri); this.fs = FileSystem.get(URI.create(namenodeUri), conf); } // 上传文件(支持大文件流式传输) public void uploadFile(String localPath, String hdfsPath) throws IOException { Path src = new Path(localPath); Path dst = new Path(hdfsPath); fs.copyFromLocalFile(false, true, src, dst); // overwrite=true, removeSrc=true System.out.println("✅ 上传完成: " + localPath + " → " + hdfsPath); } // 下载文件并打印前 10 行 public void downloadAndPreview(String hdfsPath, int lines) throws IOException { FSDataInputStream in = null; BufferedReader reader = null; try { in = fs.open(new Path(hdfsPath)); reader = new BufferedReader(new InputStreamReader(in)); String line; int count = 0; while ((line = reader.readLine()) != null && count < lines) { System.out.println((count + 1) + ": " + line); count++; } } finally { IOUtils.closeStream(reader); IOUtils.closeStream(in); } } public void close() throws IOException { fs.close(); } }
💡 最佳实践:生产环境应使用
Configuration加载core-site.xml和hdfs-site.xml,避免硬编码;对FileSystem实例建议单例复用,减少连接开销。
WordCount 是 MapReduce 的“Hello World”,但其完整实现涵盖 Mapper、Reducer、Combiner、Driver 及作业配置,是理解数据分发、Shuffle、排序与聚合机制的关键入口。
WordCountMapper.java)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.io.IOException; 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 protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { String line = value.toString().toLowerCase().replaceAll("[^a-z\\s]", ""); StringTokenizer tokenizer = new StringTokenizer(line); while (tokenizer.hasMoreTokens()) { String token = tokenizer.nextToken().trim(); if (!token.isEmpty()) { word.set(token); context.write(word, one); } } } }
🔍 增强点:增加小写转换与标点清洗,提升词频统计准确性;跳过空字符串,避免无效输出。
WordCountReducer.java)import org.apache.hadoop.io.IntWritable; import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.Reducer; import java.io.IOException; public class WordCountReducer extends Reducer<Text, IntWritable, Text, IntWritable> { private IntWritable result = new IntWritable(); @Override protected 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); } }
// 复用 Reducer 逻辑,实现 Map 端局部聚合,减少网络传输 public class WordCountCombiner extends Reducer<Text, IntWritable, Text, IntWritable> { private IntWritable result = new IntWritable(); @Override protected 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); } }
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("❌ 用法: hadoop jar WordCount.jar <输入路径> <输出路径>"); System.exit(1); } Configuration conf = new Configuration(); Job job = Job.getInstance(conf, "WordCount"); job.setJarByClass(WordCountDriver.class); // 设置 Mapper 与 Reducer job.setMapperClass(WordCountMapper.class); job.setCombinerClass(WordCountCombiner.class); // 启用 Combiner 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])); // ✅ 强制覆盖输出目录(避免重复运行报错) FileOutputFormat.setCompressOutput(job, true); FileOutputFormat.setOutputCompressorClass(job, org.apache.hadoop.io.compress.GzipCodec.class); // 提交作业并阻塞等待 boolean success = job.waitForCompletion(true); System.exit(success ? 0 : 1); } }
# 1. 编译打包(假设已配置 Maven) mvn clean package -DskipTests # 2. 上传输入数据至 HDFS hadoop fs -mkdir -p /input hadoop fs -put ./sample.txt /input/ # 3. 提交 MapReduce 作业(指定 JAR、主类、输入输出路径) hadoop jar target/wordcount-1.0.jar WordCountDriver /input /output # 4. 查看结果 hadoop fs -cat /output/part-r-00000
📈 性能提示:对于 10GB+ 数据,建议设置
mapreduce.job.reduces=10(根据集群规模调整),并启用mapreduce.map.output.compress=true。
YARN 将集群资源抽象为 Container(容器),通过 ResourceManager(RM)统一分配,NodeManager(NM)负责本地资源监控与任务执行。MapReduce 作业提交即触发 YARN 调度流程。
yarn-site.xml 中设置 yarn.scheduler.minimum-allocation-mb(最小内存)与 yarn.nodemanager.resource.memory-mb(节点总内存)http://<resourcemanager-host>:8088 访问 YARN Web UI,实时监控 Application、Container、队列资源使用率default、etl、ml),实现资源隔离与 SLA 保障| 优化维度 | 具体措施 | 效果预期 |
|---|---|---|
| MapReduce 层 | • 合理设置 mapreduce.input.fileinputformat.split.minsize 避免小文件切片• 启用 mapreduce.map.output.compress + SnappyCodec• 使用 MultipleOutputs 处理多路输出 |
减少 Map 数量 20%+,Shuffle 传输量下降 35%+ |
| HDFS 层 | • 大文件块大小设为 128MB 或 256MB(dfs.blocksize)• 启用短路本地读( dfs.client.read.shortcircuit)• 对冷数据启用 Erasure Coding 替代三副本 |
存储成本降低 50%,本地读延迟 < 1ms |
| YARN 层 | • 启用 NodeLabel 实现计算资源与存储资源拓扑感知调度 • 配置 yarn.nodemanager.vmem-pmem-ratio=4.0 优化内存超卖• 使用 GPU 资源类型支持深度学习作业 |
资源利用率提升至 75%+,GPU 作业调度延迟 < 5s |
Apache Hadoop 作为大数据技术栈的奠基性框架,其 HDFS、MapReduce 与 YARN 三位一体的架构设计,为 PB 级数据的可靠存储、高效计算与弹性调度提供了坚实底座。本章通过可运行的代码示例,系统覆盖了从环境搭建、HDFS 操作、MapReduce 编程到 YARN 调度的完整链路,并融入生产环境调优经验。
下一步建议:
- 迁移至 Hadoop 3.x:启用 Erasure Coding、S3A 多线程、GPU 支持等新特性
- 混合计算架构:将 MapReduce 与 Spark/Flink 结合,MapReduce 处理 ETL,Spark 承担迭代计算
- 云原生集成:通过 Ozone(Hadoop 对象存储)替代 HDFS,或使用 Alluxio 构建统一数据编排层
- 安全增强:集成 Kerberos 认证、Ranger 权限控制、KMS 加密服务
掌握 Hadoop 不仅是理解大数据底层原理的钥匙,更是构建企业级数据湖、实时数仓与 AI 平台不可或缺的工程能力。