☰
MapReduce初级编程实践:从WordCount到Hadoop作业调优全解析
2026/10/9 3:37:59 网站建设 项目流程

简介:这是一份面向大数据初学者的 MapReduce 初级编程实践实验报告,对应林子雨《大数据原理与技术》第三版实验5,完整演示了基于 Hadoop 3.2.2 实现文件合并与去重操作的过程。报告从 Linux/Ubuntu 16.04 实验环境配置入手,给出输入文件 A、B 的样例数据和期望输出文件 C,并附可直接运行的 Java Map/Reduce 核心代码、Job 配置参数与提交方式;代码中 Map 阶段将文本行作为键、Reduce 阶段仅输出一次,从而有效剔除重复内容。报告中还展示了如何通过 FileInputFormat 与 FileOutputFormat 指定输入输出路径,并给出 Job 提交、运行结果与样例数据比对的方法,便于读者在本地集群中复现实验。资源为单个 docx 文档打包,压缩包约 1.28MB,包含完整的环境说明、核心代码与运行结果,内容结构清晰,适合作为课程实验报告或参考模板。目前已有 14444 人学习浏览,适用于高校大数据课程作业、复习备考,以及想要入门 Hadoop 并行编程的读者。

1. 实验报告不是抄代码:MapReduce 初级编程到底在练什么

很多同学拿到“大数据实验5:MapReduce 初级编程实践”这个题目,第一反应是去网上复制一份 WordCount,改两行输出路径就算交差。但真到了答辩或面试,被问一句“Map 阶段的数据到底存在哪里、Reduce 是怎么拿到中间结果的”,现场就冷场了。这个实验表面是让你跑通一个 Hadoop 作业,实际训练的是两件事:一是把“分而治之”的并行思想写进代码,二是能讲清楚一条数据从 HDFS 读入到写出经历了哪些节点。它适合刚接触大数据的学生,也适合刚切换到分布式开发的工程师,用最小成本看清 MapReduce 的运行边界。我见过不少用伪分布集群跑通十几个作业的人,图能画、日志会看,但一问传输层协议就露馅——这套实验的价值,恰恰是逼你从“会跑”走到“懂机制”。

2. 先把环境跑通:Hadoop 单机与伪分布的最小配置

2.1 选型:为什么实验用伪分布而不是真集群

实验报告里写“集群”的,绝大多数其实跑的是 Hadoop 伪分布模式。伪分布的意思是:在一台物理机上同时启动 NameNode、DataNode、SecondaryNameNode、ResourceManager、NodeManager 五个进程,每个进程用一个独立 JVM,听起来是集群,实际上共享同一套 CPU、内存和磁盘。

我在带新手做这个实验时,会建议直接进入伪分布,理由很简单:单机模式(local mode)虽然能执行 MapReduce 作业,但读写的是本地文件系统,不会经过 HDFS,也看不到分布式调度过程。实验要求的“数据切分”、“网络传输”、“节点间通信”这些观察点,在单机模式里全都看不到。反过来,直接上三台机器的真集群,又要先解决免密登录、NTP 同步、DNS 解析、节点心跳等一系列环境问题,反而偏离了“初级编程”的目标。

伪分布是性价比最好的教学形态:既能完整展示 HDFS 文件分块、副本机制、YARN 资源调度这些核心概念,又不需要额外花钱或凑机器。唯一要注意的是资源预算。默认配置下 Hadoop 会觉得自己拥有整台机器所有内存,如果物理机只有 8GB,启动后就容易卡死。所以动手之前,先把系统内存想清楚,后面配置里专门留几个参数来控制。

环境准备我一般做这几件事:

  • 安装 JDK,版本优先选 Hadoop 官方支持范围内的,比如 JDK 8 或 JDK 11,避免后面出现类库不兼容。
  • 配置免密 SSH 到本机,因为伪分布模式也需要用 SSH 拉起远程进程。
  • 下载稳定版 Hadoop 安装包,解压到 /opt 或 /usr/local,并把目录属主改成普通用户,不要用 root 直接跑。

这里强调一下:用普通用户跑。网上不少教程为了让省略权限问题,让你直接用 root 启动,结果后面 HDFS 目录权限和本地文件权限搅在一起,报错的时候很难判断是配置错还是权限错。

2.2 配置 core-site.xml / hdfs-site.xml 的三个必改参数

第一步是改两个核心配置文件。Hadoop 安装目录下 etc/hadoop 里有模板文件,我把常见的必改项整理成最小配置。

core-site.xml 里设置默认文件系统和临时目录:

