简介:《大数据技术原理与应用》实验报告四是一份围绕MapReduce初级编程实践的完整实验文档,面向计算机科学专业学生、大数据技术初学者及对分布式数据处理感兴趣的开发者。内容以三个典型任务为主线:编程实现文件合并与去重、编写程序对输入文件排序、对给定表格进行信息挖掘,每一步均给出可运行的Java代码、HDFS路径配置、运行命令与结果截图,并基于VMWare虚拟机、Ubuntu、JDK1.8及Hadoop-3.1.3环境展开,帮助读者快速打通从环境搭建到代码调试的全流程。文档还梳理了实验过程中遇到的Hadoop配置错误、权限不足、数据倾斜等问题的解决过程,并附有心得总结,适合作为课程实验报告模板或MapReduce自学笔记。压缩包内包含1个docx文档,容量10.48MB,已有133人学习下载,便于按需查阅和反复实践。
1. 实验报告四的 MapReduce 实践:一份能少走两天弯路的资源
交实验周前一天晚上,我在本地 IDE 里把 WordCount 写得滚瓜烂熟,自以为 MapReduce 原理已经吃透,结果把 jar 丢到虚拟机里的 Hadoop 集群上,二十分钟过去进度还是 0%。查了半天才发现是 YARN 的堆内存参数和伪分布式环境不匹配。这种经历在 MapReduce 入门阶段太常见了。这份《大数据技术原理与应用》实验报告四,讲的就是 MapReduce 初级编程实践——从环境搭建、WordCount 标准实现,到自定义序列化 Bean、分组排序和多 Job 串联,把初级阶段的必过点压缩成一份可复现的操作记录。适合正在上大数据课程的同学,也适合第一次在 VMWare 里搭 Hadoop 环境的从业者照着走一遍。
2. 实验环境选型与配置:VMWare 里的伪分布式 Hadoop 怎么立起来
2.1 为什么选伪分布式而不是多节点集群
MapReduce 初级编程实践里,环境选择直接决定你后面踩坑的密度。很多初学者一上来就想搭三台节点的集群,结果光是同步配置、开防火墙、配 SSH 就花掉一周,实验本身反而没时间认真写。这份实验报告给的是 VMWare 虚拟机里的伪分布式部署:一台 Linux 虚拟机同时跑 NameNode、DataNode、ResourceManager 和 NodeManager。这样做的理由是,MapReduce 的分析逻辑在单节点上跑通和集群上跑通没有本质区别,你要掌握的 Mapper、Reducer、Partitioner、Counter 这些概念,不会因为节点数量而改变。
伪分布式环境对硬件要求也低很多。虚拟机分配 2GB~4GB 内存、20GB 磁盘就够跑基础实验,前提是 YARN 的内存参数必须跟着调,不然任务会频繁被杀。如果你机器内存只有 8GB,给虚拟机 4GB 已经是极限,YARN 默认的 8GB 内存配置就会让 Container 直接起不来,这是新手最容易翻车的地方。
2.2 伪分布式 Hadoop 初始化与关键配置文件
实验环境一般由两部分组成:宿主机负责写代码和打包,虚拟机负责跑 Hadoop。下面是一份能跑通 MapReduce 的基础配置,核心逻辑是让 Hadoop 各组件通过内网地址互相通信。
<!-- core-site.xml --> <property> <name>fs.defaultFS</name> <value>hdfs://bigdata:9000</value> </property> <!-- hdfs-site.xml --> <property> <name>dfs.replication</name> <value>1</value> </property> <!-- mapred-site.xml --> <property> <name>mapreduce.framework.name</name> <value>yarn</value> </property> <!-- yarn-site.xml --> <property> <name>yarn.nodemanager.resource.memory-mb</name> <value>2048</value> </property> <property> <name>yarn.scheduler.minimum-allocation-mb</name> <value>256</value> </property> <property> <name>yarn.scheduler.maximum-allocation-mb</name> <value>2048</value> </property>这里 fs.defaultFS 指定了 NameNode 的地址,bigdata 是虚拟机的 hostname,9000 是 RPC 通信端口。dfs.replication 设为 1 是因为伪分布式只有一台 DataNode,默认的 3 份副本策略会报副本不足的告警。mapreduce.framework.name 必须改成 yarn,否则作业默认走 local 模式,日志里根本看不到 Container 调度过程,你也就没法通过日志判断作业死活。yarn.nodemanager.resource.memory-mb 前面说过,要小于等于虚拟机实际内存,否则 NodeManager 会因资源不足反复重启。
初始化顺序我一般固定走一遍:
hdfs namenode -format start-dfs.sh start-yarn.sh jpsjps 输出里应该有 NameNode、DataNode、ResourceManager、NodeManager 四个进程,缺哪个就单独看对应日志。格式化的坑在于:格式化操作会清空 NameNode 的元数据,如果你改过 core-site.xml 里的路径,一定要先确认旧目录不存在,否则 NameNode 会起不来,报错信息又长又迷惑。
主配置文件与实际生效值的对应关系,整理成下表方便核对:
| 配置项 | 建议值 | 说明 | 常见误配 |
|---|---|---|---|
| fs.defaultFS | hdfs://bigdata:9000 | 默认文件系统地址 | 配成宿主机 IP 导致内外不一致 |
| dfs.replication | 1 | 副本数 | 新手忘改,日志出现副本告警 |
| mapreduce.framework.name | yarn | 作业调度框架 | 不配则走 local 模式 |
| yarn.nodemanager.resource.memory-mb | 2048 | NodeManager 可用内存 | 超过虚拟机内存导致容器被杀 |
| mapreduce.job.reduces | 1~2 | Reducer 数量,实验来源可看实际需求 | 调太多产生大量小文件,调太少数据倾斜 |
2.3 从宿主机到虚拟机的 jar 包传输方式
环境起来之后,下一个问题是代码怎么进去。常见的做法有三种:用共享文件夹、用 scp 命令、用 sftp 工具。如果虚拟机里没有图形界面,直接在宿主机 IDE 里把项目打成一个可执行 jar,再用 scp 传过去最省事。
scp wc.jar bigdata:/home/student/exp4/scp 适合小 jar 包,几十 KB 的 WordCount 几秒钟就传完。如果你后续要做 HBase 相关的实验,jar 会变大,也可以直接在虚拟机里配 IDE,但会占用额外内存,初级实验不推荐。传输完成后,注意确认 Linux 用户的写权限:HDFS 的 /user 目录属主是启动 Hadoop 的用户,如果当前用户不是它,就得用 hdfs dfs -chown 或 -chmod 调整,否则作业写入结果时会报 AccessControlException。
3. WordCount:把 MapReduce 跑通的那些细节
3.1 MapReduce 执行流程与 WordCount 的三个类
WordCount 是 MapReduce 的 Hello World,它的价值在于一次跑通就能理解整套流程:输入文件被 InputFormat 切片后,每行交给 Mapper 处理,Mapper 输出的键值对经过分区、排序、合并后按 key 分组传给 Reducer,Reducer 的结果再写回 HDFS。这个过程中,Map 端输出的 kv 是中间数据,Reduce 端输出的是最终结果,弄清楚这条链路比看懂任何源码都有用。
以下是一个可直接编译运行的 WordCount 实现,三个类都写在同一个文件里,方便实验报告里对照:
import java.io.IOException; 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 { 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 { // 每行内容按空白字符切分,逐词输出 <word, 1> String[] tokens = value.toString().split("\\s+"); for (String token : tokens) { if (token.length() > 0) { word.set(token); context.write(word, one); } } } } public static class IntSumReducer extends Reducer<Text, IntWritable, Text, IntWritable> { private IntWritable result = new 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(); } 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 是 LongWritable 类型的行偏移量,这是默认 TextInputFormat 决定的,不要改成其他类型。第二,map 方法里不要自行维护全局变量,每个 map 任务可能会被多个 JVM 并发执行,你无法预期变量状态。第三,Reducer 的 values 是相同 key 的迭代器,对迭代器做累加是标准操作,但不要在循环外部持有迭代器,它可能已经失效。第四,main 方法里设置了 CombinerClass,它用的是 Reducer 类,因为 WordCount 的 reduce 操作满足结合律,combiner 在 map 端提前做一次本地聚合,能大幅减少 shuffle 数据量。
3.2 提交作业与查看结果
代码编译打包之后,提交命令如下:
hadoop jar wc.jar WordCount /exp4/input /exp4/output这条命令的实际效果是:Hadoop 把 wc.jar 分发到各个节点,通过反射找到 main 方法里配置的 WordCount 类,然后启动 YARN 作业。input 和 output 都是 HDFS 上的路径,不是 Linux 本地路径。output 路径必须不存在,否则作业会直接报错退出,这是 Hadoop 防止覆盖数据的安全策略。跑完后用下面的命令查看结果:
hdfs dfs -cat /exp4/output/part-r-00000 | head -50如果产生多个 reducer,输出文件会是 part-r-00000、part-r-00001 等不同编号,每个文件对应一个 reducer 的输出。如果输出目录里还有 _SUCCESS 文件,说明作业正常结束。看到这个文件之前,无论控制台怎么显示,都不要认为作业已经成功。
3.3 控制并行度的常用参数
数据量变大后,你会发现 WordCount 跑得越来越慢,这时候该调并行度了,而不是换电脑。Map 端的并行度由输入分片数量决定,每个文件块默认 128MB,一个分片对应一个 map 任务。Reduce 端并行度由 mapreduce.job.reduces 控制,可以在提交命令里用 -D 参数临时指定:
hadoop jar wc.jar WordCount -D mapreduce.job.reduces=3 /exp4/input /exp4/outputreduces 的数量不是越大越好。输出结果会按 key 的哈希值分到不同 reducer,若 reducer 数量远大于机器核心数,大量 JVM 启动开销反而会让作业变慢。初级实验里线程数设为 1~2 即可,观察输出文件的变化规律,再决定要不要调。这里也顺带提醒,连同 CombinerClass 的使用,WordCount 是理解“Map 端预聚合”最好的载体。
4. 自定义序列化和分组排序:MapReduce 初级进阶的两道坎
4.1 为什么需要自定义 Writable
WordCount 用到的 Text、IntWritable 是 Hadoop 内置类型,但真实实验数据往往是一整行日志,比如“手机号 上行流量 下行流量”。你要在 map 阶段输出一个包含多个字段的对象,不能直接用 Java 的 Serializable,因为 Hadoop 对序列化有特殊要求:对象需要能快速序列化、快速反序列化,并且序列化后体积尽量小。自定义类实现 Writable 接口,就是要满足这套要求。
下面是一个流量统计用的 FlowBean 实现,它同时实现了 Writable 和 Comparable 两个接口,既能在 MapReduce 里传输,又能在排序阶段被比较:
import org.apache.hadoop.io.Writable; import java.io.DataInput; import java.io.DataOutput; import java.io.IOException; public class FlowBean implements Writable, Comparable<FlowBean> { private long upFlow; // 上行流量 private long downFlow; // 下行流量 private long sumFlow; // 总流量 public FlowBean() { super(); } public FlowBean(long upFlow, long downFlow) { this.upFlow = upFlow; this.downFlow = downFlow; this.sumFlow = upFlow + downFlow; } @Override public void write(DataOutput out) throws IOException { out.writeLong(upFlow); out.writeLong(downFlow); out.writeLong(sumFlow); } @Override public void readFields(DataInput in) throws IOException { this.upFlow = in.readLong(); this.downFlow = in.readLong(); this.sumFlow = in.readLong(); } @Override public int compareTo(FlowBean o) { // 按总流量降序排列 return Long.compare(o.sumFlow, this.sumFlow); } @Override public String toString() { return upFlow + "\t" + downFlow + "\t" + sumFlow; } }这个类有三个要点。第一,必须有一个无参构造方法,因为 Hadoop 反序列化时要通过反射创建对象实例。第二,write 和 readFields 的字段顺序必须完全一致,先写哪个就先读哪个,顺序错位会造成数据错乱且不容易定位。第三,compareTo 里返回的是降序排序,这是通过交换比较顺序实现的,很多初学者写成 this.sumFlow - o.sumFlow,升序降序搞混了,结果输出的结果是反的,还以为是 Hadoop 的问题。
4.2 自定义 Partitioner 控制分区
当你要把不同手机号段的流量数据分到不同 reduce 任务时,默认的哈希分区满足不了需求。默认分区器对 key 做 hashCode 取模,结果随机性太强,你没法控制某类数据进入某个 reduce。自定义 Partitioner 可以精确控制分区逻辑:
import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.Partitioner; public class PhonePartitioner extends Partitioner<Text, FlowBean> { @Override public int getPartition(Text key, FlowBean value, int numPartitions) { String phone = key.toString(); String prefix = phone.substring(0, 3); if ("136".equals(prefix)) { return 0; } else if ("137".equals(prefix)) { return 1; } else if ("138".equals(prefix)) { return 2; } else if ("139".equals(prefix)) { return 3; } else { return 4; } } }这里用手机号前三位做条件分区,返回的分区号范围是 0 到 numPartitions-1。使用时要确保 numPartitions 与自定义分区器返回的最大分区号一致,否则会抛 IllegalStateException。设置方式是在 Driver 里调用 job.setPartitionerClass(PhonePartitioner.class),同时 job.setNumReduceTasks(5)。记住一个原则:分区器决定了“哪个 key 去哪台 reduce”,分组比较器决定了“到了同一台 reduce 后,哪些 key 算一组”,两者不要混为一谈。
4.3 GroupingComparator 实现分组排序
分组排序的常见场景是“每个号段中找出流量最高的手机号”。Reducer 接收的数据是经过分区、排序后的,默认把完全相同的 key 分为一组,但我们需要把相同号段的不同手机号分为一组,这时就要自定义 GroupingComparator:
import org.apache.hadoop.io.WritableComparable; import org.apache.hadoop.io.WritableComparator; public class PhoneGroupComparator extends WritableComparator { protected PhoneGroupComparator() { super(Text.class, true); } @Override public int compare(WritableComparable a, WritableComparable b) { // 只按手机号前三位比较,前三位相同视为同一组 String aPrefix = a.toString().substring(0, 3); String bPrefix = b.toString().substring(0, 3); return aPrefix.compareTo(bPrefix); } }这个比较器的 compare 方法只在分组时生效,不会改变排序顺序。排序阶段用的是 FlowBean 的 compareTo,分组阶段用的是这个自定义比较器,两者作用阶段不同。设置了自定义分组比较器后,Reducer 的 values 里并不是同一个 key 的所有数据,而是同一个号段下所有手机号的数据,你可以在同一组内遍历找出最大流量值。如果只设置了分区器而不设置分组比较器,reduce 会为每个手机号分别调用一次,需求就无法实现。
4.4 多 Job 串联的常见做法
复杂的实验需求经常要拆成两个 Job 来完成,比如第一个 Job 清洗数据,第二个 Job 做统计。串联方式很简单:第一个 Job 的输出路径作为第二个 Job 的输入路径。
Job job1 = Job.getInstance(conf, "clean"); // 配置 job1 的 Mapper、Reducer、输入输出路径 job1.waitForCompletion(true); Job job2 = Job.getInstance(conf, "statistics"); FileInputFormat.addInputPath(job2, new Path("/exp4/output1")); FileOutputFormat.setOutputPath(job2, new Path("/exp4/output2")); job2.waitForCompletion(true);这里 job1 的输出目录就成了 job2 的输入目录,中间结果保存在 HDFS 上,不需要落到本地文件系统。要注意的是,job1 的 reducer 数量直接影响中间文件的数量,如果 job2 有大量小文件输入,输入效率会明显下降。我一般会在 job1 设置 reducer 数量跟分区需求匹配,同时考虑后续 job 的读取效率。如果你后续想把结果写入 HBase,可以把 job2 的输出类换成 TableOutputFormat,但初级实验用文本输出检查数据更直接,排查问题时也更方便。
5. 运行排查:五个高频翻车点与应对解法
5.1 YARN 容器反复被杀,任务卡在 0%
现象:作业提交后,进度一直停在 map 0% reduce 0%,查看日志发现 NodeManager 反复报 Container killed。
原因:虚拟机内存只有 2GB,但 YARN 默认的 nodemanager.resource.memory-mb 是 8GB,注册资源时就已经超过实际内存,容器一启动物理内存就超标,直接被节点管理器杀掉。
解决:把 yarn-site.xml 里的 yarn.nodemanager.resource.memory-mb 改到 2048,同时把 yarn.scheduler.minimum-allocation-mb 和 maximum-allocation-mb 分别改成 256 和 2048。改完重启 YARN,确认 NodeManager 不再重启,再看任务进度。
这个坑我踩过不止一次,后来每开一台新虚拟机,第一件事就是检查 YARN 内存参数和 free -m 输出是否匹配,不再依赖默认值。
5.2 HDFS 写入权限不足,报 AccessControlException
现象:作业能提交,map 也执行了,但 reduce 写结果时抛出 org.apache.hadoop.security.AccessControlException: Permission denied。
原因:启动 Hadoop 的用户是 hadoop,而当前登录用户在 HDFS 上没有 /exp4/output 相关目录的写权限。HDFS 权限模型和 Linux 类似,它不是摆设。
解决:先用 hdfs dfs -ls /exp4 确认目录属主,再把输出目录的属主改成当前用户,或者直接把整个实验目录赋予写权限。命令如下:
hdfs dfs -chown -R student:student /exp4 hdfs dfs -chmod -R 755 /exp4如果只是临时测试,也可以直接删掉旧输出路径,换一个新的未创建路径,让作业自己创建目录,作业用户的写权限只取决于父目录是否可写。
5.3 作业报 ClassNotFoundException
现象:编译不报错,但提交作业时报 java.lang.ClassNotFoundException,指向你自己写的某个类。
原因:提交命令里指定的主类名和 jar 包里的实际类名不匹配,或者主类没有通过 setJarByClass 指定当前类。更隐蔽的情况是,项目里引用了第三方依赖,打包时没打进去。
解决:提交命令用全限定类名,jar 包内用 jar tf 检查类文件是否存在。如果确实引用了第三方库,打包时打成 fat jar,把所有依赖打进同一个 jar。很多 IDE 默认打包不包含依赖,这是新手反复踩的坑。
jar tf wc.jar | grep WordCount这条命令查到的路径应该和你提交命令里写的类名一致。
5.4 输出目录已存在,作业秒退
现象:第二次运行同一命令,控制台立刻报 FileAlreadyExistsException,并提示输出目录已存在。
原因:Hadoop 出于数据安全考虑,不允许作业覆盖已有输出目录。这不是真的“异常”,而是设计策略。很多人第一次跑通后不改路径,第二次就卡在这里。
解决:每次运行前手动删除已有输出目录,或者通过 shell 脚本取当前时间戳动态生成输出路径。
hdfs dfs -rm -r /exp4/output建议在 Driver 里不写死输出路径,而是从 args 读取,这样每次运行不必重新编译代码。
5.5 日志显示成功,输出文件却是空的
现象:控制台输出 Job complete successfully,_SUCCESS 文件也在,但打开 part-r-00000 没有任何数据。
原因:Reducer 的 reduce 方法没有把 context.write 写在正确的循环位置,或者 map 端输出的 key 和 reduce 端使用的 key 类型不一致,导致 reduce 方法被跳过。另一种可能是输入文件本来就为空,HDFS 上落了一个空文件。
解决:检查输入文件大小,hdfs dfs -du -h /exp4/input 确认有数据。再查看 Counter 输出的 map input records 数量,如果是 0,问题出在输入路径;如果 map 有记录但是 reduce output records 是 0,问题出在 reduce 逻辑。Counter 是定位这种问题的第一工具,不要靠猜。
6. 验证方法:从日志、计数器和 HDFS 文件中确认作业真的成功
作业跑完,最忌讳只看控制台黑框里那一行“SUCCESS”。我后来不管实验多急,验证都固定走三件套。第一,看 HDFS 输出文件,cat 出来的数据量和内容至少要符合对输入数据的常识判断。比如输入有 1000 行日志,WordCount 统计结果至少有几百行,如果只有几行,那肯定有问题。第二,看 Counter 里的关键指标,map input records、reduce input records、reduce output records 三个数的变化趋势能直接反映数据处理过程。第三,看应用程序日志,用 yarn application -status 或 yarn logs 查看 ApplicationMaster 的完整输出,而不是只看控制台。如果日志里出现 Info 级别的 Container 启动记录,说明作业确实走了 YARN 调度流程,而不是被某个异常吞掉。
有一个小技巧是给 Driver 增加一个日志输出。在 main 方法最后把 job.getCounters() 里感兴趣的计数器打印出来,这样每次跑完不用去翻网页界面,直接在终端就能确认任务是否达到预期。我习惯把 map input records 和 reduce output records 的比例关系也打印出来:如果输入 1 万条,输出也是 1 万条,说明 reduce 几乎是逐条输出,可能存在数据未聚合的问题。从那以后我每次提交作业前,都会强制把这三件套过一遍,确认无误才在实验报告上写“运行成功”。希望帮到你。
本文还有配套的精品资源,点击获取