本文系统梳理大数据技术栈的核心基础——Hadoop生态,涵盖大数据基础概念、Hadoop架构演进、Java反射机制(Hadoop底层依赖)、HDFS分布式文件系统、MapReduce计算框架五大模块,包含原理讲解、命令实操、Java API开发、实战案例全流程,适合大数据入门学习者系统复习。
第一章 大数据与Hadoop概述
1.1 大数据基础概念
1.1.1 大数据的定义与5V特征
大数据(Big Data) 指无法在一定时间范围内用常规软件工具进行捕捉、管理和处理的数据集合,是需要新处理模式才能具备更强决策力、洞察发现力和流程优化能力的海量、高增长率和多样化的信息资产。
业界通用5V特征描述大数据:
| 大量 | Volume | 采集、存储、管理和分析的数据规模庞大,且数据量持续高速增长 |
| 高速 | Velocity | 数据产生和增长速度快,对存储和处理的时效性要求高 |
| 多样 | Variety | 数据类型和数据源具备多样性,涵盖结构化、半结构化、非结构化多种格式 |
| 低价值密度 | Value | 海量数据中有价值信息的占比低,需要通过算法挖掘有效信息 |
| 真实性 | Veracity | 数据质量反映真实业务情况,是数据分析结果可信的基础 |
研究大数据的核心意义在于预测:数据是对过去和现在的归纳总结,本身不具备趋势性,但通过分析数据可以总结事物发展的客观规律,建立数据思维模型,从而对未来趋势进行预测和判断。
1.1.2 大数据的三类数据格式
大数据按照结构特性可分为三类:
结构化数据
采用标准化格式、具备明确定义结构的数据,存储和排列有固定规律,易于程序读取和处理。典型应用是关系型数据库中的表数据,存储需求包括高速读写、数据备份、数据共享、容灾等。
半结构化数据
不遵循严格的数据模型,没有固定表结构,但包含标签标记来分隔语义、对数据分层组织,也被称为自描述结构,通常以树或图的形式存储。
-
典型格式:XML、JSON
-
特点:每条记录的属性个数可以不固定,灵活性远高于结构化数据
-
常见场景:邮件系统、Web页面、资源库、档案系统等
非结构化数据
数据结构不规则、不完整,没有预定义的数据模型,无法用二维逻辑表表示。企业中80%以上的数据都是非结构化数据,也是大数据最主要的组成部分。
-
典型格式:办公文档、纯文本、图片、音频、视频
-
特点:格式多样、标准不统一,存储、检索、分析的技术门槛更高
-
常见场景:医疗影像、视频点播、监控系统、文件服务器、媒体资源管理等
1.2 Hadoop核心概述
Hadoop是一个用于处理海量数据的分布式框架,核心能力分为两部分:为海量数据提供可靠的分布式存储,以及为海量数据提供高效的并行计算处理。
1.2.1 Hadoop的优缺点
| 高扩展性 | 可通过新增服务器节点,横向扩展集群的存储和计算能力 |
| 高效率 | 基于并行计算思想,采用"移动计算而非移动数据"的设计,提升计算效率 |
| 低成本 | 可基于廉价的普通服务器组建集群,大幅降低海量数据处理的硬件成本 |
| 高可靠性 | 自动维护数据的多份副本,单节点故障不会导致数据丢失 |
| 高容错性 | 任务执行过程中节点宕机时,框架会自动将任务转移到其他节点重试,保障任务完成 |
| 不适合处理小文件 | 设计目标是处理大文件,大量小文件会占用NameNode大量内存,寻址开销远超读取开销 |
| 不支持实时计算 | 核心是离线计算引擎,无法保证毫秒级/秒级的低延迟结果返回 |
| 原生安全性较低 | 存储和网络传输层面默认缺乏数据加密,存在数据泄露风险,生产环境需额外配置安全方案 |
1.2.2 Hadoop三大核心组件
Hadoop的核心由三部分组成:
-
HDFS(Hadoop Distributed File System):分布式文件系统,负责海量数据的可靠存储
-
MapReduce:分布式计算编程框架,负责海量数据的并行计算处理
-
YARN(Yet Another Resource Negotiator):集群资源管理器与任务调度器,负责集群资源的统一管理和任务调度
广义上的Hadoop也指代整个Hadoop生态体系,包含Hive、HBase、Spark、Flink等一系列基于Hadoop的大数据开源组件。
1.2.3 Hadoop架构演进
Hadoop 1.x 架构
-
MapReduce同时承担资源管理和数据处理两大职责,负载较重,扩展性差
-
HDFS仅负责分布式文件存储
-
仅支持MapReduce一种计算框架
Hadoop 2.x 架构
-
将资源管理功能从MapReduce中剥离,由YARN统一负责集群资源管理和任务调度
-
MapReduce仅负责数据处理,负载大幅降低
-
YARN支持为多种计算框架(MapReduce、Spark、Flink等)提供资源管理,生态兼容性大幅提升
-
HDFS仍负责分布式文件存储
Hadoop 3.x 架构优化
在2.x的基础上对四大模块进行了性能优化与功能增强:
-
Hadoop Common:通用工具包优化,支持更多操作系统与硬件架构
-
MapReduce:任务优化,提升小任务执行效率
-
YARN:资源调度策略增强,支持更细粒度的资源隔离
-
HDFS:支持纠删码技术,降低存储开销;支持NameNode高可用,解决单点故障问题
第二章 Java反射机制(Hadoop底层基础)
Hadoop框架的底层大量依赖Java反射机制实现类的动态加载、对象的动态创建与方法调用,是理解Hadoop源码的前置基础。
2.1 反射机制概述
常规编程中,我们通过类创建对象;而反射是将这个过程反转:在程序运行时,通过对象获取其所属类的完整信息,或动态创建对象、调用方法。
2.1.1 反射的核心作用
Java反射机制的核心能力是动态性,具体包含4点:
运行时动态构造任意一个类的对象
运行时获取任意一个对象所属类的完整信息
运行时调用任意一个类的成员变量和成员方法
运行时修改任意一个对象的属性和方法访问权限
2.1.2 反射的优势
反射最大的优点是实现动态创建对象和动态编译,赋予程序极高的灵活性。在Java EE、大数据框架开发中,反射是实现配置化、插件化的核心技术。
2.2 Class类详解
JVM编译.java文件生成.class字节码文件,加载字节码时会在内存中生成一个对应的Class对象,这个对象包含了类的全部结构信息。反射操作的本质就是获取Class对象,再通过它操作类的成员。
2.2.1 Class类核心方法
| forName(String className) | 根据全限定类名获取对应的Class对象 |
| getConstructors() | 获取类中所有public修饰的构造方法 |
| getDeclaredFields() | 获取本类所有成员变量(含private/protected/default/public),不包含父类继承的字段 |
| getFields() | 获取所有public修饰的成员变量,包含从父类继承的字段 |
| getMethods() | 获取所有public修饰的成员方法,包含从父类继承的方法 |
| getMethod(String name, Class… parameterTypes) | 根据方法名和参数类型,获取指定的public方法 |
| getInterfaces() | 获取当前类实现的全部接口 |
| getClass() | 获取实例对象对应的Class对象 |
| getName() | 获取类的全限定名(包含包名) |
| getSuperclass() | 获取类的父类Class对象 |
| newInstance() | 通过无参构造创建该类的实例对象 |
| isArray() | 判断该Class对象是否代表数组类型 |
2.2.2 获取Class对象的三种方式
Class类没有公共构造方法,获取Class对象有三种标准方式:
全限定类名获取:Class.forName("全限定类名")
对象实例获取:对象名.getClass()
类名直接获取:类名.class
代码示例
先定义一个实体类:
package cn.edu.aust;
public class Person {
private String name;
private int age;
private double height;
public Person() {}
public Person(String name, int age, double height) {
this.name = name;
this.age = age;
this.height = height;
}
// getter、setter、toString方法省略
}
测试三种获取方式:
package cn.edu.aust;
public class PersonTest {
public static void main(String[] args) {
Class<?> c1 = null;
Class<?> c2 = null;
Class<?> c3 = null;
// 方式1:全限定类名
try {
c1 = Class.forName("cn.edu.aust.Person");
} catch (ClassNotFoundException e) {
e.printStackTrace();
}
// 方式2:对象实例
c2 = new Person().getClass();
// 方式3:类名直接获取
c3 = Person.class;
// 三种方式获取的是同一个Class对象
System.out.println(c1.getName());
System.out.println(c2.getName());
System.out.println(c3.getName());
}
}
2.2.3 三种方式的区别
-
类名.class:JVM仅将类加载入内存,不执行类的初始化,返回Class对象
-
Class.forName("类名"):加载类的同时执行静态初始化,返回Class对象
-
对象.getClass():返回实例对象运行时实际所属类的Class对象
2.3 反射创建对象
2.3.1 通过无参构造创建对象
调用Class对象的newInstance()方法,会调用类的无参构造方法创建实例。
⚠️ 注意:该方式要求目标类必须存在公共无参构造方法,否则会抛出异常。
package cn.edu.aust;
public class PersonInstanceTest {
public static void main(String[] args) {
Class<?> c = null;
try {
c = Class.forName("cn.edu.aust.Person");
} catch (ClassNotFoundException e) {
e.printStackTrace();
}
Person person = null;
try {
// 通过无参构造创建对象
person = (Person) c.newInstance();
} catch (Exception e) {
e.printStackTrace();
}
person.setName("张三");
person.setAge(20);
person.setHeight(175.0);
System.out.println(person);
}
}
2.3.2 通过有参构造创建对象
如果需要调用有参构造,需要先获取对应的Constructor对象,再通过它实例化对象。
操作步骤:
通过getConstructors()获取类的所有构造方法
定位到目标有参构造对应的Constructor对象
调用Constructor的newInstance()方法传入参数创建对象
Constructor类常用方法
| getModifiers() | 获取构造方法的权限修饰符(数字编码) |
| getName() | 获取构造方法名称 |
| getParameterTypes() | 获取构造方法的所有参数类型 |
| newInstance(Object… initargs) | 传入参数,通过该构造方法创建对象 |
代码示例:
package cn.edu.aust;
import java.lang.reflect.Constructor;
public class ConstructorTest {
public static void main(String[] args) {
Class<?> c = null;
try {
c = Class.forName("cn.edu.aust.Person");
} catch (ClassNotFoundException e) {
e.printStackTrace();
}
// 获取所有public构造方法
Constructor<?>[] cons = c.getConstructors();
for (Constructor<?> con : cons) {
System.out.println(con);
}
Person person = null;
try {
// 调用有参构造创建对象
person = (Person) cons[1].newInstance("张三", 20, 175.0);
} catch (Exception e) {
e.printStackTrace();
}
System.out.println(person);
}
}
💡 补充:getModifiers()返回的是数字编码,如需转为可读的权限关键字,可调用java.lang.reflect.Modifier.toString(修饰符数字)转换。
2.4 反射获取类的完整结构
通过反射可以获取类的全部结构信息,核心依赖java.lang.reflect包下的三个类:
-
Constructor:封装构造方法信息
-
Field:封装成员属性信息
-
Method:封装成员方法信息
2.4.1 获取实现的接口与父类
-
获取接口:getInterfaces(),返回Class数组
-
获取父类:getSuperclass(),返回父类Class对象
2.4.2 获取全部成员方法
调用getMethods()获取所有public方法(含父类继承),返回Method数组。
Method类常用方法
| getModifiers() | 获取方法的权限修饰符 |
| getName() | 获取方法名称 |
| getParameterTypes() | 获取方法的所有参数类型 |
| getReturnType() | 获取方法的返回值类型 |
| getExceptionTypes() | 获取方法抛出的所有异常类型 |
| invoke(Object obj, Object… args) | 调用目标对象的该方法,传入参数 |
代码示例:
package cn.edu.aust;
import java.lang.reflect.Method;
import java.lang.reflect.Modifier;
public class MethodTest {
public static void main(String[] args) {
Class<?> c = null;
try {
c = Class.forName("cn.edu.aust.Person");
} catch (ClassNotFoundException e) {
e.printStackTrace();
}
Method[] methods = c.getMethods();
for (Method method : methods) {
// 权限修饰符
System.out.print(Modifier.toString(method.getModifiers()) + " ");
// 返回值类型
System.out.print(method.getReturnType().getSimpleName() + " ");
// 方法名
System.out.print(method.getName() + "(");
// 参数列表
Class<?>[] params = method.getParameterTypes();
for (int i = 0; i < params.length; i++) {
System.out.print(params[i].getSimpleName() + " arg" + i);
if (i < params.length – 1) System.out.print(", ");
}
System.out.println(")");
}
}
}
运行结果会包含Person类自身的方法,以及从Object类继承的equals、hashCode、toString等方法。
2.4.3 获取全部成员属性
获取属性有两种方式:
-
getFields():获取所有public属性,包含父类继承的
-
getDeclaredFields():获取本类所有属性(含私有),不包含父类的
Field类常用方法
| getModifiers() | 获取属性的权限修饰符 |
| getName() | 获取属性名称 |
| getType() | 获取属性的类型 |
| setAccessible(boolean flag) | 设置属性是否可访问(暴力反射,用于访问私有属性) |
| get(Object obj) | 获取指定对象中该属性的值 |
| set(Object obj, Object value) | 设置指定对象中该属性的值 |
代码示例:
package cn.edu.aust;
import java.lang.reflect.Field;
import java.lang.reflect.Modifier;
public class FieldTest {
public static void main(String[] args) {
Class<?> c = null;
try {
c = Class.forName("cn.edu.aust.Person");
} catch (ClassNotFoundException e) {
e.printStackTrace();
}
// 获取本类所有属性(含私有)
Field[] fields = c.getDeclaredFields();
for (Field field : fields) {
System.out.print(Modifier.toString(field.getModifiers()) + " ");
System.out.print(field.getType().getSimpleName() + " ");
System.out.println(field.getName());
}
}
}
2.5 反射的进阶操作
2.5.1 反射调用普通方法
步骤:
通过getMethod()获取指定的Method对象
调用Method的invoke()方法,传入目标对象和参数,执行方法
package cn.edu.aust;
import java.lang.reflect.Method;
public class MethodInvokeTest {
public static void main(String[] args) {
Class<?> c = null;
try {
c = Class.forName("cn.edu.aust.Person");
// 创建对象
Object obj = c.newInstance();
// 获取setName方法,参数为String类型
Method setName = c.getMethod("setName", String.class);
// 调用方法:对象 + 参数
setName.invoke(obj, "李四");
// 获取getName方法,无参数
Method getName = c.getMethod("getName");
// 调用方法并获取返回值
String name = (String) getName.invoke(obj);
System.out.println("姓名:" + name);
} catch (Exception e) {
e.printStackTrace();
}
}
}
2.5.2 反射通用调用getter/setter
封装通用工具方法,通过属性名自动拼接get/set方法名,动态调用getter/setter:
package cn.edu.aust;
import java.lang.reflect.Method;
public class ReflectUtil {
// 将属性名首字母大写
private static String initStr(String old) {
return old.substring(0, 1).toUpperCase() + old.substring(1);
}
// 通用setter调用
public static void setter(Object obj, String attr, Object value, Class<?> type) {
try {
Method method = obj.getClass().getMethod("set" + initStr(attr), type);
method.invoke(obj, value);
} catch (Exception e) {
e.printStackTrace();
}
}
// 通用getter调用
public static void getter(Object obj, String attr) {
try {
Method method = obj.getClass().getMethod("get" + initStr(attr));
System.out.println(attr + " : " + method.invoke(obj));
} catch (Exception e) {
e.printStackTrace();
}
}
public static void main(String[] args) throws Exception {
Class<?> c = Class.forName("cn.edu.aust.Person");
Object obj = c.newInstance();
setter(obj, "name", "王五", String.class);
setter(obj, "age", 22, int.class);
getter(obj, "name");
getter(obj, "age");
}
}
2.5.3 暴力反射:直接操作私有属性
反射可以绕过Java的访问权限控制,直接修改私有属性(不推荐在业务代码中使用,框架底层常用):
package cn.edu.aust;
import java.lang.reflect.Field;
public class FieldAccessTest {
public static void main(String[] args) throws Exception {
Class<?> c = Class.forName("cn.edu.aust.Person");
Object obj = c.newInstance();
// 获取私有属性name
Field nameField = c.getDeclaredField("name");
// 开启访问权限(暴力反射)
nameField.setAccessible(true);
// 直接设置属性值
nameField.set(obj, "赵六");
// 直接获取属性值
System.out.println("姓名:" + nameField.get(obj));
}
}
⚠️ 注意:直接操作私有属性破坏了类的封装性,实际开发中应优先通过getter/setter操作属性,该方式仅用于框架底层开发。
第三章 HDFS分布式文件系统
3.1 HDFS核心概述
HDFS(Hadoop Distributed File System)是Hadoop的分布式文件系统,用于存储海量文件,通过目录树结构定位文件;由多台服务器联合组成集群,不同节点承担不同角色。
适用场景:一次写入、多次读取的离线存储场景;文件创建写入关闭后不支持随机修改,仅支持追加。
3.1.1 HDFS的优缺点
| 高容错性 | 数据自动保存多副本,副本丢失后自动恢复 |
| 适合大数据 | 支持GB、TB、PB级数据存储,可管理百万级以上文件 |
| 低成本 | 可构建在廉价普通服务器上,通过多副本保证可靠性 |
| 不适合低延迟访问 | 无法满足毫秒级的实时数据访问需求 |
| 不适合大量小文件 | 小文件会占用NameNode大量内存存储元数据,且寻址时间超过读取时间 |
| 不支持并发写入与随机修改 | 一个文件同一时间仅支持一个写入者;仅支持数据追加,不支持文件随机位置修改 |
3.1.2 HDFS组成架构
HDFS采用主从(Master/Slave)架构,包含四大核心角色:
NameNode(NN,主节点):集群的管理者
-
管理HDFS的名称空间(目录树)
-
配置文件副本策略
-
管理数据块(Block)的映射信息
-
处理客户端的读写请求
DataNode(DN,从节点):实际存储数据的工作节点
-
存储实际的数据块
-
执行数据块的读写操作
-
定期向NameNode汇报自身状态和块信息
Client(客户端)
-
文件上传时将文件切分为固定大小的Block,逐个上传
-
与NameNode交互,获取文件的位置信息
-
与DataNode交互,执行实际的数据读写
-
提供命令行、API等方式管理和访问HDFS
Secondary NameNode(2NN):NameNode的辅助节点
-
定期合并Fsimage(镜像文件)和Edits(编辑日志),推送给NameNode,分担NameNode压力
-
紧急情况下可辅助恢复NameNode,但不是NameNode的热备,无法直接替换故障的NameNode
3.1.3 HDFS文件块大小设计
HDFS中的文件在物理上被切分为固定大小的数据块(Block)存储:
-
Hadoop 2.x/3.x 默认块大小:128MB
-
Hadoop 1.x 默认块大小:64MB
-
可通过配置参数dfs.blocksize自定义调整
块大小设计原理:
寻址时间(查找目标Block的时间)约为10ms,业界最优标准是寻址时间为传输时间的1%,因此理想传输时间约为1s。
普通机械磁盘传输速率约为100MB/s,因此理论最优块大小约为100MB,Hadoop官方取整设定为128MB。
块大小并非越大越好:块过大,数据传输时间变长,并行度降低;块过小,寻址开销占比升高,NameNode元数据压力增大。
3.2 Hadoop Web UI 查看集群状态
Hadoop启动后提供Web可视化管理界面:
-
HDFS Web UI:默认端口 9870(Hadoop 3.x)/ 50070(Hadoop 2.x),查看文件系统、集群节点状态
-
YARN Web UI:默认端口 8088,查看任务运行状态、资源使用情况
前置配置:
systemctl stop firewalld
systemctl disable firewalld
本地电脑配置hosts映射(C:\\Windows\\System32\\drivers\\etc\\hosts),添加集群节点IP与主机名映射
浏览器访问:http://NameNode主机名:9870 查看HDFS,http://ResourceManager主机名:8088 查看YARN
3.3 快速入门:官方词频统计案例
使用Hadoop自带的示例Jar包快速体验MapReduce任务:
vi /home/test.txt
# 输入测试文本,例如:
# hello hadoop
# hello world
# hadoop bigdata
hdfs dfs -mkdir -p /wordcount/input
hdfs dfs -put /home/test.txt /wordcount/input
cd $HADOOP_HOME/share/hadoop/mapreduce
hadoop jar hadoop-mapreduce-examples-*.jar wordcount /wordcount/input /wordcount/output
hdfs dfs -cat /wordcount/output/part-r-00000
常见问题:虚拟内存不足报错
解决方案:修改yarn-site.xml,调大虚拟内存与物理内存的比值:
<property>
<name>yarn.nodemanager.vmem-pmem-ratio</name>
<value>4.1</value>
</property>
同时可在mapred-site.xml中调小Map/Reduce任务的内存占用:
<property>
<name>mapreduce.map.memory.mb</name>
<value>256</value>
</property>
<property>
<name>mapreduce.reduce.memory.mb</name>
<value>256</value>
</property>
3.4 HDFS Shell 常用命令
HDFS提供类Linux Shell的命令操作文件系统,语法格式:
hdfs dfs [命令选项] [参数]
3.4.1 目录操作
# 查看根目录
hdfs dfs -ls /
# 递归查看所有子目录
hdfs dfs -ls -R /
# 人性化显示文件大小
hdfs dfs -ls -h /wordcount
# 创建单层目录
hdfs dfs -mkdir /test
# 递归创建多级目录
hdfs dfs -mkdir -p /a/b/c
# 查看目录下每个文件的大小
hdfs dfs -du /wordcount
# 查看目录总大小
hdfs dfs -du -s -h /wordcount
3.4.2 文件操作
# 上传本地文件到HDFS指定目录
hdfs dfs -put 本地文件路径 HDFS目录路径
# 强制覆盖已存在的文件
hdfs dfs -put -f 本地文件路径 HDFS目录路径
# 下载HDFS文件到本地
hdfs dfs -get HDFS文件路径 本地目录路径
# 强制覆盖本地已存在的文件
hdfs dfs -get -f HDFS文件路径 本地目录路径
hdfs dfs -cat /wordcount/input/test.txt
# 移动文件到目标目录
hdfs dfs -mv /a.txt /test/
# 重命名文件
hdfs dfs -mv /a.txt /b.txt
hdfs dfs -cp /source/file.txt /target/
# 删除文件
hdfs dfs -rm /test.txt
# 递归删除目录及所有内容
hdfs dfs -rm -r /test
# 跳过回收站直接删除
hdfs dfs -rm -r -skipTrash /test
3.5 实战:Shell脚本定时采集日志到HDFS
生产环境中,服务器每天产生大量日志,通常通过定时脚本将日志周期性上传到HDFS存储。
3.5.1 编写采集脚本
创建uploadHDFS.sh脚本:
#!/bin/bash
# 配置Hadoop环境变量
export HADOOP_HOME=/usr/hadoop/hadoop-2.7.3
export PATH=$PATH:$HADOOP_HOME/bin:$HADOOP_HOME/sbin
# Hadoop日志目录
hadoop_log_dir=/usr/hadoop/hadoop-2.7.3/logs/
# 待上传日志的临时目录
log_toupload_dir=/usr/data/logs/toupload/
# 按时间生成HDFS存储目录
date=`date +%Y_%m_%d_%H_%M`
hdfs_dir=/hadoop_log/$date/
# 创建临时目录
if [ -d $log_toupload_dir ];then
echo "$log_toupload_dir exist"
else
mkdir -p $log_toupload_dir
fi
# 收集各节点日志到临时目录
ls $hadoop_log_dir | while read fileName
do
if [[ $fileName == *.log ]];then
echo "moving hadoop log to $log_toupload_dir"
cp $hadoop_log_dir/*.log $log_toupload_dir
scp root@hadoop2:$hadoop_log_dir/*.log $log_toupload_dir
scp root@hadoop3:$hadoop_log_dir/*.log $log_toupload_dir
break
fi
done
# 在HDFS创建存储目录
echo "create $hdfs_dir"
hdfs dfs -mkdir -p $hdfs_dir
# 上传日志文件到HDFS
ls $log_toupload_dir | while read fileName
do
echo "upload hadoop log $fileName to $hdfs_dir"
hdfs dfs -put ${log_toupload_dir}${fileName} $hdfs_dir
done
# 清理临时目录
rm -rf $log_toupload_dir
3.5.2 配置定时任务
使用Linux Crontab实现定时执行:
rpm -qa | grep crontab
# 未安装则执行
yum -y install vixie-cron crontabs
systemctl start crond.service
systemctl enable crond.service
chmod 777 /usr/data/uploadHDFS.sh
crontab -e
# 添加配置:每10分钟执行一次
*/10 * * * * /usr/data/uploadHDFS.sh
3.6 HDFS Java API 操作
HDFS提供Java API实现编程式操作文件系统,核心类位于org.apache.hadoop.fs包下。
3.6.1 核心API介绍
| FileSystem | HDFS文件系统核心类,提供文件的增删改查等所有操作 |
| FileStatus | 封装文件/目录的元数据(大小、块大小、副本数、修改时间等) |
| FSDataInputStream | HDFS输入流,用于读取HDFS文件 |
| FSDataOutputStream | HDFS输出流,用于写入数据到HDFS文件 |
| Path | 代表HDFS中的文件或目录路径 |
FileSystem核心方法
| copyFromLocalFile(Path src, Path dst) | 将本地文件上传到HDFS |
| copyToLocalFile(Path src, Path dst) | 将HDFS文件下载到本地 |
| mkdirs(Path f) | 创建目录,支持递归创建 |
| rename(Path src, Path dst) | 重命名文件/目录 |
| delete(Path f, boolean recursive) | 删除文件/目录,recursive为true时递归删除 |
| listFiles(Path f, boolean recursive) | 列出目录下的所有文件信息 |
3.6.2 环境准备
创建Maven项目,引入Hadoop依赖:
<dependencies>
<dependency>
<groupId>org.apache.hadoop</groupId>
<artifactId>hadoop-common</artifactId>
<version>3.3.0</version>
</dependency>
<dependency>
<groupId>org.apache.hadoop</groupId>
<artifactId>hadoop-hdfs</artifactId>
<version>3.3.0</version>
</dependency>
<dependency>
<groupId>org.apache.hadoop</groupId>
<artifactId>hadoop-client</artifactId>
<version>3.3.0</version>
</dependency>
<dependency>
<groupId>junit</groupId>
<artifactId>junit</artifactId>
<version>4.12</version>
<scope>test</scope>
</dependency>
</dependencies>
3.6.3 常用操作实操
package cn.edu.aust;
import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.fs.*;
import org.junit.Before;
import org.junit.Test;
import java.io.IOException;
import java.net.URI;
public class HDFSOperationTest {
private FileSystem fs;
@Before
public void init() throws Exception {
Configuration conf = new Configuration();
// 指定HDFS地址
URI uri = new URI("hdfs://hadoop1:9000");
// 指定操作的用户身份,避免权限报错
String user = "root";
fs = FileSystem.get(uri, conf, user);
}
// 1. 创建目录
@Test
public void testMkdir() throws IOException {
fs.mkdirs(new Path("/api/test"));
fs.close();
}
// 2. 上传文件
@Test
public void testPut() throws IOException {
fs.copyFromLocalFile(
new Path("D:\\\\test\\\\input.txt"),
new Path("/api/test/")
);
fs.close();
}
// 3. 下载文件
@Test
public void testGet() throws IOException {
// 参数:是否删除源文件、源路径、目标路径、是否使用本地校验
fs.copyToLocalFile(
false,
new Path("/api/test/input.txt"),
new Path("D:\\\\test\\\\download"),
true
);
fs.close();
}
// 4. 重命名
@Test
public void testRename() throws IOException {
fs.rename(
new Path("/api/test/input.txt"),
new Path("/api/test/word.txt")
);
fs.close();
}
// 5. 删除
@Test
public void testDelete() throws IOException {
// 递归删除目录
fs.delete(new Path("/api"), true);
fs.close();
}
// 6. 查看文件详情
@Test
public void testListFiles() throws IOException {
RemoteIterator<LocatedFileStatus> files = fs.listFiles(new Path("/"), true);
while (files.hasNext()) {
LocatedFileStatus file = files.next();
System.out.println("文件名:" + file.getPath().getName());
System.out.println("文件大小:" + file.getLen() + "字节");
System.out.println("副本数:" + file.getReplication());
System.out.println("权限:" + file.getPermission());
// 获取块的位置信息
BlockLocation[] blocks = file.getBlockLocations();
for (BlockLocation block : blocks) {
String[] hosts = block.getHosts();
System.out.print("块所在节点:");
for (String host : hosts) {
System.out.print(host + " ");
}
System.out.println();
}
System.out.println("————————");
}
fs.close();
}
// 7. 判断是文件还是目录
@Test
public void testIsFile() throws IOException {
FileStatus[] statuses = fs.listStatus(new Path("/"));
for (FileStatus status : statuses) {
if (status.isFile()) {
System.out.println("文件:" + status.getPath().getName());
} else {
System.out.println("目录:" + status.getPath().getName());
}
}
fs.close();
}
}
3.6.4 配置参数优先级
HDFS配置参数的优先级从高到低为:
客户端代码中通过configuration.set()设置的值
项目ClassPath下的自定义配置文件(如hdfs-site.xml)
服务器端的自定义配置文件(xxx-site.xml)
服务器端的默认配置文件(xxx-default.xml)
3.7 HDFS核心原理:读写流程
3.7.1 HDFS写文件流程
客户端向NameNode发起写文件请求
NameNode校验权限、路径是否存在,通过后返回可写入的DataNode列表
客户端将文件切分为Block,按顺序逐个写入
客户端通过Pipeline管道机制,将数据块依次写入多个DataNode副本
每个DataNode写入完成后向上游返回确认
所有副本写入完成后,客户端通知NameNode关闭文件
3.7.2 HDFS读文件流程
客户端向NameNode发起读文件请求
NameNode返回文件的Block位置信息,按距离客户端的远近排序
客户端就近选择DataNode读取第一个Block
读取完成后校验数据完整性,继续读取下一个Block
所有Block读取完成后,客户端关闭流
第四章 MapReduce分布式计算框架
MapReduce是一个分布式运算程序的编程框架,是Hadoop核心的计算层。它将用户编写的业务逻辑与框架默认组件整合,自动分发到集群上并行运行,开发者无需关心底层分布式细节,只需专注业务逻辑。
4.1 核心基础
4.1.1 MapReduce运行架构
一个完整的MapReduce程序运行时有三类进程:
MrAppMaster:负责整个程序的调度和状态协调,一个Job对应一个AppMaster
MapTask:负责Map阶段的数据处理,并行执行多个
ReduceTask:负责Reduce阶段的数据汇总处理,并行执行多个
4.1.2 Hadoop序列化机制
序列化:将内存中的Java对象转换为字节序列,用于持久化存储或网络传输。
反序列化:将字节序列恢复为内存中的Java对象。
Java原生的Serializable是重量级序列化,会附带大量额外信息(校验头、继承体系等),网络传输效率低。因此Hadoop自研了Writable序列化机制,特点是:
-
紧凑:存储空间利用率高
-
快速:读写额外开销小
-
互操作:支持多语言交互
Hadoop常用序列化类型
| Boolean | BooleanWritable |
| Byte | ByteWritable |
| Int | IntWritable |
| Float | FloatWritable |
| Long | LongWritable |
| Double | DoubleWritable |
| String | Text |
| Map | MapWritable |
| 数组 | ArrayWritable |
| null | NullWritable |
4.2 MapReduce编程规范
用户编写的代码分为三部分:Mapper、Reducer、Driver。
4.2.1 Mapper阶段
-
自定义Mapper类继承Mapper<KEYIN, VALUEIN, KEYOUT, VALUEOUT>
-
输入数据是键值对(KV)形式
-
业务逻辑写在map()方法中
-
输出数据也是键值对形式
-
map()方法对每一个输入KV调用一次
4.2.2 Reducer阶段
-
自定义Reducer类继承Reducer<KEYIN, VALUEIN, KEYOUT, VALUEOUT>
-
输入KV类型对应Mapper的输出KV类型
-
业务逻辑写在reduce()方法中
-
reduce()方法对每一组相同Key的KV调用一次
4.2.3 Driver驱动类
相当于YARN的客户端,负责封装Job的运行参数,将整个程序提交到YARN集群运行。
4.3 入门案例:WordCount词频统计
统计输入文件中每个单词出现的次数。
4.3.1 项目依赖与日志配置
pom.xml依赖:
<dependencies>
<dependency>
<groupId>org.apache.hadoop</groupId>
<artifactId>hadoop-client</artifactId>
<version>3.3.0</version>
</dependency>
<dependency>
<groupId>junit</groupId>
<artifactId>junit</artifactId>
<version>4.12</version>
</dependency>
<dependency>
<groupId>org.slf4j</groupId>
<artifactId>slf4j-log4j12</artifactId>
<version>1.7.30</version>
</dependency>
</dependencies>
resources目录下log4j.properties:
log4j.rootLogger=INFO, stdout
log4j.appender.stdout=org.apache.log4j.ConsoleAppender
log4j.appender.stdout.layout=org.apache.log4j.PatternLayout
log4j.appender.stdout.layout.ConversionPattern=%d %p [%c] – %m%n
4.3.2 Mapper类实现
package cn.edu.aust.wordcount;
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;
/**
* 输入KEY:文件字节偏移量 LongWritable
* 输入VALUE:一行文本 Text
* 输出KEY:单词 Text
* 输出VALUE:次数1 IntWritable
*/
public class WordCountMapper extends Mapper<LongWritable, Text, Text, IntWritable> {
private final Text outK = new Text();
private final IntWritable outV = new IntWritable(1);
@Override
protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException {
// 1. 获取一行数据
String line = value.toString();
// 2. 按空格切分单词
String[] words = line.split(" ");
// 3. 逐个输出 <单词, 1>
for (String word : words) {
outK.set(word);
context.write(outK, outV);
}
}
}
4.3.3 Reducer类实现
package cn.edu.aust.wordcount;
import org.apache.hadoop.io.IntWritable;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.mapreduce.Reducer;
import java.io.IOException;
/**
* 输入KEY:单词 Text
* 输入VALUE:次数集合 Iterable<IntWritable>
* 输出KEY:单词 Text
* 输出VALUE:总次数 IntWritable
*/
public class WordCountReducer extends Reducer<Text, IntWritable, Text, IntWritable> {
private final IntWritable outV = new IntWritable();
@Override
protected void reduce(Text key, Iterable<IntWritable> values, Context context) throws IOException, InterruptedException {
int sum = 0;
// 累加同一单词的所有次数
for (IntWritable count : values) {
sum += count.get();
}
outV.set(sum);
context.write(key, outV);
}
}
4.3.4 Driver驱动类
package cn.edu.aust.wordcount;
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;
import java.io.IOException;
public class WordCountDriver {
public static void main(String[] args) throws IOException, InterruptedException, ClassNotFoundException {
// 1. 获取配置和Job对象
Configuration conf = new Configuration();
Job job = Job.getInstance(conf);
// 2. 关联Driver类(找Jar包)
job.setJarByClass(WordCountDriver.class);
// 3. 关联Mapper和Reducer
job.setMapperClass(WordCountMapper.class);
job.setReducerClass(WordCountReducer.class);
// 4. 设置Mapper输出KV类型
job.setMapOutputKeyClass(Text.class);
job.setMapOutputValueClass(IntWritable.class);
// 5. 设置最终输出KV类型
job.setOutputKeyClass(Text.class);
job.setOutputValueClass(IntWritable.class);
// 6. 设置输入输出路径
FileInputFormat.setInputPaths(job, new Path(args[0]));
FileOutputFormat.setOutputPath(job, new Path(args[1]));
// 7. 提交Job,等待运行结束
boolean result = job.waitForCompletion(true);
System.exit(result ? 0 : 1);
}
}
4.3.5 打包运行
在pom.xml中添加打包插件,打包含依赖的Jar包:
<build>
<plugins>
<plugin>
<artifactId>maven-compiler-plugin</artifactId>
<version>3.8.1</version>
<configuration>
<source>1.8</source>
<target>1.8</target>
</configuration>
</plugin>
<plugin>
<artifactId>maven-assembly-plugin</artifactId>
<configuration>
<descriptorRefs>
<descriptorRef>jar-with-dependencies</descriptorRef>
</descriptorRefs>
</configuration>
<executions>
<execution>
<id>make-assembly</id>
<phase>package</phase>
<goals>
<goal>single</goal>
</goals>
</execution>
</executions>
</plugin>
</plugins>
</build>
打包后上传到集群,通过hadoop jar命令提交运行。
4.4 进阶案例:流量统计(自定义Bean序列化)
统计每个手机号的上行流量、下行流量、总流量。
4.4.1 自定义FlowBean实现Writable接口
自定义Bean作为Value传输时,必须实现Writable接口:
package cn.edu.aust.flow;
import org.apache.hadoop.io.Writable;
import java.io.DataInput;
import java.io.DataOutput;
import java.io.IOException;
public class FlowBean implements Writable {
private long upFlow; // 上行流量
private long downFlow; // 下行流量
private long sumFlow; // 总流量
// 必须有空参构造,反序列化时反射调用
public FlowBean() {}
public FlowBean(long upFlow, long downFlow) {
this.upFlow = upFlow;
this.downFlow = downFlow;
this.sumFlow = upFlow + downFlow;
}
// 序列化方法:顺序写
@Override
public void write(DataOutput out) throws IOException {
out.writeLong(upFlow);
out.writeLong(downFlow);
out.writeLong(sumFlow);
}
// 反序列化方法:顺序必须和序列化完全一致
@Override
public void readFields(DataInput in) throws IOException {
upFlow = in.readLong();
downFlow = in.readLong();
sumFlow = in.readLong();
}
// 重写toString,方便输出结果
@Override
public String toString() {
return upFlow + "\\t" + downFlow + "\\t" + sumFlow;
}
// getter、setter省略
public long getUpFlow() { return upFlow; }
public void setUpFlow(long upFlow) { this.upFlow = upFlow; }
public long getDownFlow() { return downFlow; }
public void setDownFlow(long downFlow) { this.downFlow = downFlow; }
public long getSumFlow() { return sumFlow; }
public void setSumFlow(long sumFlow) { this.sumFlow = sumFlow; }
}
⚠️ 注意:反序列化的字段顺序必须和序列化完全一致,否则会出现数据错乱。
4.4.2 Mapper与Reducer实现
FlowMapper
package cn.edu.aust.flow;
import org.apache.hadoop.io.LongWritable;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.mapreduce.Mapper;
import java.io.IOException;
public class FlowMapper extends Mapper<LongWritable, Text, Text, FlowBean> {
private final Text outK = new Text();
private final FlowBean outV = new FlowBean();
@Override
protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException {
// 按制表符切分一行数据
String line = value.toString();
String[] fields = line.split("\\t");
// 提取手机号、上行流量、下行流量
String phone = fields[1];
long up = Long.parseLong(fields[fields.length – 3]);
long down = Long.parseLong(fields[fields.length – 2]);
outK.set(phone);
outV.setUpFlow(up);
outV.setDownFlow(down);
outV.setSumFlow(up + down);
context.write(outK, outV);
}
}
FlowReducer
package cn.edu.aust.flow;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.mapreduce.Reducer;
import java.io.IOException;
public class FlowReducer extends Reducer<Text, FlowBean, Text, FlowBean> {
private final FlowBean outV = new FlowBean();
@Override
protected void reduce(Text key, Iterable<FlowBean> values, Context context) throws IOException, InterruptedException {
long totalUp = 0;
long totalDown = 0;
// 累加同一手机号的所有流量
for (FlowBean bean : values) {
totalUp += bean.getUpFlow();
totalDown += bean.getDownFlow();
}
outV.setUpFlow(totalUp);
outV.setDownFlow(totalDown);
outV.setSumFlow(totalUp + totalDown);
context.write(key, outV);
}
}
4.4.3 驱动类
package cn.edu.aust.flow;
import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.fs.Path;
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;
import java.io.IOException;
public class FlowDriver {
public static void main(String[] args) throws IOException, InterruptedException, ClassNotFoundException {
Configuration conf = new Configuration();
Job job = Job.getInstance(conf);
job.setJarByClass(FlowDriver.class);
job.setMapperClass(FlowMapper.class);
job.setReducerClass(FlowReducer.class);
job.setMapOutputKeyClass(Text.class);
job.setMapOutputValueClass(FlowBean.class);
job.setOutputKeyClass(Text.class);
job.setOutputValueClass(FlowBean.class);
FileInputFormat.setInputPaths(job, new Path(args[0]));
FileOutputFormat.setOutputPath(job, new Path(args[1]));
System.exit(job.waitForCompletion(true) ? 0 : 1);
}
}
4.5 核心原理:Shuffle机制
Shuffle是MapReduce的核心,作用是将Mapper输出的数据整理后,有序地分发给Reducer。它是整个MapReduce最耗时的环节,也是性能优化的重点。
4.5.1 Shuffle整体流程
Mapper输出 → 分区(Partition) → 排序(Sort) → 溢写磁盘 → 合并(Merge) →
→ 拉取数据 → 归并排序 → 分组(Group) → 交给Reducer
4.5.2 分区(Partition)
决定每一条KV数据发送给哪个Reducer处理。
-
默认分区规则:根据Key的哈希值对ReduceTask数量取模
-
核心规则:相同的Key一定会进入同一个Reducer
-
可自定义Partitioner类实现自定义分区逻辑
4.5.3 排序(Sort)
MapReduce默认会对所有数据按Key进行字典序排序:
-
Map端:溢写前对缓冲区数据排序,溢写文件合并时归并排序
-
Reduce端:拉取所有Map端数据后,进行归并排序
4.5.4 分组(Group)
将排序后相同Key的Value聚合为一个Iterable集合,作为reduce()方法的输入。
4.5.5 Combiner局部聚合
Combiner是在Map端进行的局部聚合,作用是减少网络传输的数据量,提升执行效率。
-
本质上是一个Reducer,运行在MapTask节点上
-
使用前提:运算必须满足交换律和结合律,不能影响最终结果
-
适用场景:求和、求最值
-
不适用场景:求平均值(局部平均的平均 ≠ 全局平均)
Shuffle慢的原因:涉及多次排序(耗CPU)、多次磁盘读写(耗IO)、跨节点网络传输(耗带宽),通常占整个任务执行时间的70%以上。
4.6 实战案例:TopN排序
按分数从高到低排序,输出前10名学生信息。自定义Bean作为Key,需实现WritableComparable接口。
4.6.1 自定义排序Bean
package cn.edu.aust.topn;
import org.apache.hadoop.io.WritableComparable;
import java.io.DataInput;
import java.io.DataOutput;
import java.io.IOException;
public class ScoreBean implements WritableComparable<ScoreBean> {
private String name;
private int score;
public ScoreBean() {}
// 倒序排序:分数从高到低
@Override
public int compareTo(ScoreBean o) {
return Integer.compare(o.score, this.score);
}
@Override
public void write(DataOutput out) throws IOException {
out.writeUTF(name);
out.writeInt(score);
}
@Override
public void readFields(DataInput in) throws IOException {
this.name = in.readUTF();
this.score = in.readInt();
}
@Override
public String toString() {
return name + "\\t" + score;
}
// getter、setter省略
public String getName() { return name; }
public void setName(String name) { this.name = name; }
public int getScore() { return score; }
public void setScore(int score) { this.score = score; }
}
4.6.2 Mapper与Reducer实现
TopNMapper
package cn.edu.aust.topn;
import org.apache.hadoop.io.LongWritable;
import org.apache.hadoop.io.NullWritable;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.mapreduce.Mapper;
import java.io.IOException;
public class TopNMapper extends Mapper<LongWritable, Text, ScoreBean, NullWritable> {
private final ScoreBean outK = new ScoreBean();
@Override
protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException {
String line = value.toString();
String[] fields = line.split("\\t");
outK.setName(fields[0]);
outK.setScore(Integer.parseInt(fields[1]));
context.write(outK, NullWritable.get());
}
}
TopNReducer
package cn.edu.aust.topn;
import org.apache.hadoop.io.NullWritable;
import org.apache.hadoop.mapreduce.Reducer;
import java.io.IOException;
public class TopNReducer extends Reducer<ScoreBean, NullWritable, ScoreBean, NullWritable> {
private int count = 0;
private static final int TOP_N = 10;
@Override
protected void reduce(ScoreBean key, Iterable<NullWritable> values, Context context) throws IOException, InterruptedException {
for (NullWritable v : values) {
if (count < TOP_N) {
context.write(key, NullWritable.get());
count++;
} else {
return;
}
}
}
}
4.6.3 驱动类
package cn.edu.aust.topn;
import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.fs.Path;
import org.apache.hadoop.io.NullWritable;
import org.apache.hadoop.mapreduce.Job;
import org.apache.hadoop.mapreduce.lib.input.FileInputFormat;
import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat;
import java.io.IOException;
public class TopNDriver {
public static void main(String[] args) throws IOException, InterruptedException, ClassNotFoundException {
Configuration conf = new Configuration();
Job job = Job.getInstance(conf);
job.setJarByClass(TopNDriver.class);
job.setMapperClass(TopNMapper.class);
job.setReducerClass(TopNReducer.class);
job.setMapOutputKeyClass(ScoreBean.class);
job.setMapOutputValueClass(NullWritable.class);
job.setOutputKeyClass(ScoreBean.class);
job.setOutputValueClass(NullWritable.class);
FileInputFormat.setInputPaths(job, new Path(args[0]));
FileOutputFormat.setOutputPath(job, new Path(args[1]));
System.exit(job.waitForCompletion(true) ? 0 : 1);
}
}





