本文系统梳理大数据技术栈的核心基础——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 大数据的三类数据格式

大数据按照结构特性可分为三类:

  1. 结构化数据
    采用标准化格式、具备明确定义结构的数据,存储和排列有固定规律,易于程序读取和处理。典型应用是关系型数据库中的表数据,存储需求包括高速读写、数据备份、数据共享、容灾等。

  2. 半结构化数据
    不遵循严格的数据模型,没有固定表结构,但包含标签标记来分隔语义、对数据分层组织,也被称为自描述结构,通常以树或图的形式存储。

    • 典型格式:XML、JSON

    • 特点:每条记录的属性个数可以不固定,灵活性远高于结构化数据

    • 常见场景:邮件系统、Web页面、资源库、档案系统等

  3. 非结构化数据
    数据结构不规则、不完整,没有预定义的数据模型,无法用二维逻辑表表示。企业中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架构演进
  1. Hadoop 1.x 架构

    • MapReduce同时承担资源管理数据处理两大职责,负载较重,扩展性差

    • HDFS仅负责分布式文件存储

    • 仅支持MapReduce一种计算框架

  2. Hadoop 2.x 架构

    • 将资源管理功能从MapReduce中剥离,由YARN统一负责集群资源管理和任务调度

    • MapReduce仅负责数据处理,负载大幅降低

    • YARN支持为多种计算框架(MapReduce、Spark、Flink等)提供资源管理,生态兼容性大幅提升

    • HDFS仍负责分布式文件存储

  3. Hadoop 3.x 架构优化
    在2.x的基础上对四大模块进行了性能优化与功能增强:

    • Hadoop Common:通用工具包优化,支持更多操作系统与硬件架构

    • MapReduce:任务优化,提升小任务执行效率

    • YARN:资源调度策略增强,支持更细粒度的资源隔离

    • HDFS:支持纠删码技术,降低存储开销;支持NameNode高可用,解决单点故障问题


第二章 Java反射机制(Hadoop底层基础)

Hadoop框架的底层大量依赖Java反射机制实现类的动态加载、对象的动态创建与方法调用,是理解Hadoop源码的前置基础。

2.1 反射机制概述

常规编程中,我们通过类创建对象;而反射是将这个过程反转:在程序运行时,通过对象获取其所属类的完整信息,或动态创建对象、调用方法。

2.1.1 反射的核心作用

Java反射机制的核心能力是动态性,具体包含4点:

  1. 运行时动态构造任意一个类的对象

  2. 运行时获取任意一个对象所属类的完整信息

  3. 运行时调用任意一个类的成员变量和成员方法

  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对象有三种标准方式:

  1. 全限定类名获取Class.forName("全限定类名")

  2. 对象实例获取对象名.getClass()

  3. 类名直接获取类名.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对象,再通过它实例化对象。
操作步骤:

  1. 通过getConstructors()获取类的所有构造方法

  2. 定位到目标有参构造对应的Constructor对象

  3. 调用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 反射调用普通方法

步骤:

  1. 通过getMethod()获取指定的Method对象

  2. 调用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)架构,包含四大核心角色:

  1. NameNode(NN,主节点):集群的管理者

    • 管理HDFS的名称空间(目录树)

    • 配置文件副本策略

    • 管理数据块(Block)的映射信息

    • 处理客户端的读写请求

  2. DataNode(DN,从节点):实际存储数据的工作节点

    • 存储实际的数据块

    • 执行数据块的读写操作

    • 定期向NameNode汇报自身状态和块信息

  3. Client(客户端)

    • 文件上传时将文件切分为固定大小的Block,逐个上传

    • 与NameNode交互,获取文件的位置信息

    • 与DataNode交互,执行实际的数据读写

    • 提供命令行、API等方式管理和访问HDFS

  4. 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,查看任务运行状态、资源使用情况

前置配置

  1. 关闭服务器防火墙:
systemctl stop firewalld
systemctl disable firewalld
  1. 本地电脑配置hosts映射(C:\Windows\System32\drivers\etc\hosts),添加集群节点IP与主机名映射

  2. 浏览器访问:http://NameNode主机名:9870 查看HDFS,http://ResourceManager主机名:8088 查看YARN

3.3 快速入门:官方词频统计案例

使用Hadoop自带的示例Jar包快速体验MapReduce任务:

  1. 准备本地测试文件
vi /home/test.txt
# 输入测试文本,例如:
# hello hadoop
# hello world
# hadoop bigdata
  1. 在HDFS创建输入目录并上传文件
hdfs dfs -mkdir -p /wordcount/input
hdfs dfs -put /home/test.txt /wordcount/input
  1. 运行官方词频统计示例
cd $HADOOP_HOME/share/hadoop/mapreduce
hadoop jar hadoop-mapreduce-examples-*.jar wordcount /wordcount/input /wordcount/output
  1. 查看结果
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 目录操作
  1. 查看目录-ls
