跳到正文
Joeplover
学习笔记·2025-09-11·约 4 分钟阅读

Hadoop MapReduce 入门:WordCount 与 HDFS 操作

Hadoop 大数据学习实践,涵盖 HDFS 文件操作、MapReduce WordCount 编程模型和分布式计算基础。

Hadoop

Hadoop MapReduce 入门:WordCount 与 HDFS

项目背景

这是我最早接触大数据技术时的学习笔记。Hadoop 是大数据生态的基石,虽然现在生产环境很少直接写 MapReduce(用 Hive/Spark 替代),但理解它的核心思想对学习大数据技术体系非常关键。

Hadoop 体系结构

Hadoop 的核心组件:

  • HDFS(分布式文件系统): 数据存储层,文件被切块存储在多台机器上
  • MapReduce(分布式计算框架): 数据计算层,计算逻辑发送到数据所在的节点
  • YARN(资源调度): 集群资源管理,分配 CPU 和内存

核心思想:移动计算而不是移动数据。把计算程序分发到数据所在的节点执行,避免大数据量在网络中传输。

HDFS 架构

NameNode(主节点)
  ├── 管理文件系统的元数据(文件目录树、数据块映射)
  └── 不存数据本身

DataNode(从节点,多台)
  ├── 实际存储数据块(默认 128MB 一块)
  └── 定期向 NameNode 汇报心跳和块信息

Secondary NameNode
  └── 辅助 NameNode 合并编辑日志(不是热备!)

关键概念:

  • 文件被切成 block(默认 128MB),每个 block 有 3 个副本(默认)
  • 副本分散在不同机架的 DataNode 上,保证容错
  • NameNode 是单点故障——HDFS 2.x 引入了 HA 模式(Active/Standby NameNode)

MapReduce 计算模型

MapReduce 将计算分为两个阶段:

Map 阶段(分): 把输入数据拆成独立的块并行处理

public static class TokenizerMapper
    extends Mapper<Object, Text, Text, IntWritable> {

    private final static IntWritable one = new IntWritable(1);
    private Text word = new Text();

    public void map(Object key, Text value, Context context)
            throws IOException, InterruptedException {
        // value 是文件中的一行文本
        StringTokenizer itr = new StringTokenizer(value.toString());
        while (itr.hasMoreTokens()) {
            word.set(itr.nextToken());
            context.write(word, one);  // 输出 (word, 1)
        }
    }
}

Map 输出:(hello, 1), (world, 1), (hello, 1), ...

Shuffle 阶段(自动): Map 的输出按键排序、分组,发送到对应的 Reducer。这是 MapReduce 最复杂的阶段,也是性能瓶颈所在。

Reduce 阶段(合): 将相同 key 的 value 汇总

public static class IntSumReducer
    extends Reducer<Text, IntWritable, Text, IntWritable> {

    private 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();  // 累加同一个 word 的所有计数
        }
        result.set(sum);
        context.write(key, result);  // 输出 (hello, 2)
    }
}

Reduce 输出:(hello, 2), (world, 1)

数据流完整流程

输入文件
    ↓
InputFormat(将文件拆分成 Split)
    ↓
RecordReader(将 Split 解析成 key-value 对)
    ↓
Mapper.map() ← 每个键值对调用一次
    ↓
Shuffle(分区 → 排序 → 合并 → 拷贝到 Reducer)
    ↓
Reducer.reduce() ← 每个 key 及其 value 列表调用一次
    ↓
OutputFormat(写入 HDFS)

WordCount 的局限性

WordCount 只是入门 demo,真实生产中的 MapReduce 要处理:

  • 数据倾斜: 某个 key 的数据量远大于其他 key(比如"的"字在中文语料中出现几十万次),一个 Reducer 成了瓶颈
  • 大量小文件: 大量小文件会让 Mapper 启动过多,每个 Mapper 处理的数据太少
  • 迭代计算: MapReduce 每次都需要写磁盘(MR 之间的中间结果写到 HDFS),迭代算法效率极低——Spark 的 DAG 计算引擎就是为了解决这个问题

Hadoop 环境搭建

在 Ubuntu VM 上搭建 Hadoop 3.x 伪分布式模式:

  1. 安装 JDK 8+、配置 SSH 免密登录(localhost)
  2. 下载 Hadoop 包,配置 core-site.xml、hdfs-site.xml、mapred-site.xml
  3. 格式化 NameNode:hdfs namenode -format
  4. 启动:start-dfs.sh + start-yarn.sh
  5. 提交 WordCount:hadoop jar WordCount.jar /input /output

踩坑: Windows 上跑 Hadoop 需要 winutils.exe,不然格式化 NameNode 会报权限错误。

总结

MapReduce 是理解大数据计算模型的起点。虽然生产环境已经很少直接写 MR 作业,但"分而治之"的思想(Map=分、Reduce=合、Shuffle=排序分组)贯穿整个大数据生态。后续的 Spark、Flink 本质上是在 MapReduce 模型上的优化——减少磁盘 IO、增加内存计算、支持 DAG 更灵活的编排。