<configuration> <property> <name>fs.defaultFS</name> <value>hdfs://localhost:9000</value> </property> <property> <name>hadoop.tmp.dir</name> <value>/home/hadoop/data/tmp</value> </property> </configuration>

fs.defaultFS 决定了 Hadoop 默认访问的文件系统。写成 hdfs://localhost:9000 就表示作业里写的路径都默认是 HDFS 路径,而不是本地路径。hadoop.tmp.dir 是 NameNode 和 DataNode 存放元数据和数据块的根目录,这个目录千万不能放在默认的 /tmp 下,因为系统重启会清空。我一般会在用户目录下建一个 data/tmp,并在开机脚本里手动创建。

hdfs-site.xml 里设置副本数和块大小、NameNode 元数据目录:

<configuration> <property> <name>dfs.replication</name> <value>1</value> </property> <property> <name>dfs.namenode.name.dir</name> <value>/home/hadoop/data/name</value> </property> <property> <name>dfs.datanode.data.dir</name> <value>/home/hadoop/data/data</value> </property> <property> <name>dfs.permissions.enabled</name> <value>false</value> </property> </configuration>

dfs.replication 在伪分布下必须设成 1。默认值是 3,单节点上存 3 份副本既没意义,还会直接把磁盘写满。dfs.namenode.name.dir 和 dfs.datanode.data.dir 分别指定元数据目录和数据块目录,这一步是给后面格式化做准备。如果这两个目录没有提前创建,NameNode 启动会报目录不存在。dfs.permissions.enabled 设成 false 是教学环境的偷懒选项,生产环境绝不能这么干,实验环境里可以省掉很多文件权限报错。

最后还要改 hadoop-env.sh 里的 JAVa_HOME,不然启动时找不到 JDK。

export JAVA_HOME=/usr/lib/jvm/java-8-openjdk-amd64

这里建议直接写绝对路径,不要依赖系统环境变量。有同学明明安装了 JDK,但 start-dfs.sh 起不来,就是因为 hadoop-env.sh 里的 JAVA_HOME 没写或写错了路径。

改完配置后,必须格式化 NameNode:

hdfs namenode -format

格式化会生成初始的元数据,同时清空之前的数据,所以只能在第一次使用或确认数据不要了之后执行。我见过旁边同学每次启动失败就 format 一次,结果 HDFS 里的数据全没了,实验做到后期才发现这个问题,后悔药都买不到。

2.3 启动与检查:jps 怎么看,日志在哪找

执行启动命令的顺序是:

start-dfs.sh start-yarn.sh

如果一切正常,jps 命令会列出五个进程:NameNode、DataNode、SecondaryNameNode、ResourceManager、NodeManager。如果只有四个或三个,比如就缺 DataNode,多半是格式化目录和当前进程使用的目录不一致,或者上次非正常退出留下了临时文件。

看到完整进程列表后,别急着跑作业,先做一个最小读写验证:

hdfs dfs -mkdir -p /input hdfs dfs -put /etc/hadoop/core-site.xml /input/ hdfs dfs -cat /input/core-site.xml | head -n 5

能读到文件内容,说明 HDFS 链路没问题。接下来再验证 YARN 是否就绪:

yarn node -list

如果显示 NodeManager 信息正常,环境就算通了。

日志是排查问题的第一现场。Hadoop 的日志分散在几个地方:

  • NameNode 和 DataNode 日志在安装目录的 logs/ 下,常见文件名是 hadoop-hadoop-namenode-主机名.log。
  • YARN 的 ResourceManager 日志也在 logs/ 下,文件叫 hadoop-hadoop-resourcemanager-主机名.log。
  • MapReduce 作业跑起来之后的具体任务日志,在 NodeManager 的 logs/userlogs/ 目录里,每一项作业都有自己的子目录。

很多同学作业失败后只盯着屏幕上的红字,其实红色只是异常摘要,真正的原因是藏在 userlogs 里的 syslog。养成这个习惯:先看日志,再猜原因。日志会告诉你哪个 Container 被杀、内存超了多少、哪个类找不到,这些都是后面章节里排查清单的原始素材。

3. 从 WordCount 到自定义 Mapper/Reducer:两个必须手写的程序

3.1 WordCount 逐行拆解:Map 端与 Reduce 端的数据流

