Hadoop 编程入门:WordCount 完整链路从环境配置到离线计算调优
2026/9/17 21:32:10 网站建设 项目流程

简介:面向大数据初学者与高校相关课程学生,这份实验报告围绕 Hadoop 编程实现 WordCount 单词统计程序,完整记录了从环境搭建到代码运行的各个环节,可用于课程实验、期末报告或自学参考。报告基于 Window11、Hadoop 虚拟机与 JDK1.8,详细展示 Eclipse 与 Hadoop 连接配置过程,并给出可直接运行的 Java 源码,包含 Mapper、Reducer 及 Job 主类代码,帮助理解 MapReduce 核心编程模型。包体为单个 doc 文档,约758KB,方便携带查看。目前已有1725人学习。内容按实验目的、环境、内容、步骤、代码和结果组织,步骤具体到命令和界面操作,适合遇到环境配置困难或首次接触 MapReduce 的读者按图索骥。

1. 大数据实验课上的 WordCount,远不止“会跑”这么简单

如果你在实验报告上写过“Hadoop编程实现WordCount”,大概率是照着教材敲一遍代码,跑通后截图保存,然后整理流程收尾。可到了自己搭环境、自己提交任务、自己调参数的时候,很多人会卡在入门的第一道坎上:伪分布式模式下明明能从日志里看到SUCCEEDED,却说不清 Mapper 和 Reducer 各自的数据流到底发生了什么;到了集群部署时,更是连yarn命令报错都看不出是哪里没配。WordCount 作为 Hadoop 生态的第一课,真正的价值不在代码本身,而在于它把 HDFS 存储、YARN 资源调度、MapReduce 计算模型这三个核心模块串成了一条完整链路。本文围绕 Hadoop 编程实现 WordCount 这条主线,从环境搭建、源码剖析、参数调到问题排查,完整走一遍离线计算中最典型的开发路径。无论是交实验报告、准备面试,还是初次接触大数据开发,这篇文章都能当一份可以直接落地的手册来用。

2. Hadoop 编程的基础:先理解存储与计算的边界

2.1 为什么 WordCount 能验证 Hadoop 环境是否正常

WordCount 的核心逻辑极度简单:读进一批文本文件,按空格等分隔符切分单词,然后统计每个单词的出现次数。越是简单的任务,越能暴露一套分布式系统的基础能力是否正常。它需要 NameNode 管理元数据、DataNode 存储真实数据块、YARN 负责分配容器、NodeManager 执行任务,一条链路里任何一环出问题,任务都会直接失败。所以几乎所有 Hadoop 课程的实验一,都把这当成“环境部署 + 开发环境搭建 + 提交作业”三联体验收动作。

做这个实验之前,需要先明确当前运行模式。伪分布式模式(Pseudo-Distributed Mode)意味着所有角色进程都在一台机器上,但 DataNode、NameNode、ResourceManager 等以独立进程运行,数据仍然按 HDFS 协议读写,依然经过完整的 MapReduce 生命周期。这个模式下跑通 WordCount,再去碰真实集群,只是把主机名、副本数和资源量放大而已,原理不发生任何变化。

2.2 离线环境下 Hadoop 与 Java 版本如何匹配

这里先讲一个常见返工点:Java 版本不匹配。Hadoop 3.x 必须在 Java 8 或 Java 11 环境下运行,而绝大多数教学实验用的是 Hadoop 3.3.x 系列,建议统一准备 JDK 1.8。如果你在 2024 年之后的实验环境里拿到的是最新 LTS 版 JDK 17,别说跑hadoop jar,连 HDFS 的start-dfs.sh都会抛UnsupportedClassVersionError。这是整个实验的第一个硬性边界。

安装完成之后,直接验证二进制版本:

java -version hadoop version

输出里同时出现openjdk version "1.8.0_xxx"Hadoop 3.3.6这类信息时,说明基础环境就绪。为了后续操作统一,需要把 HADOOP_HOME 写进/etc/profile,这是后续所有命令能直接执行的前提条件。

