简介:大数据实验5实验报告聚焦 MapReduce 初级编程实践,对应林子雨《大数据原理与技术(第三版)》实验5,适合刚接触 Hadoop 的大数据学习者。实验以文件合并与去重为例,演示了 MapReduce 在分布式环境下的并行处理流程。压缩包为 docx 格式,共 1 个文件,大小约 1.28MB,内含实验环境、实验内容与完成情况、代码及运行验证等完整记录。报告基于 Hadoop 3.2.2,给出了 Java 编写的映射阶段(Map)与归约阶段(Reduce)核心代码、输入输出文件样例,以及 MapReduce 作业的提交与配置步骤;通过 A、B 两文件合并去重生成 C 文件的过程,明确说明了映射阶段如何输出键值、归约阶段如何完成去重,并附带实验结果供读者对照验证。学习者既可据此复现该实验,也能将去重思路迁移到其他数据处理需求中。已有 14439 人学习,适合需要完成类似实验或理解 MapReduce 初阶应用的读者参考。
1. MapReduce 初级编程实践到底在做什么:一份实验报告背后的完整落地路径
上学期第一次在实验报告里写“MapReduce 初级编程实践”时,我以为就是登录实训平台跑通一个 WordCount,然后截图交差。结果同一个代码在老师电脑上秒出结果,换到自己搭的伪分布式集群上就反复报错:作业一直 ACCEPTED、输出目录冲突、日志里找不到完整堆栈。最后发现大多是配置和类型这类基础细节。这篇笔记就是当时填坑后整理的落地流程,从理解 Map 和 Reduce 各自拿到什么数据、产出什么数据,到写代码、打包、上传、提交、查日志、核对输出文件,再到现在做实验前必查的几个参数。适合第一次做 MapReduce 实验、要交实验报告但不想停留在界面操作的同学,也适合准备大数据岗位面试、想在简历上写“独立完成 HDFS 与 MapReduce 实验”的入门工程师。
2. 把 WordCount 拆开看:MapReduce 编程模型里的数据约定与设计取舍
MapReduce 初级编程的第一关不是写代码,而是理解“数据以什么形式进出”。MapReduce 把计算过程固定成 map、shuffle、reduce 三个阶段,map 的输入是一条一条记录,shuffle 帮你做排序和分组,reduce 拿到的是“同一个 key 的一组 value”。WordCount 正好把这三个阶段全部走了一遍,而且数据含义足够简单,适合当第一个 mapreduce 编程实例来对照学习。下面先讲清楚 map 和 reduce 的职责边界,再讲类型约束和 shuffle 里的隐藏规则。
2.1 为什么初级实践几乎都用 WordCount:它的数据流能讲清整套约定
我见过不少实训平台,比如头歌上的 mapreduce 基础编程题,第一关基本都拿 WordCount 当模板。原因是它能把“数据流”讲得特别直观:输入是一整个文本文件,但 Mapper 收到的是一条一行,Reducer 收到的不是一个列表,而是“同一个单词的所有出现次数”。很多同学在初学阶段卡住的点,其实是不知道 map 方法里的 value 到底是一行文件还是一整个文件。答案是一行。Hadoop 的 TextInputFormat 默认会把文件按行切开,每行封装成一个<key, value>,key 是该行在文件中的字节偏移量,value 才是这一行的文本内容。所以你在 map 方法里做的第一件事,永远是value.toString(),然后按空格或标点切分成单词,再逐个输出<单词, 1>。
map 端输出之后,框架会对 key 排序、分区,再把相同 key 的所有 value 合并到一起交给 reduce。这个过程你不用写一行代码,但它的存在决定了两个实验思路:第一,reduce 的输入不是散乱的 value,而是“同一个 key 排好序的一组 value”,所以你不能指望它保持文件原始顺序;第二,因为框架帮你做了分组,所以 WordCount 的 reduce 只需要把迭代器里的整数累加就行,复杂度非常低。理解了这条数据流,后面排错才有个方向,比如报类型不匹配时,你会本能地先去查 map 输出的 key/value 类型,而不是在代码里乱试。
2.2 Mapper 与 Reducer 的类型约束:源码层面对输入输出有哪些硬性要求
在写第一个完整 mapreduce 程序之前,建议先记住这张类型映射表,它能解释掉一半报错。Hadoop 的序列化机制用的是Writable,代替 Java 原生类型,因为它的序列化结果更紧凑、跨语言兼容性更好。常见对应关系如下:
| Java 类型 | Hadoop Writable 类型 | 适用场景 |
|---|---|---|
| String | Text | 文本内容、单词、路径 |
| Long | LongWritable | 行偏移量、时间戳 |
| Integer | IntWritable | 计数器、次数 |
| Float | FloatWritable | 评分、比例 |
| null 占位 | NullWritable | 只需要 key 或只需要 value 的场景 |
一旦你确定 map 输出的 key 用 Text、value 用 IntWritable,那么整个作业里必须有几处保持类型一致:job.setMapOutputKeyClass、job.setMapOutputValueClass、job.setOutputKeyClass、job.setOutputValueClass。很多初学翻车的原因就是只设了 job 的输出类型,没设 map 输出类型。map 输出的类型默认跟着作业最终输出类型走,但 reduce 的输出类型和 map 的输出类型在自定义 Mapper 类里经常不一致,一旦不一致,运行到 reduce 拉取数据时就会报java.io.IOException: Type mismatch in value from map。
还有一个细节来自 Reducer 的 API 设计:reduce 方法收到的Iterable<IntWritable> values并不是一个可以安全缓存到 List 再遍历的数据结构。迭代器每次返回的是同一个对象的引用,Hadoop 会复用这个 Writable 对象来减少 GC 压力。如果初学者在 reducer 里加了一句list.add(val),最终 list 里全是最后一个值。这是我在代码评审里看过的经典问题,也是新手从能跑到结果错之间最隐蔽的一步。要保存多个值,必须new IntWritable(val.get())复制后再放入集合。
2.3 shuffle 和 Combiner:决定这个实验能不能用真实数据的两个隐藏点
shuffle 是 MapReduce 里最像“黑匣子”的部分。它发生在 map 输出之后、reduce 输入之前,框架帮你完成分区、排序、归并、压缩。初级编程实践不需要修改 shuffle 逻辑,但你要知道它的两个直接后果:第一,reduce 拿到的 key 是有序的,所以不同 map 任务产出的相同单词会被合并到同一个 reduce 里;第二,map 输出的中间结果会先写本地磁盘,再通过网络拉取,数据量越大,shuffle 耗时就越高。
这就是 Combiner 存在的意义。Combiner 是 map 端的局部聚合器,可以在 map 输出落盘之前先做一次求和。WordCount 的 Combiner 可以直接复用 Reducer 类,因为局部求和和全局求和的逻辑完全一样,都是把一组相同 key 的 value 累加。但换一个场景就要小心:如果 reduce 计算的是平均值,Combiner 就不能直接复用 Reducer,否则局部平均会被再次平均,导致最终结果错误。一个安全的判断标准是:Combiner 的输入和输出类型必须与 map 输出类型一致,且运算满足交换律和结合律。求和、取最大值、取最小值都可以,平均值绝对不行。这个点经常被当作实验报告里的讨论题:“为什么 WordCount 可以复用 Combiner?”
3. 在 Hadoop 上跑通第一个 mapreduce 程序:从环境准备到完整的 WordCount 代码
确定了编程模型之后,下一步是把代码真正跑起来。这一步最花时间的往往不是 Java 代码本身,而是环境。我会先说明伪分布式和完全分布式的选型判断,再给出一份我常用的最小配置,最后贴出完整的 WordCount 实现和打包命令。
3.1 环境选型:伪分布式还是容器化自建,决定了你后面几小时的调试体验
做这份实验时,我见过三种环境做法:第一种,直接用头歌这类网页实训平台,代码编辑、上传、运行都在网页里完成,适合验证算法逻辑,但实验报告里写不出环境配置过程;第二种,自己用虚拟机装一个 Hadoop 伪分布式集群,所有进程都在一台机器上,能完整走一遍部署、启动、提交作业、查日志的流程;第三种,有条件的话用 Docker 起一个单节点 Hadoop,隔离性更好,改坏了直接删容器重启。
我建议的优先级是:如果目标是理解 MapReduce 编程实践,平台和自建都可以;如果简历里想写大数据集群部署策略,必须至少自己做过一次伪分布式。伪分布式只需要二到四 GB 内存就能跑,DataNode、NameNode、ResourceManager、NodeManager 都跑在同一个进程组里,但身份各自独立。它最重要的价值在于让你能练习“启动、停止、看日志、清理临时文件”这套操作。报告里也可以多写一段:为什么实验选择伪分布式而不是完全分布式,常见的理由包括资源受限、作业规模小、便于观察各节点日志。
3.2 最小配置:core-site.xml、yarn-site.xml 与内存参数
我一般把 Hadoop 安装在/opt/hadoop,数据目录单独放到/data/hadoop。伪分布式下需要重点改三个配置:core-site.xml里指定 NameNode 地址,hdfs-site.xml里把副本数设为 1,yarn-site.xml里开启 shuffle 辅助服务。副本数设为 1 不是因为生产环境也这么干,而是单节点上设 3 个副本会白白增加落盘等待,还可能因为磁盘空间不足导致 block 放置失败,实验里没有任何收益。
下面是一份最小可用的yarn-site.xml片段:
<configuration> <property> <name>yarn.nodemanager.aux-services</name> <value>mapreduce_shuffle</value> </property> <property> <name>yarn.nodemanager.resource.memory-mb</name> <value>4096</value> </property> <property> <name>yarn.scheduler.maximum-allocation-mb</name> <value>3072</value> </property> </configuration>yarn.nodemanager.aux-services是新手最容易漏掉的一项。没有它,作业可以正常提交,ResourceManager 也能看到任务,但 NodeManager 不知道如何给 MapReduce 提供 shuffle 服务,容器启动后立刻失败,表现为作业反复重试。resource.memory-mb是 NodeManager 能使用的物理内存总量,伪分布式下不要超过机器实际可用内存;maximum-allocation-mb限制了单个容器最大内存。HDFS 端配置相对简单,hdfs-site.xml里加一个dfs.replication=1即可。配完之后用start-dfs.sh和start-yarn.sh启动,再用jps确认存在 NameNode、DataNode、ResourceManager、NodeManager 四个进程,缺一个都先别急着提交作业。
3.3 WordCount 完整 Java 代码:类型、计数器与提交参数
这里给出的是一个可以直接编译提交的版本,代码里保留了完整的包名和导入,方便在本地 IDE 里直接调试。
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.LongWritable; 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 { // key: 行偏移量, value: 一行文本 public static class TokenizerMapper extends Mapper<LongWritable, Text, Text, IntWritable> { private final static IntWritable one = new IntWritable(1); private Text word = new Text(); @Override protected void map(LongWritable 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> { @Override protected void reduce(Text key, Iterable<IntWritable> values, Context context) throws IOException, InterruptedException { int sum = 0; for (IntWritable val : values) { sum += val.get(); } context.write(key, new IntWritable(sum)); } } public static void main(String[] args) throws Exception { if (args.length != 2) { System.err.println("Usage: WordCount <input path> <output path>"); System.exit(-1); } 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方法里StringTokenizer按空白字符切分文本,它不会把标点符号单独剔除,所以hello,和hello会被当作两个单词。这是实验报告里可以展开讨论的一个点:如果想做更干净的词频统计,要用正则[\\p{Punct}\\s]+做切分。reduce 方法里的sum += val.get()是对迭代器逐条取出整数值累加。注意前面提到的坑:不要把这个迭代器里的val直接放进 List,因为取出来的是同一个对象的引用。
编译时不要自己一个个拼 classpath,直接用 Hadoop 提供的命令展开依赖。以下是一套完整的编译、打包、准备输入数据、提交作业的命令:
cd ~/workspace/wordcount mkdir -p classes javac -classpath "$(hadoop classpath)" -d classes WordCount.java jar -cvf wordcount.jar -C classes . hdfs dfs -mkdir -p /user/student/input hdfs dfs -put words.txt /user/student/input/words.txt hdfs dfs -cat /user/student/input/words.txt hadoop jar wordcount.jar WordCount /user/student/input /user/student/outputhadoop classpath会输出当前环境下所有依赖 JAR 的完整路径,javac用-classpath指定它之后,代码里 import 的所有 Hadoop 类都能被找到。-d classes表示把编译产物放到 classes 目录,最后用jar -cvf打成普通 JAR 包。FileInputFormat.addInputPath和FileOutputFormat.setOutputPath都用 HDFS 路径,这里 /user/student/input 是 HDFS 上的目录。启动前先hdfs dfs -cat确认文本确实上传成功,省得提交后才发现输入目录为空。
4. 提交实验作业并验证结果:yarn 状态查询与 HDFS 输出文件的核对方法
提交作业只是开始,能解释作业状态、能从日志里找到失败原因、能从输出文件结构说出 job 是否成功,才是实验报告里更值钱的部分。这一章按实际操作的顺序展开:先提交、再查状态、最后核对结果并整理记录。
4.1 数据准备与作业提交:hadoop jar 命令背后的路径约束
hadoop jar wordcount.jar WordCount /user/student/input /user/student/output这行命令里,第一个参数是本地文件系统上的 JAR 包路径,后两个参数是 HDFS 路径。很多人第一步就弄混:输入文件必须已经存放在 HDFS 上,JAR 包则可以在 Linux 本地任意目录。hadoop jar 会把 JAR 上传到 HDFS 的/.staging临时目录,再分发给各 NodeManager,所以提交作业的机器上只要存在这个 JAR 文件即可。
输出路径有一个硬性约定:必须不存在。Hadoop 出于安全考虑不会覆盖已有输出目录。第一次跑完想重跑,必须先把上次的输出删除:
hdfs dfs -rm -r /user/student/output如果你想保留每次结果,可以给输出路径带上时间戳,比如/user/student/output_20240520。这样既能保留实验过程数据,也方便对比不同数据量下的耗时。
作业提交后,终端会进入进度等待状态。你按 Ctrl+C 只会取消本地等待,不会杀掉 YARN 上的作业,作业会在后台继续跑。想随时看任务状态,要另开一个终端用yarn application -list查。对实验报告来说,保留 JOB ID 比保留终端滚动日志更重要,后面所有排查都要靠它。
4.2 从日志验证任务状态:yarn application 与 container 日志的读法
作业跑起来后,状态查询和日志读取是两套命令。状态查询用yarn application,任务级日志用yarn logs。我实验时固定会执行下面三条命令:
yarn application -list yarn application -status application_1716000000000_0001 yarn logs -applicationId application_1716000000000_0001第一条看当前有哪些作业;第二条看指定作业的运行状态、申请容器数量、当前阶段;第三条拉取该作业所有 container 的原始日志。三条命令的输出各有侧重:application -status里最重要的字段是Tracking URL和State,State 为 FINISHED 且 FinalStatus 为 SUCCEEDED 才算真正跑通;yarn logs默认会把 stdout、stderr、syslog 全部输出,信息量大,初看去像一堆乱码。
更高效的读法是先定位失败任务编号。在application -status输出里找到Map-Reduce Framework相关的计数器,或者在 Web UI 上点开失败任务,拿到container_xxxxx编号,然后执行:
yarn logs -applicationId <app_id> -containerId <container_id> | grep -A 20 Exception这里-A 20是让 grep 把异常堆栈之后的 20 行一起打出来,通常一眼就能看到Caused by。实践经验是:Hadoop 的完整堆栈永远在 syslog 里,不要只看 stdout。stdout 主要打印 JVM 的输出,异常根本不会写在那里。
4.3 输出结果核对:part-r-00000、_SUCCESS 和实验记录表
作业显示成功之后,不要直接关掉终端,还要到 HDFS 上核对输出目录的结构。正常成功后,输出目录里会有两个关键标记:一个part-r-00000文件,一个_SUCCESS空文件。
hdfs dfs -ls /user/student/output hdfs dfs -cat /user/student/output/part-r-00000part-r-00000是 reduce 阶段生成的输出文件,文件名里的r表示来自 reducer,00000是 reduce 任务编号。没有_SUCCESS文件说明作业没有正常结束,哪怕有这个前缀文件也可能是残缺结果。hdfs dfs -cat输出的每行格式是word加制表符加计数,例如hadoop 42,如果文本量较大可以用hdfs dfs -tail只看末尾部分。
把下面这张表贴进实验报告,比单纯截图更像一份完整的跑通记录。每次实验我都会填一次,字段包括作业 ID、输入数据量、Map 任务数、Reduce 任务数、耗时、输出行数。做横向对比时还能看出数据规模变大后耗时增长是否接近线性。
| 作业 ID | 输入文件大小 | Map 任务数 | Reduce 任务数 | 总耗时 | 输出结果行数 |
|---|---|---|---|---|---|
| application_1716000000000_0001 | 256 KB | 1 | 1 | 42 秒 | 4 |
| application_1716000000000_0002 | 2.1 MB | 2 | 1 | 51 秒 | 63 |
表格里还需要注明环境参数:Hadoop 版本、JDK 版本、伪分布式内存限制。这些信息决定了你的实验是否可复现,也是报告评分时容易拿分的点。
5. MapReduce 初级编程避坑指南:提交作业阶段最常见的五个故障
这一章写的是我自己做实验时反复踩过的坑,每条都按“现象、原因、解决”来呈现。如果你在某个地方卡住,直接对照当前现象看对应章节,大部分问题几分钟内能定位。
5.1 作业一直 ACCEPTED 却不 RUNNING:NodeManager 缺失导致的任务假死
现象:yarn application -list能看到作业,State 一直停留在 ACCEPTED,等十几分钟都不进入 RUNNING。
原因:ResourceManager 收到作业后,需要找一台有 NodeManager 的节点来运行容器。如果你只启动了 HDFS,没有启动 YARN,或者yarn-site.xml里漏掉了mapreduce_shuffle,NodeManager 虽然活着但无法为 MapReduce 任务提供辅助服务,容器反复启动失败。伪分布式常见的另一个原因是机器内存不足,NodeManager 申请不到足够内存,Container 一直处于等待分配状态。
解决:先执行jps,确认 ResourceManager 和 NodeManager 进程都在。再用yarn node -list -all查看节点状态,如果节点状态是 UNHEALTHY,基本就是内存耗尽。检查内存参数是否小于机器实际可用内存,然后重启 yarn:stop-yarn.sh再start-yarn.sh。
5.2 报错 Type mismatch in value from map:类型设置前后不一致,看堆栈最高效
现象:Map 阶段全部成功,Reduce 阶段开始拉取数据时直接报错,堆栈里有java.io.IOException: Type mismatch in value from map: expected org.apache.hadoop.io.IntWritable, received org.apache.hadoop.io.Text。
原因:Mapper 输出的 value 类型是 Text,但 job 里设置的 map 输出 value 类型是 IntWritable,或者反过来。Reduce 在拉取 map 输出时,会按 job 配置的类型做反序列化,类型对不上就抛异常。这个错误常常出现在复制他人代码、只改了 Mapper 逻辑却没更新 setMapOutputValueClass 的场景。
解决:打开 Mapper 类声明的泛型参数,逐个确认三个类型的实际值,再检查 main 方法里四个 set 方法是否一致。我自己的经验是,把所有 set 方法集中写在一起,不要拆到不同方法里,检查起来最快。
5.3 FileAlreadyExistsException:输出目录不是自动覆盖而是报错
现象:第二次提交同一个作业时,作业提交后立刻报org.apache.hadoop.mapreduce.lib.output.FileAlreadyExistsException: Output directory ... already exists。
原因:Hadoop 不覆盖已存在输出目录,这是保护机制,防止误删历史结果。很多从传统文件系统转过来的同学默认覆盖写入,但 HDFS 上必须手动处理。
解决:提交前删除旧目录,或者把输出路径改成带时间戳的新路径。批量实验时我还会写一个很小的 shell 片段:
OUTPUT=/user/student/output hdfs dfs -test -d ${OUTPUT} if [ $? -eq 0 ]; then hdfs dfs -rm -r ${OUTPUT} fi这样每次提交前自动清理,不会因为漏删而浪费时间。
5.4 中文内容计数错乱或乱码:输入编码不一致导致的分词异常
现象:输入文件是中文文章,跑出来的词频结果里有大量乱码,或者一个完整中文词被切成单字,计数完全不对。
原因:WordCount 默认按空格和标点切分,中文文本之间没有空格,整句会被当成一个长串,而且如果文件以 GBK 编码保存,Hadoop 默认按 UTF-8 读取,中文部分全部变成乱码。这不是 MapReduce 框架的问题,是输入数据编码和切分策略的问题。
解决:先把文本统一转为 UTF-8:
iconv -f gbk -t utf-8 words.txt > words_utf8.txt hdfs dfs -put -f words_utf8.txt /user/student/input/words.txt如果实验要求统计中文词语而不是单字,需要对切分逻辑做升级,引入分词器,而不是在 WordCount 里硬写。避坑的关键是:上传 HDFS 前先file words.txt看一下编码。
5.5 想看异常却找不到完整堆栈:syslog 才是唯一的完整日志
现象:作业失败了,yarn logs打出来一堆内容,但找不到Caused by或关键 Exception,只能看到进程退出码。
原因:Hadoop 容器日志分三层,stdout 只记录标准输出,stderr 只记录 JVM 崩溃信息,完整的 Java 异常在 syslog 里。很多人盯着前两个文件找线索,当然一无所获。
解决:先用yarn logs -applicationId <app_id> | grep "Caused by",再用grep -A 20看上下文。如果日志量太大,先找到失败 container 的 ID,再单独拉某个 container 的日志。这个习惯能帮你把排错时间从半小时压缩到两分钟,是 MapReduce 编程实践里最值得养成的技能之一。
6. 把 WordCount 升级成自己的实验工具:计数器、Combiner 与本地复现的验证技巧
初级实践做完后,如果还停留在“会跑一个模板”的程度,实验报告的含金量会打折。我会再做两件事:给代码加一个自定义计数器,统计非法输入行数;然后把同样的逻辑切到本地模式跑一遍,验证代码不依赖集群环境。
自定义计数器的代码非常轻量,定义一个枚举类型,然后在 map 里对异常行累加:
public static enum LineCounters { BAD_LINES } // map 方法里对未通过校验的行执行: context.getCounter(LineCounters.BAD_LINES).increment(1L);作业结束后,不需要解析输出文件,直接在终端查:
yarn application -status <app_id> | grep BAD_LINES计数器的好处是它由框架统一汇总,比在输出文件里统计更可靠,也不会受到 reduce 数据倾斜影响。实验报告里可以拿这个数据论证输入质量,比干巴巴的词频表更有说服力。
第二个技巧是本地模式复现。把main里的 HDFS 路径改成本地路径,不启动任何 Hadoop 进程,直接用 IDEA 运行 main 方法:
inputPath=file:///home/student/workspace/input outputPath=file:///home/student/workspace/output这能瞬间启动一个单进程 MapReduce,适合验证代码逻辑和调试切分规则。等本地跑通了,再把路径换回 HDFS 路径提交到集群,这样能把环境问题和代码问题分开。这两年做实验时我都按这个顺序来:先本地复现,再集群提交,最后用计数器核对结果。本地模式跑出来的结果和集群保持一致,说明代码本身没问题;如果集群上失败,就集中查配置、权限、容器内存这些外围因素,不再怀疑业务逻辑。这套流程帮我避免了很多次无效返工,希望帮到你。
本文还有配套的精品资源,点击获取