环境通了之后,第一个要手写的程序就是 WordCount。不要复制粘贴,要一行一行敲进去才能理解每个类在干什么。我给出一个典型实现,注释里写清关键点:

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 { // 继承 Mapper,输入 key 是行偏移量,value 是整行文本 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); } } } // 继承 Reducer,输入 key 是单词,values 是同一个单词的所有计数 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); } }

这段代码里最容易被忽略的是泛型四元组。Mapper 的四个类型分别是:输入 key 类型、输入 value 类型、输出 key 类型、输出 value 类型。输入 key 是行偏移量,在文件切分后由 InputFormat 自动生成;输入 value 是一行文本;输出 key 是单词;输出 value 是计数 1。Reducer 的输入类型必须和 Mapper 的输出类型完全一致,否则运行期会报类型转换异常或默默丢数据。

shuffle 阶段是你看不见但必须讲清楚的部分。Map 端 context.write 之后,数据并不是直接进 Reduce,而是先做一个分区(partition),决定哪一部分给哪个 Reducer。然后按 key 排序,同一个 key 的 value 被合并成一个迭代器,再通过网络拷贝到对应的 Reduce 节点。这个过程跨越了节点边界,是 MapReduce 最消耗时间的地方。实验报告里的数据流图,一定要把这段过程画出来,从“map 输出”到“spill 落盘”再到“merge 合并”。

main 方法里值得一提的 setCombinerClass 这一行。我见过不少人把它注释掉,理由是“本地跑也能出结果”。确实能出结果,但 Combiner 的作用是在 Map 端先做一次局部合并,减少网络传输的数据量。WordCount 的累加逻辑满足交换律和结合律,所以可以直接复用 Reducer 类作为 Combiner。如果你自定义的 Reduce 逻辑不满足这两个性质,就不能这样偷懒。

参数方面,默认一个 Map 任务处理一个 HDFS 块,默认块大小在 Hadoop 2.x 和 3.x 里分别是 128MB 和 128MB(可按参数调)。如果输入文件特别小且数量特别多,每个文件都会启动一个 Map 任务,会产生大量任务开销。这种情况下有两个办法:一是用 combineTextInputFormat 把多个小文件合并成一个分片,二是提前在数据导入阶段做合并。实验数据通常不大,但养成“小文件会拖垮 NameNode 内存”的意识很重要。

3.2 自定义输出格式与多输出:实验加分项怎么写

实验要求如果只是统计词频,上面的 WordCount 够交差。但多数实验会留一个“选做”或“进阶”环节,比如按首字母把结果分别输出到不同文件,或者输出一个自定义对象。这时候就用得上 MultipleOutputs。

常见场景是:输入是订单数据,要按地区输出每个地区的统计结果。默认一个 Reduce 的输出是唯一的 part-r-00000 文件,但你希望得到 part-r-00000 存东部、part-r-00001 存西部。用 MultipleOutputs 可以按 key 或按 value 的条件动态决定输出文件。

代码可以这样改:

import org.apache.hadoop.mapreduce.lib.output.MultipleOutputs; public static class RegionReducer extends Reducer<Text, IntWritable, Text, IntWritable> { private MultipleOutputs<Text, IntWritable> mos; protected void setup(Context context) { mos = new MultipleOutputs<>(context); } public void reduce(Text key, Iterable<IntWritable> values, Context context) throws IOException, InterruptedException { int sum = 0; for (IntWritable v : values) sum += v.get(); String region = key.toString().substring(0, 2); mos.write(key, new IntWritable(sum), "region_" + region); } protected void cleanup(Context context) throws IOException, InterruptedException { mos.close(); } }

这一段里有两个细节。setup 方法在 Reduce 任务启动时执行一次,用来初始化 MultipleOutputs,不要在 reduce 方法里反复 new,否则每个 key 都会创建一个输出流,文件句柄很快就爆了。cleanup 方法必须显式 close,否则缓冲区的数据可能没落盘,作业显示 success 但结果文件是空的。

MultipleOutputs 能写任意数量的文件,但每个输出文件都会消耗一个文件句柄,实验数据量小没事,生产环境如果每个 key 一个文件,NameNode 的内存会被大量小文件塞满。所以这个功能的正确用法是控制输出文件数量,比如按业务维度聚合,而不是按原始 key 无限拆分。

另一个加分项是自定义 WritableComparable。默认的 Text 排序是按字典序,如果你想按数值排序,就要自己写一个类实现 WritableComparable 接口:

public class NumWritable implements WritableComparable<NumWritable> { private int value; public NumWritable() {} public NumWritable(int value) { this.value = value; } public void write(DataOutput out) throws IOException { out.writeInt(value); } public void readFields(DataInput in) throws IOException { value = in.readInt(); } public int compareTo(NumWritable o) { return Integer.compare(this.value, o.value); } public int hashCode() { return value; } public boolean equals(Object o) { if (!(o instanceof NumWritable)) return false; return this.value == ((NumWritable)o).value; } }

写这类代码时最容易翻车的点有两个:一个是 write 和 readFields 的字段顺序必须严格一致,先写 value 就一定要先读 value;另一个是 hashCode 方法,MapReduce 分区依赖它,如果你不覆盖 hashCode,不同对象可能分到不同区,甚至导致同样的 key 跑到不同 Reducer,结果就会变得不可预测。我见过有人写完自定义类后,成绩全对但迟迟不出结果,最后发现是一个字段没写进 equals 导致的去重失效。

到了这一步,你已经不是“会跑 WordCount”的水平了,而是在理解 MapReduce 序列化和排序机制。实验报告里如果能写清楚“为什么要写 write/readFields,为什么顺序不能乱”,老师一眼就能看出你是不是真动手了。

4. 作业提交、调试与常见报错排查:避坑清单

4.1 作业提交命令与参数:jar 包与输入输出路径的坑

代码编译之后打成 jar 包,提交作业的标准命令是:

hadoop jar wordcount.jar WordCount /input /output

这里有几个新手必踩的坑。

第一个坑是路径理解错误。 /input 和 /output 都指的是 HDFS 路径,不是本地路径。很多人在本地随便建了一个 input 目录,然后直接在命令行里把本地路径传进去,结果报“文件不存在”。正确做法是先把数据 put 到 HDFS 上,再用 HDFS 路径提交作业。

第二个坑是输出目录不能存在。MapReduce 要求输出目录在作业启动前是不存在的,它会在成功后自动创建。如果你的 /output 已经存在,直接报错“Output directory hdfs://localhost:9000/output already exists”。遇到这种情况,用下面的命令删掉再跑:

hdfs dfs -rm -r /output

这个设计是为了防止误覆盖上一次的结果。所以每次重跑作业前,要么改个新的输出路径,要么先清理旧目录。

第三个坑是参数传值。有时候需要临时调整 MapTask 数量或内存大小,可以这样传:

hadoop jar wordcount.jar WordCount \ -D mapreduce.map.memory.mb=1024 \ -D mapreduce.reduce.memory.mb=2048 \ /input /output

注意 -D 参数必须放在 jar 包名和类名之后、输入输出路径之前。放在其他位置会被当成新的输入路径,然后报“路径不存在”。

第四个坑是运行时类找不到。如果你在代码里引用了第三方库,比如常见的 JSON 解析库,打包时没有把这些依赖打进去,提交作业后会出现 NoClassDefFoundError。解决办法有两个:一是用 maven-shade-plugin 打 fat jar,把依赖一起打进去;二是用 job.setJarByClass 指定主类所在的 jar,然后把依赖 jar 放到 Hadoop 的 classpath 下。实验环境里我建议直接打 fat jar,省事且不容易漏。

作业提交后,可以在 YARN 的 Web 界面(默认 http://localhost:8088)查看进度。点进 Application ID,能看到每个 Map 和 Reduce 任务的状态、开始时间、结束时间、日志链接。这个界面是调试的指挥中心。

4.2 常见问题排查:现象、原因与解决

下面是我在这个实验里见别人踩过、自己也踩过的高频问题,按“现象 → 原因 → 解决”写清楚。

现象一:作业卡在 RUNNING 十几分钟,进度停在 33% 不动,最后 Container 被杀死,日志里写 “Container killed on request. Exit code is 143”。

原因是容器内存超限。伪分布模式下 YARN 默认给每个 Container 分配的内存很小,而物理机内存可能不足,或者代码里一次性加载了大量数据。我当时查了 userlogs 里的 syslog 才发现是物理内存超限。

解决:调整 YARN 和 MapReduce 的内存参数。先改 yarn-site.xml:

<property> <name>yarn.nodemanager.resource.memory-mb</name> <value>4096</value> </property> <property> <name>yarn.scheduler.maximum-allocation-mb</name> <value>4096</value> </property>

再在提交命令里调整 map 和 reduce 的内存。一般 Map 1GB、Reduce 2GB 就够了,别再往上调,因为你物理机总共就那么多内存。

现象二:任务失败,控制台只显示 “Job aborted due to stage FAILED”,但具体原因没写。

原因:错误被吞进了任务日志,控制台只输出最外层状态。我们组当时盯着红色报错看了半小时,什么线索都没有,其实去 NodeManager 找到对应的 application_xxx 目录,里面的 stdout 和 syslog 写得很清楚。

解决:作业失败后,立即去 logs/userlogs 下找最新的 application 目录,用 grep 搜 Exception 或错误关键字:

grep -r "Exception" /opt/hadoop/logs/userlogs/ | tail -n 50

现象三:Map 任务全部成功,Reduce 任务全部失败,错误是数组越界或空指针。

原因:Reduce 端处理了 Map 端没出现的边界值。比如我写过一个统计最大值的小程序,Map 输出 key 都正常,但某个输入文件是完全空文件,Reduce 拿到的 value 列表为空,直接取第一个元素就崩了。

解决:在 Reduce 里先判空或判断 value 数量,再执行逻辑。还有,检查你的 Reducer 输出类型和 job.setOutputKeyClass 是否一致。不一致时会在序列化阶段抛出复杂异常,表面上看是 ClassCastException,实际是类型设置错位。

现象四:启动 NameNode 时,日志报 “Incompatible clusterIDs”。

原因:你格式化过 NameNode,但 DataNode 的 data 目录里还保留着旧集群的 clusterID,或者反过来。我同学因为重装系统后直接复制了旧数据目录,导致两边的 clusterID 对不上。

解决:停止服务,删除 dfs.datanode.data.dir 里所有临时数据,重新格式化 NameNode。注意这只是教学环境的解法,生产环境严禁这样操作。

现象五:MapReduce 作业运行起来,但处理的数据总量比原始文件小,像是丢了一半数据。

原因:最常见的是使用文件分片时,某些文件是隐藏文件或临时文件被算进去了,或者文本编码问题导致行读取错误。另一个常见原因是自定义 Mapper 里只写了部分分支,没有 else 兜底。

解决:先用 hdfs dfs -count 确认输入目录里实际文件数和字节数,再在 Mapper 里加一个计数器:

context.getCounter("MyGroup", "processed_lines").increment(1);

作业跑完在控制台或 HistoryServer 的 Counter 里看到实际处理了多少行,和文件行数对比就能定位。

现象六:作业成功,输出文件是空的,或者只有 0 字节。

原因:Reduce 没有写任何数据。常见于条件判断写反了,或者输出路径被覆盖过。有一次我为了调试在 cleanup 里调了 context.write,结果写到了临时内存,cleanup 一结束数据就丢了。

解决:先看 job 成功的日志里 Reduce 输入输出计数,如果输出计数为 0,直接检查 reducer 的遍历逻辑。不要轻易在 cleanup 里写常规数据,cleanup 只适合 flush 缓冲区。

这六条是从几十次实验里筛出来的高频问题。每一条背后对应的都是某个机制理解不到位:内存管理、日志级别、序列化一致性、集群 ID 校验。你踩过一次,后面看报错就会形成直觉。所谓的“玄学”问题,九成都是日志没看全。

5. 把实验报告写成一页能讲清的数据流图:三种验证方法与一个技巧

实验报告交上去之前,我习惯做三件事来验证程序真的懂了。

第一,用参数扫描验证正确性。改输入文件大小,从 1KB 改成 10MB,观察 Map 任务数量是否从 1 变成 2 或更多。如果 Map 任务数量始终不涨,说明小文件合并逻辑生效了,或者文件分片没按预期切分。写报告时把这个变化画成表格,比写一百字都管用。

第二,用自定义计数器验证数据流。在 Mapper 里加一个计数器统计每个 Map 处理的记录数,在 Reducer 里加一个计数器统计每个 Reduce 收到的 key 数量。运行后在 YARN 的计数器页面看这两组数字。如果某些 Map 处理了 10 万条,某个 Map 只处理了 10 条,说明数据倾斜很严重。这时候实验报告里就可以写“我用自定义计数器识别到倾斜,然后用增加分区数或预聚合缓解”——这就是中级水平。

第三,画一页数据流图。图上从 HDFS 读入文件开始,经过 InputFormat 生成键值对,再到 Mapper 输出中间结果,标出分区、排序、合并、拷贝、归约,最后到 Reducer 输出。我见过很多人的报告里贴了十页代码,却没有一张图,答辩时讲半天别人也听不明白。一张图把所有环节列清楚,老师一眼就知道你懂了。

还有一个技巧:把输出结果拿到本地再排一次序,验证 Reduce 输出的 key 是否按默认字典序排列。MapReduce 的默认排序是 key 的字典序,但实验环境里可能因为分区方式产生多个输出文件,每个文件内部有序,文件间并不整体有序。你学会用 hadoop fs -getmerge 把结果合并回本地再 sort,就能发现这个细节,报告里写一句“Reduce 输出分区有序而非全局有序”,立刻显得深入。

我最开始做这个实验时,也以为跑通就是学会了。后来被一个“为什么输出文件多了一个”的问题追问到哑口无言,才回头把 shuffle 过程重新读了一遍。从那以后,我每做一个作业,都要把数据流图画一遍再动手写代码。希望帮到你。

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

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

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

立即咨询