export HADOOP_HOME=/usr/local/hadoop-3.3.6 export PATH=$PATH:$HADOOP_HOME/bin:$HADOOP_HOME/sbin export JAVA_HOME=/usr/local/jdk1.8.0_202 export PATH=$PATH:$JAVA_HOME/bin

HADOOP_HOME指定了 Hadoop 二进制文件的根目录;HADOOP_HOME/sbin提供了start-dfs.shstart-yarn.sh等启停脚本;JAVA_HOME是 Hadoop 启动 Java 进程时寻找 JVM 的入口。

2.3 三个核心配置文件这样改,避免“能启动但报错”

伪分布式部署的关键配置文件有三个:

  • core-site.xml:指定 NameNode 的地址
  • hdfs-site.xml:指定 HDFS 副本数与 NameNode 元数据目录
  • yarn-site.xml:指定 ResourceManager 的通信地址与辅助服务

最精简的一组配置,直接写好放进去就能启动。复制这三份内容到对应配置文件的<configuration>标签之间即可:

<!-- core-site.xml --> <property> <name>fs.defaultFS</name> <value>hdfs://localhost:9000</value> </property>

fs.defaultFS表示 HDFS 的默认命名服务地址,后续的所有hdfs://操作都会默认指向本机 9000 端口。这个端口号需要同时和 NameNode 的 RPC 端口保持一致,如果你改过 8020,这里必须同步改。

<!-- hdfs-site.xml --> <property> <name>dfs.replication</name> <value>1</value> </property> <property> <name>dfs.namenode.name.dir</name> <value>/home/hadoop/data/namenode</value> </property>

伪分布式环境下dfs.replication只能设为 1,因为物理上只有一份数据。如果保持默认的 3,上传文件也不会失败,但会一直出现块副本不足的告警,很多新人会被这条 WARN 带偏,以为环境坏了。dfs.namenode.name.dir用于控制 NameNode 持久化元数据的落盘目录,第一次启动前必须保证这个目录不存在,或者用hdfs namenode -format格式化过。

<!-- yarn-site.xml --> <property> <name>yarn.nodemanager.aux-services</name> <value>mapreduce_shuffle</value> </property> <property> <name>yarn.nodemanager.env-whitelist</name> <value>JAVA_HOME,HADOOP_HOME,HOME</value> </property>

yarn.nodemanager.aux-services配置为mapreduce_shuffle是 MapReduce 运行的必经链路。Map 阶段产出的中间结果要经过 shuffle 发给 Reduce,这个动作由 NodeManager 上加载的辅助服务完成。你没配这一个字段,容器启动过程会直接失败,YARN UI 里看不到任何可用的 NodeManager 节点。

2.4 初始化 HDFS 与启动集群的标准动作

配置完成后,先格式化文件系统,再启动所有角色,然后用jps检查进程。这套顺序是固定且必须严格遵循的。

hdfs namenode -format start-dfs.sh start-yarn.sh jps

正常情况下jps能看到 6 个进程:NameNode、DataNode、SecondaryNameNode、ResourceManager、NodeManager 以及 jps 本身。如果缺了 DataNode,最常见的原因是格式化之后重新配置了dfs.namenode.name.dir,导致 NameNode 元数据与 DataNode 的注册信息不一致。处理方法是删除本地数据目录,重新格式化重启。

3. 编写与打包 WordCount:从 Mapper 到 Reducer 的完整源码实现

3.1 Java 源码:Mapper、Reducer、Main 三层结构

WordCount 的标准实现由三个类组成:一个继承Mapper的 TokenizerMapper、一个继承Reducer的 IntSumReducer、一个包含main方法的作业入口。先直接给出一版能用的完整代码,Maven 编译版本,JDK 1.8 直接通过。