# 查看根目录
hdfs dfs -ls /
# 递归查看所有子目录
hdfs dfs -ls -R /
# 人性化显示文件大小
hdfs dfs -ls -h /wordcount
  1. 创建目录-mkdir
# 创建单层目录
hdfs dfs -mkdir /test
# 递归创建多级目录
hdfs dfs -mkdir -p /a/b/c
  1. 查看目录大小-du
# 查看目录下每个文件的大小
hdfs dfs -du /wordcount
# 查看目录总大小
hdfs dfs -du -s -h /wordcount
3.4.2 文件操作
  1. 上传文件-put
# 上传本地文件到HDFS指定目录
hdfs dfs -put 本地文件路径 HDFS目录路径
# 强制覆盖已存在的文件
hdfs dfs -put -f 本地文件路径 HDFS目录路径
  1. 下载文件-get
# 下载HDFS文件到本地
hdfs dfs -get HDFS文件路径 本地目录路径
# 强制覆盖本地已存在的文件
hdfs dfs -get -f HDFS文件路径 本地目录路径
  1. 查看文件内容-cat
hdfs dfs -cat /wordcount/input/test.txt
  1. 移动/重命名-mv
# 移动文件到目标目录
hdfs dfs -mv /a.txt /test/
# 重命名文件
hdfs dfs -mv /a.txt /b.txt
  1. 复制文件-cp
hdfs dfs -cp /source/file.txt /target/
  1. 删除文件/目录-rm
# 删除文件
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实现定时执行:

  1. 检查并安装Crontab
rpm -qa | grep crontab
# 未安装则执行
yum -y install vixie-cron crontabs
  1. 启动Crontab服务
systemctl start crond.service
systemctl enable crond.service
  1. 给脚本添加执行权限
chmod 777 /usr/data/uploadHDFS.sh
  1. 编辑定时任务
crontab -e
# 添加配置:每10分钟执行一次
*/10 * * * * /usr/data/uploadHDFS.sh

3.6 HDFS Java API 操作

HDFS提供Java API实现编程式操作文件系统,核心类位于org.apache.hadoop.fs包下。

3.6.1 核心API介绍
类名功能说明
FileSystemHDFS文件系统核心类,提供文件的增删改查等所有操作
FileStatus封装文件/目录的元数据(大小、块大小、副本数、修改时间等)
FSDataInputStreamHDFS输入流,用于读取HDFS文件
FSDataOutputStreamHDFS输出流,用于写入数据到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配置参数的优先级从高到低为:

  1. 客户端代码中通过configuration.set()设置的值

  2. 项目ClassPath下的自定义配置文件(如hdfs-site.xml

  3. 服务器端的自定义配置文件(xxx-site.xml

  4. 服务器端的默认配置文件(xxx-default.xml

3.7 HDFS核心原理:读写流程

3.7.1 HDFS写文件流程
  1. 客户端向NameNode发起写文件请求

  2. NameNode校验权限、路径是否存在,通过后返回可写入的DataNode列表

  3. 客户端将文件切分为Block,按顺序逐个写入

  4. 客户端通过Pipeline管道机制,将数据块依次写入多个DataNode副本

  5. 每个DataNode写入完成后向上游返回确认

  6. 所有副本写入完成后,客户端通知NameNode关闭文件

3.7.2 HDFS读文件流程
  1. 客户端向NameNode发起读文件请求

  2. NameNode返回文件的Block位置信息,按距离客户端的远近排序

  3. 客户端就近选择DataNode读取第一个Block

  4. 读取完成后校验数据完整性,继续读取下一个Block

  5. 所有Block读取完成后,客户端关闭流


第四章 MapReduce分布式计算框架

MapReduce是一个分布式运算程序的编程框架,是Hadoop核心的计算层。它将用户编写的业务逻辑与框架默认组件整合,自动分发到集群上并行运行,开发者无需关心底层分布式细节,只需专注业务逻辑。

4.1 核心基础

4.1.1 MapReduce运行架构

一个完整的MapReduce程序运行时有三类进程:

  1. MrAppMaster:负责整个程序的调度和状态协调,一个Job对应一个AppMaster

  2. MapTask:负责Map阶段的数据处理,并行执行多个

  3. ReduceTask:负责Reduce阶段的数据汇总处理,并行执行多个

4.1.2 Hadoop序列化机制

序列化:将内存中的Java对象转换为字节序列,用于持久化存储或网络传输。
反序列化:将字节序列恢复为内存中的Java对象。

Java原生的Serializable是重量级序列化,会附带大量额外信息(校验头、继承体系等),网络传输效率低。因此Hadoop自研了Writable序列化机制,特点是:

  • 紧凑:存储空间利用率高

  • 快速:读写额外开销小

  • 互操作:支持多语言交互

Hadoop常用序列化类型

Java原生类型Hadoop Writable类型
BooleanBooleanWritable
ByteByteWritable
IntIntWritable
FloatFloatWritable
LongLongWritable
DoubleDoubleWritable
StringText
MapMapWritable
数组ArrayWritable
nullNullWritable

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);
    }
}

更多推荐