Hadoop MapReduce 入门:WordCount 与 HDFS 操作
Hadoop 大数据学习实践,涵盖 HDFS 文件操作、MapReduce WordCount 编程模型和分布式计算基础。
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 伪分布式模式:
- 安装 JDK 8+、配置 SSH 免密登录(localhost)
- 下载 Hadoop 包,配置
core-site.xml、hdfs-site.xml、mapred-site.xml - 格式化 NameNode:
hdfs namenode -format - 启动:
start-dfs.sh+start-yarn.sh - 提交 WordCount:
hadoop jar WordCount.jar /input /output
踩坑: Windows 上跑 Hadoop 需要 winutils.exe,不然格式化 NameNode 会报权限错误。
总结
MapReduce 是理解大数据计算模型的起点。虽然生产环境已经很少直接写 MR 作业,但"分而治之"的思想(Map=分、Reduce=合、Shuffle=排序分组)贯穿整个大数据生态。后续的 Spark、Flink 本质上是在 MapReduce 模型上的优化——减少磁盘 IO、增加内存计算、支持 DAG 更灵活的编排。