import java.io.IOException; import java.util.StringTokenizer; 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.Mapper; import org.apache.hadoop.mapreduce.Reducer; import org.apache.hadoop.mapreduce.lib.input.FileInputFormat; import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat; public class WordCount { 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 { StringTokenizer itr = new StringTokenizer(value.toString()); while (itr.hasMoreTokens()) { word.set(itr.nextToken()); context.write(word, one); } } } 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(); } result.set(sum); context.write(key, result); } } public static void main(String[] args) throws Exception { Configuration conf = new Configuration(); Job job = Job.getInstance(conf, "word count"); job.setJarByClass(WordCount.class); job.setMapperClass(TokenizerMapper.class); job.setCombinerClass(IntSumReducer.class); job.setReducerClass(IntSumReducer.class); job.setOutputKeyClass(Text.class); job.setOutputValueClass(IntWritable.class); FileInputFormat.addInputPath(job, new Path(args[0])); FileOutputFormat.setOutputPath(job, new Path(args[1])); System.exit(job.waitForCompletion(true) ? 0 : 1); } }

这段代码的关键在于 Map 阶段的输出类型和 Reduce 阶段的输入类型必须完全对齐。map方法输出的键为Text(单词本身),值为IntWritable(数字1),Reducer 接收到的values是一个可以迭代的Iterable<IntWritable>,框架保证同一个键的所有值会被送入同一次reduce调用。job.setCombinerClass(IntSumReducer.class)这行不能省,它的作用是把每个 Map 任务内部已经累加过的部分结果先做一轮合并,减少网络传输量。

3.2 Maven 打包的最小 pom.xml 与编译参数

实验环境通常没有现成的 Maven 依赖仓库,所以最好用 Maven 管理依赖,并把 Hadoop 相关依赖标记为provided,避免把整套 Hadoop 运行库打进 JAR 里造成冲突。

<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd"> <modelVersion>4.0.0</modelVersion> <groupId>com.example</groupId> <artifactId>wordcount</artifactId> <version>1.0</version> <packaging>jar</packaging> <properties> <maven.compiler.source>1.8</maven.compiler.source> <maven.compiler.target>1.8</maven.compiler.target> </properties> <dependencies> <dependency> <groupId>org.apache.hadoop</groupId> <artifactId>hadoop-client</artifactId> <version>3.3.6</version> <scope>provided</scope> </dependency> </dependencies> </project>

执行mvn clean package之后,target 目录下会生成wordcount-1.0.jar,把该 JAR 和配套的input目录一起管理好,后续不管是伪分布式还是真实集群,提交命令都一样。首次执行时 Maven 会拉取hadoop-client及其传递依赖,这个过程需要访问中央仓库,提前执行一次即可。

3.3 本地输入数据准备与 HDFS 上传

先准备本地实验文件,比如在/home/hadoop/wc_input目录下创建两个文本文件,内容为多行英文句子。这里用一个国内机房常见的模拟方式,用循环生成几百行文本,观察数据量变大会对任务时间产生什么影响:

mkdir -p /home/hadoop/wc_input for i in $(seq 1 200); do echo "hadoop mapreduce spark flink hadoop hdfs yarn" >> /home/hadoop/wc_input/file_$i.txt; done

文件准备号之后,在 HDFS 上创建输入目录并上传:

hdfs dfs -mkdir -p /user/hadoop/wordcount/input hdfs dfs -put /home/hadoop/wc_input/*.txt /user/hadoop/wordcount/input/ hdfs dfs -ls /user/hadoop/wordcount/input/

-mkdir -p的作用是递归创建多级目录,与 Linux 的mkdir -p行为一致。如果直接跑一个不存在的路径,客户端会报FileNotFoundException,这一点在实验报告里很容易被忽略,但对理解 HDFS 的路径模型有直接帮助。

4. 提交作业与运行参数调整:从日志和 Web UI 定位问题

4.1 用 hadoop jar 提交任务并解读执行日志

打包完成、输入数据也放到位之后,提交作业就是一条命令的事。这里有一个细节:如果 JAR 包打过多次包,务必确认提交的是最新构建产物,否则容易出现改了几行代码但运行结果不变的现象。

hadoop jar /home/hadoop/wc.jar WordCount /user/hadoop/wordcount/input /user/hadoop/wordcount/output

输出路径/user/hadoop/wordcount/output必须不存在。MapReduce 框架在启动时会检查输出路径是否已存在,存在则直接抛FileAlreadyExistsException,这是为了防覆盖。真正做重复实验时,应该用一个带时间戳的新路径,或者运行前删掉旧目录:

hdfs dfs -rm -r /user/hadoop/wordcount/output

任务执行过程中,终端会滚动输出两类关键日志:一类是INFO mapreduce.Job: map 100% reduce 100%,表示整体执行进度;另一类是最后一段计数器列表。需要着重看的是Map input recordsReduce output records这两个计数器。前者显示输入文件被拆分成多少条<key, value>记录,后者是最终写出的结果条目数。如果 Map 阶段结果正常,但 Reduce 输出为 0,问题百分之百出在 Reducer 代码逻辑上。

4.2 YARN Web UI 判断资源分配是否合理

任务提交后,打开http://localhost:8088/cluster可以看到正在运行的 Application。进入详情页能看到每个 Map 任务与 Reduce 任务的资源分配情况。伪分布式环境默认每个容器只有 1 核 1GB 内存,如果输入文件较大,Map 任务会长时间处于 PENDING 状态。原因很简单:整个机器的可用资源只够同时起 1 到 2 个容器,其他的都在排队。

在真实集群里,这些资源参数就要按节点规格调整。在mapred-site.xml中显式配置以下参数:

<property> <name>mapreduce.map.memory.mb</name> <value>1024</value> </property> <property> <name>mapreduce.reduce.memory.mb</name> <value>2048</value> </property> <property> <name>yarn.app.mapreduce.am.resource.mb</name> <value>1024</value> </property>

mapreduce.map.memory.mb控制每个 Map 容器能用的物理内存上限,当 Mapper 对象内部缓存过多时,超过这个值会被 NodeManager 直接杀掉,现象是任务失败且报Container killed by the ApplicationMastermapreduce.reduce.memory.mb同理,只是作用在 Reducer 上。

4.3 三个必调参数:Combiner、输入分片大小与 JVM 复用

以下三个参数在实际生产与实验中都有明显效果,值得单独展开。

第一个是mapreduce.job.reduces,表示 Reduce 任务个数。默认值是 1,但如果数据量上来了,1 个 Reducer 会成为整个任务的瓶颈。经验值可以设为节点 CPU 核数的 0.95 倍左右:

hadoop jar wc.jar WordCount -D mapreduce.job.reduces=4 /input /output

提交命令中-D参数运行时注入配置,优先级高于 jar 包内配置。Reducer 数不是越多越好,因为每个 Reducer 的输出结果会写到单独文件part-r-00000part-r-00001等,后续如果做全量聚合,反而要多一步合并操作。

第二个是mapreduce.input.fileinputformat.split.minsize。HDFS 默认块大小是 128MB,输入文件如果远小于这个值,每个文件会独占一个 Map 任务。当实验数据包含几千个小文件时,启动 Map 任务的开销占到总时长的 50% 以上。把最小分片调大,强制多个小文件合并进同一个 Map 任务:

-D mapreduce.input.fileinputformat.split.maxsize=67108864

64MB 作为分片最大值,一下能减少大量 Map 启动开销。但要注意它和 HDFS 块大小之间的配合关系,分片跨越块边界时会多一次网络读取,这在小集群里可以接受。

第三个是mapreduce.job.jvm.numtasks,控制同一个 JVM 可以重复运行的 task 数量。默认值是 1,意味着每个 Map 或 Reduce 任务启动一个新的 JVM。数据量大时 JVM 启动开销非常可观:

-D mapreduce.job.jvm.numtasks=5

设为 5 之后,同一个 JVM 会串行执行最多 5 个 task,适合计算逻辑不重的作业。这个参数在最理想的情况下能把耗时压到原来的 1/3 左右。

4.4 常见报错与对应解法:ClassNotFound、端口占用与磁盘空间

实验报告里最常见的失败场景是ClassNotFoundException: WordCount,它通常不是代码问题,而是提交命令写错了类名。hadoop jar wc.jar WordCount里的WordCount是包含main方法的完整类名,如果类在带包名的目录下,就得写全限定名:

hadoop jar wc.jar com.example.WordCount /input /output

第二个高频问题是NameNode 9000 端口被占用。如果之前/tmp目录下残留了旧版本 Hadoop 的进程,或者本机有其他服务占用 9000,NameNode 会启动失败。判断方法是:

netstat -nltp | grep 9000

如果被占用,要么杀掉旧进程,要么改fs.defaultFS的端口号。改端口后需要同步修改hdfs-site.xmldfs.namenode.rpc-address,否则会陷入改了配置但看不见效果的循环。

第三个问题在真实集群或服务器上跑实验时常出现:DataNode 数据目录磁盘写满。伪分布式模式下数据都在 Linux 根分区,HDFS 写数据不会产生任何预警,直到块写入失败。处理方式是先扩容或清理旧数据,再手工触发均衡:

hdfs balancer

5. Hadoop 中 Combiner 与 Partitioner 的配合逻辑

WordCount 的进阶改法离不开两个类:CombinerPartitioner。Combiner 本质上是运行在 Map 端的迷你 Reducer,从代码上看只是把job.setReducerClass增加一行job.setCombinerClass。但它的执行位置决定了它的能力边界:Combiner 的输入输出类型必须与 Mapper 的输出类型完全一致,同时还要满足交换律与结合律。统计单词次数这种求和操作具备这两个性质,所以 WordCount 可以安全地复用 Reducer 作 Combiner。你要是换成一个求平均值逻辑,复用 IntSumReducer 就会产生严重错误:每个切片里的平均值再做一次平均,最终结果完全错误。

Partitioner 则决定了 Map 输出的每个键流向哪个 Reducer。默认实现是哈希取模,(key.hashCode() & Integer.MAX_VALUE) % numReduceTasks。在 WordCount 里,如果你想让某些单词集中到同一个 Reducer,或者想把以指定前缀开头的单词分到同一个分区,就需要自定义 Partitioner。以下代码按首字母把 A 开头的单词全部分到第 0 号分区,其他单词进入第 1 号分区:

public static class FirstLetterPartitioner extends Partitioner<Text, IntWritable> { @Override public int getPartition(Text key, IntWritable value, int numPartitions) { if (numPartitions == 0) return 0; String word = key.toString(); if (word.isEmpty()) return 0; if (word.charAt(0) == 'A' || word.charAt(0) == 'a') { return 0; } else { return 1 % numPartitions; } } }

然后把分区器与 Reducer 数量绑定,提交任务时写明 Reducer 数为 2:

hadoop jar wc-part.jar WordCountWithPartitioner -D mapreduce.job.reduces=2 /input /output

这样输出目录下会出现part-r-00000(A 开头单词)和part-r-00001(其他单词)两个文件。结合 Combiner 和 Partitioner 调整,WordCount 就不再只是“抄一遍跑通”的题了,它变成了一组可控的数据分片实验,能精确看到每个环节的数据分流路径。

6. 用一行 Linux 命令验证统计结果,再理解 MapReduce 的边界

任务跑完后,先别急着写实验报告。验证结果最能暴露你的理解程度。HDFS 上的输出文件可以通过hdfs dfs -cat直接查看:

hdfs dfs -cat /user/hadoop/wordcount/output/part-r-00000 | head -n 20

为了做交叉验证,可以同时把本地输入文件合并后交给awk统计,两边的数字一对比,能立即确认任务没有丢数据:

cat /home/hadoop/wc_input/*.txt | tr ' ' '\n' | sort | uniq -c

这两组输出的结果应当逐行对齐。如果数量差很多,跑一遍所有输入文件的校验和,问题大概率出在 HDFS 上传时部分文件写入了损坏副本。另外,把 HDFS 上的输出文件搬到本地再检查一遍是常见操作:

hdfs dfs -getmerge /user/hadoop/wordcount/output /home/hadoop/wc_output_local/result.txt

-getmerge会把输出目录下所有 part 文件合并成一个本地文件,适合把结果集当作后续统计分析输入。还有一个小技巧值得记下来:在提交命令里指定-D mapreduce.output.fileoutputformat.compress=true,可以压缩输出结果减小落盘量,需要同时配套指定压缩格式为org.apache.hadoop.io.compress.GzipCodec,特别适合节点磁盘空间紧张的场景。

WordCount 能验证的环境配置与代码逻辑,到这里才算真正闭环。

本文还有配套的精品资源,点击获取

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询