MapReduce初级编程实践:三个Java程序详解合并去重与全局排序
2026/9/24 12:32:34 网站建设 项目流程

简介:这是一份面向高校大数据课程学习者的实验报告资源,基于林子雨《大数据原理与技术》第三版第五章内容,完整呈现MapReduce初级编程实践过程。报告以文件合并与去重为实战目标,给出Hadoop 3.2.2环境下的可运行Java代码,涵盖Map、Reduce阶段的关键实现、作业配置与提交方式,并附有输入输出样例供对照验证,适合正在完成同类实验或复习MapReduce原理的读者参考。资源包共1个文件,为docx格式文档,大小约1.28MB,排版清晰、内容完整,可直接作为实验报告撰写与代码调试的参照。已有14000余人浏览学习,是同类资源中较受关注的实验案例。通过阅读该报告,读者能够快速理解如何利用MapReduce解决去重与合并问题,掌握编写、配置和运行MapReduce作业的基本思路,为后续分布式数据处理实践打下扎实基础。

1. 大数据实验5 MapReduce 初级编程实践:三个跑通的 Java 程序一次讲清

大数据实验五的 MapReduce 初级编程实践,给的是三个已经跑通的 Hadoop Java 程序:文件合并与去重、多文件全局排序、单表关联挖掘祖孙关系,配套信息是林子雨《大数据技术原理与应用》实验5,Hadoop 版本 3.2.2。它适合三类人:要交实验报告的学生、想在本地伪分布式环境跑通第一个 Job 的从业者,以及想搞明白自定义 Partitioner 到底解决什么问题的面试准备者。直接说结论:合并去重那个作业 20 行核心代码就能跑通;排序那个作业如果不写自定义分区,输出大概率不是全局有序,卡在这儿的同学不在少数。

2. 环境准备与提交流程:Hadoop 3.2.2 伪分布式下的编译、打包与运行

拿到别人的实验源码,第一件事不是看代码,而是先把环境对齐。原报告写的是 Linux(建议 Ubuntu 16.04)+ Hadoop 3.2.2,这个组合在实际复现时有个前提:JDK 必须用 8。Hadoop 3.x 对 JDK 版本有硬性要求,JDK 11 在某些发行版上能启动,但跑 MapReduce 作业时偶尔会冒出奇怪的类加载异常,保守起见直接 JDK 8。

2.1 伪分布式环境怎么配:四个 XML 和一条启动链

这三个实验的数据量是 KB 级别,完全没必要搭三台机器的集群。伪分布式模式就是单台机器上同时跑 NameNode、DataNode、ResourceManager、NodeManager,四个进程各司其职,既能完整走一遍 HDFS 读写和 YARN 调度,又方便直接看日志定位问题。原代码里conf.set("fs.defaultFS","hdfs://localhost:9000")就是伪分布式下的 NameNode 地址,说明实验默认你用的是单机模式。

核心配置集中在四个文件:core-site.xml 设置默认文件系统地址,hdfs-site.xml 设置副本数,mapred-site.xml 指定调度框架,yarn-site.xml 配置 NodeManager 的辅助服务。实验场景下副本数设 1 就够了,设 3 在伪分布式里会产生大量副本等待日志。

配置文件配置项实验推荐值作用
core-site.xmlfs.defaultFShdfs://localhost:9000默认 HDFS 地址,和代码里 conf.set 的值对应
hdfs-site.xmldfs.replication1单机伪分布式副本数,设 3 没有意义
mapred-site.xmlmapreduce.framework.nameyarn让 MapReduce 作业跑在 YARN 上
yarn-site.xmlyarn.nodemanager.aux-servicesmapreduce_shuffleShuffle 阶段依赖的辅助服务
# 装好 JDK8 和 Hadoop 3.2.2 后,先配环境变量 export JAVA_HOME=/usr/lib/jvm/java-8-openjdk-amd64 export HADOOP_HOME=/opt/hadoop-3.2.2 export PATH=$PATH:$HADOOP_HOME/bin:$HADOOP_HOME/sbin # 首次启动前必须格式化 NameNode,注意只需执行一次 hdfs namenode -format # 启动 HDFS 和 YARN start-dfs.sh start-yarn.sh # 检查五类进程是否都在:NameNode DataNode SecondaryNameNode ResourceManager NodeManager jps # 创建实验目录,input 目录名对应代码里的 otherArgs[0] hdfs dfs -mkdir -p /user/ubuntu/input # 把本地实验数据传上去 hdfs dfs -put A.txt B.txt /user/ubuntu/input/

格式化 NameNode 是个不能手滑的操作,它会清空元数据,生产环境里重复执行等于删库。伪分布式无所谓,但养成习惯:只在首次搭建时格式化,之后重启集群只用start-dfs.shstart-yarn.sh。上传数据后可以用hdfs dfs -ls /user/ubuntu/input确认文件完整,也可以在浏览器打开localhost:9870(Hadoop 3.x 的 Web 端口,2.x 是 50070)直接看 HDFS 上的文件分布。

2.2 编译打包与提交:从 .java 到 part-r-00000

原报告用的是 Eclipse 导出 Runnable JAR 的方式,操作路径是右键项目 → Export → Runnable JAR file → 在 Launch Configuration 里选择主类。如果用的是 Maven 工程,pom.xml 里引入 hadoop-client 依赖,scope 设为 provided,然后mvn clean package。两种方式都行,重点是job.setJarByClass(Merge.class)这一行,它告诉 Hadoop 到哪个 jar 里找 Mapper 和 Reducer 类,漏掉这行作业会直接报 ClassNotFoundException。

# 提交作业,三个参数分别是 jar 包、主类名、输入目录、输出目录 hadoop jar merge.jar Merge input output # 查看运行结果,part-r-00000 是 Reducer 的输出文件 hdfs dfs -cat output/part-r-00000 # 作业失败时查日志,applicationId 从控制台输出里复制 yarn logs -applicationId application_xxx

hadoop jarhdfs dfs是两套命令,前者提交 MapReduce 作业,后者操作 HDFS 文件,新手常把这俩搞混。output 目录必须是 HDFS 上不存在的路径,这是 FileOutputFormat 的硬性规定,防止覆盖旧数据。跑完后目录里会出现两个文件:_SUCCESS标记作业成功,part-r-00000是真正的输出。如果作业失败,优先看yarn logs里的堆栈信息,比在控制台瞎猜有效得多。

3. 合并去重与全局排序:第一个 Job 和自定义 Partitioner 背后的为什么

这一章讲前两个实验:文件合并去重和全排序。这两个实验放在一起看特别合适,因为它们的 Map 阶段几乎一样简单,但 Reduce 和分区策略完全不同,对比着看能理解 Shuffle 阶段到底替你做了什么。

3.1 合并去重:Map 阶段把整行当 key,Reduce 阶段只透传一次

代码的巧妙之处在于它没有做任何"去重"操作,只是把整行文本作为 key 输出,value 全部置空。Shuffle 阶段会按 key 分组,相同 key 的所有 value 会合并成一个 list 送到同一个 Reducer。Reducer 拿到 key 后不管 values 里有几条,只输出一次 key,去重就完成了。

public static class Map extends Mapper<Object, Text, Text, Text> { private static Text text = new Text(); // 直接将输入的 value(一整行文本)复制到输出 key 上 // value 输出为空字符串,因为我们只关心“哪些行出现过” public void map(Object key, Text value, Context content) throws IOException, InterruptedException { text = value; content.write(text, new Text("")); } } public static class Reduce extends Reducer<Text, Text, Text, Text> { // 同一个 key(即同一行文本)只输出一次,重复内容自然消失 public void reduce(Text key, Iterable<Text> values, Context context) throws IOException, InterruptedException { context.write(key, new Text("")); } }

注意 job 里设置了job.setCombinerClass(Reduce.class),这个设置很关键。Combiner 是在 Map 端提前做一次本地合并,减少 Shuffle 传输的数据量。这里的 Reducer 逻辑是"每个 key 输出一次",Combiner 在 Map 端执行相同逻辑不会改变语义,所以可以安全复用。如果 Reducer 逻辑是求和求平均,直接复用 Reducer 作 Combiner 会算错结果,这是经典误用场景。

输入文件 A 和 B 的每一行是日期加字母的组合,比如20170101 x。Mapper 输出 key 就是这整行字符串,两个文件里重复的行会被 Shuffle 分到同一个 key,Reducer 只输出一次,最终得到的就是合并且去重的 C 文件。

3.2 全局排序:三个 Reduce 各排各的,自定义分区才能全局有序

第二个实验的输出要求是每行两个数字:排序位次和原始整数。如果只靠默认的 HashPartitioner,多个 Reducer 各自有序,但把多个 part-r-0000x 文件拼起来看,整体是乱序的。实验要求输出到一个文件里的数据是升序排列,所以必须自定义 Partitioner,让数值按大小区间划分到不同分区,每个分区只负责一个数值段。

默认情况下 MapReduce 会启动一个 Reducer,这个数量不足以体现分区的意义。贴出来的代码里没有显式设置 Reducer 数量,走的是默认值 1,这种情况下自定义 Partitioner 没有实际效果,但它展示了自定义分区的完整写法,值得逐行拆解。

// 自定义分区:根据数值大小划分到不同分区,保证分区之间数值范围严格分离 public static class Partition extends Partitioner<IntWritable, IntWritable> { public int getPartition(IntWritable key, IntWritable value, int num_Partition) { int Maxnumber = 65223; // 输入数据的最大边界 int bound = Maxnumber / num_Partition + 1; // 每个分区覆盖的数值范围 int keynumber = key.get(); for (int i = 0; i < num_Partition; i++) { if (keynumber < bound * (i + 1) && keynumber >= bound * i) { return i; // 返回分区编号 } } return -1; // 理论上走不到这里 } }

bound = Maxnumber / num_Partition + 1的逻辑值得注意。加 1 是为了处理整数除法的取整误差。假设 Maxnumber 是 65223,分区数是 3,bound 是 21742。那么 0 到 21741 落在分区 0,21742 到 43483 落在分区 1,43484 到 65223 落在分区 2。如果输入数据里出现负数或者超过 Maxnumber,返回 -1 会直接报错,所以这个常量必须大于等于输入最大值。

Reduce 阶段用一个全局变量line_num记录当前位次,每输出一个 key 就自增 1。由于分区之间数值范围严格分离,只要分区编号从小到大排列,整个输出的位次就是正确的全局排序。

3.3 运行与判定:part-r-00000 里看到的就是答案

作业跑完后,去 HDFS 上查看输出文件,建议用下面这个检查表逐项核对:

检查项方法预期结果
作业是否成功查看控制台或 ResourceManager UI出现Job complete字样
输出文件个数hdfs dfs -ls output_SUCCESS和一个或多个part-r-0000x
去重结果对比 A、B 文件的行数之和与输出行数输出行数 = A∪B 的唯一行数
排序结果检查输出的第二个数字必须是非递减序列
位次正确性检查输出的第一个数字从 1 开始连续递增

第一次跑通后,把 Reducer 数量改成 2 或 3 再跑一次排序实验,观察输出文件变成多个,再用hdfs dfs -getmerge合并后检查是否还是全局有序。这样做一遍,你对 Partitioner 的理解会扎实很多。

4. 单表关联与祖孙关系挖掘:左右表标志位和笛卡尔积的拼装方法

第三个实验输入是一张 child-parent 两列表,要输出 grandchild-grandparent 关系。它和普通 Join 最大的不同是:Join 的两张表是同一张表,这就是经典的"自连接"场景。MapReduce 框架没有提供直接的 Join 原语,自连接的本质是把一张表通过 Mapper 拆成左右两张逻辑表,再在 Reducer 里按 key 匹配。

4.1 自连接的拆分思路:左表右表都靠 Mapper 打出来

对于输入中的每一行child parent,Mapper 需要输出两条记录。第一条以 parent 为 key,value 里包含 child 信息,这构成了"找孙子"的左表;第二条以 child 为 key,value 里包含 parent 信息,这构成了"找爷爷"的右表。为了区分这两条记录,value 前面加了一个标志位12

输入行输出 key输出 value语义
Steven LucyLucy1+Steven+Lucy左表:Steven 是 Lucy 的孩子
Steven LucySteven2+Steven+Lucy右表:Lucy 是 Steven 的父辈
public static class Map extends Mapper<Object, Text, Text, Text> { public void map(Object key, Text value, Context context) throws IOException, InterruptedException { String line = value.toString(); int i = 0; while (line.charAt(i) != ' ') { // 按空格拆出 child 和 parent i++; } String[] values = { line.substring(0, i), line.substring(i + 1) }; if (values[0].compareTo("child") != 0) { // 跳过表头 String child_name = values[0]; String parent_name = values[1]; // 左表:以 parent 为 key,说明 “这条记录里的 child 是孙子候选” context.write(new Text(values[1]), new Text("1+" + child_name + "+" + parent_name)); // 右表:以 child 为 key,说明 “这条记录里的 parent 是爷爷候选” context.write(new Text(values[0]), new Text("2+" + child_name + "+" + parent_name)); } } }

这种"一条输入打两次"的写法是自连接的标准做法。标志位放在 value 的最前面,为的是在 Reducer 里用charAt(0)就能取到,不需要再做字符串分割。value 的格式是标志位+child+parent,加号是自定义的分隔符,这里不像表头那样受空格干扰,只要保证 Mapper 输出和 Reducer 解析用的是同一个分隔符就行。

4.2 Reducer 里做笛卡尔积:拆串、分组、全组合

Reducer 收到的是相同 key 的 value-list。以Lucy这个 key 为例,输入数据里有Steven LucyLucy Mary两行。Mapper 对第一行输出了左表记录1+Steven+Lucy,对第二行输出了右表记录2+Lucy+Mary。这两个 value 会同时落到 key=Lucy 的 reduce 组里,Reducer 把左表记录里的 child(Steven)放进 grand_child 数组,把右表记录里的 parent(Mary)放进 grand_parent 数组,然后双重循环输出Steven Mary,即 Steven 是 Mary 的孙子。

// 拆出 value-list 里的 child 和 parent,按标志位分别存数组 char relation_type = record.charAt(0); // 标志位:1 左表,2 右表 int i = 2; // 跳过 "1+" 或 "2+" while (record.charAt(i) != '+') { // 解析 child_name child_name = child_name + record.charAt(i); i++; } i = i + 1; while (i < len) { // 解析 parent_name parent_name = parent_name + record.charAt(i); i++; } // 左表:child 放进 grand_child 数组 if (relation_type == '1') { grand_child[grand_child_num] = child_name; grand_child_num++; } else { // 右表:parent 放进 grand_parent 数组 grand_parent[grand_parent_num] = parent_name; grand_parent_num++; }

这里的数组容量写死了 10,是原实验代码的一个隐患。如果某个 key 关联的记录数超过 10,会抛数组越界。我一般会改成ArrayList,或者至少把容量提到一个明显够用的值。笛卡尔积双重循环输出grand_child[m]grand_parent[n]的所有组合。

4.3 表头输出和静态变量 time 的坑

代码里用静态变量time控制只在第一次 Reduce 调用时输出表头,static int time = 0在 Reducer 类里定义。在单 Reducer 场景下没问题,但如果有多个 Reducer,每个 Reducer 实例都有可能触发time == 0的条件,导致输出多个表头。静态变量在 MapReduce 里是不可靠的全局状态,不同 Task 跑在不同 JVM 里,根本不共享。这个坑在这份代码里没暴露,是因为默认只有一个 Reducer,但要意识到这个写法有边界。更稳妥的做法是把表头输出放在 main 函数里,或者用MultipleOutputs单独写表头文件。

5. 避坑实录:输出目录、LICENSE.txt 和空字符串引发的翻车现场

这一章把原报告里的三个真实报错,加上我复现时踩过的两个补充坑,按"现象 → 原因 → 解决"拆开写。这些问题都很典型,看完能省下不少排查时间。

5.1 翻车一:输出文件变成一长串 Apache License 2.0

现象:跑合并去重作业,任务显示成功,但打开 part-r-00000,里面不是预期数据,而是一大段英文文本,开头是 Apache License 2.0。

原因:input 目录里残留了一个 Hadoop 自带的 LICENSE.txt 文件。Hadoop 的 FileInputFormat 会读取输入目录下的所有文件,这个 LICENSE.txt 被当成普通数据交给 Mapper 处理了,里面的每一行都成了 key,输出结果自然混进来一堆无关英文。

解决:先hdfs dfs -ls input看看目录里到底有什么,把 LICENSE.txt 删掉或者移到 input 目录外,再重新提交作业。从那以后我每次传数据都要看一眼输入目录,养成随手清理的习惯。

5.2 翻车二:输出目录已存在导致作业直接失败

现象:第一次作业正常跑完,第二次再提交同样命令,报错信息里出现Output directory output already existsFileAlreadyExistsException

原因:FileOutputFormat 硬性要求输出目录在作业启动时不存在,目的是防止覆盖上一次的结果。这个设计避免了很多误操作,但也意味着每次重新跑都要手动清理。

解决:手动执行一次hdfs dfs -rm -r output,或者把自动删除写进 main 方法。第二种方式更省事,具体代码在下一章给出,这里先记住原理。

5.3 翻车三:For input string = "" 的 NumberFormatException

现象:排序实验提交后,作业在 Map 阶段频繁失败,日志里看到java.lang.NumberFormatException: For input string: ""

原因:某个输入文件末尾多打了一个换行,Hadoop 的 TextInputFormat 按行切分时,末尾那行空内容也被当成一条记录传给 Mapper。Integer.parseInt("")直接抛异常。

解决:原报告的做法是删掉多余换行。但更健壮的办法是在 map 方法里加一层防御,判断value.toString()去除空格后是否为空字符串,为空直接 return,不输出任何键值对。这样无论输入文件末尾有没有空行,作业都能正常跑。

5.4 翻车四:改了代码,跑出来结果还是旧的

现象:修改了逻辑,重新导出 jar 提交作业,输出结果和上一版一模一样,怀疑自己的修改没生效。

原因:两个可能性。一是 HDFS 上有旧 jar 没覆盖,提交时用的还是旧包;二是修改的代码本身没进编译产物,Eclipse 导出 jar 时勾选了旧 class 文件。

解决:重新导出 jar 后先看本地文件大小和修改时间,确认 jar 有变化;提交前强制走一遍删除输出目录,保留旧输出带来的干扰也一并排除。这个检查清单我在多个项目里吃过亏才固定下来。

5.5 翻车五:单表关联输出里表头重复出现

现象:跑单表关联时,输出文件里出现多行grand_child grand_parent表头,以为是 Shuffle 把表头当数据重新分发了。

原因:Reducer 类里的static int time在多个 Reducer 实例下不共享。即使设置了多个 Reducer,每个实例都是独立的time = 0,所以每个 Reducer 的输出里都带了一个表头。单 Reducer 默认配置下问题不出现,属于隐藏的定时炸弹。

解决:保持默认一个 Reducer 跑通实验;如果确需多个 Reducer,把表头输出移到 main 函数里,作业提交前用hdfs dfs -mkdir单独建一个表头文件,或者在 Reducer 代码里用作业级计数器判断是否第一个输出。最省事的还是别用静态变量控制表头输出。

6. 进阶:把删输出目录写进 main 方法,再验证一遍分区边界

重复跑实验最烦的就是每次都要手动删输出目录,手动删一两次还能忍,调试参数时跑几十次,纯属浪费生命。解决方式是把这个操作直接写进 main 方法:作业启动前,检查输出路径是否存在,存在就递归删除。这样每次提交作业前不用再手工清目录。

Path in = new Path(args[0]); Path out = new Path(args[1]); FileSystem fileSystem = FileSystem.get(new URI(in.toString()), new Configuration()); if (fileSystem.exists(out)) { fileSystem.delete(out, true); // true 表示递归删除,目录里有文件也能删干净 }

这段代码写在Job.getInstance之前即可。注意FileSystem.delete(path, true)第二个参数表示是否递归删除,伪分布式下目录里只有一个 part-r-00000 和 _SUCCESS,不递归也能删掉,但写成 true 更保险,避免目录层级变化时踩坑。我自己的习惯是逻辑代码写成一个小工具方法,三个实验的主类共用一个清理逻辑,避免每份代码里重复粘贴。

输出目录自动清理解决的是"重复跑"的麻烦,但 MapReduce 作业的验证不能只看是否跑通。自定义 Partitioner 的排序作业,建议把 Reducer 数量显式设成 3,再验证一遍分区边界是否符合预期。代码里设 3 个 Reducer 的方式是job.setNumReduceTasks(3),显式设置后每个分区对应一个输出文件。检查逻辑按表格来:

输入数值范围分区编号对应输出文件排序位置
0 到 bound-10part-r-00000最前面
bound 到 2*bound-11part-r-00001中间
2*bound 到 Maxnumber2part-r-00002最后

验证完成后,记得把job.setNumReduceTasks注释掉或改回 1,因为实验要求的输出文件格式是单文件,三个 part 文件合并后结果虽然全局有序,但和样例输出格式不完全一致。

Combiner 的边界验证也值得顺手做一次。把合并去重作业里job.setCombinerClass(Reduce.class)这一行注释掉再跑一遍,对比两次 Shuffle 传输的数据量和作业总耗时。数据量小看不出差别,但能直观理解 Combiner 的作用。如果是求和类作业,千万别复用 Reducer 当 Combiner,会得到错误的平均值,这是面试高频坑。

从那以后,我每次提交 MapReduce 作业前都强制走一遍固定流程:检查 input 目录只有预期文件,确认输出目录自动清理代码已写入,重新导出 jar 并核对大小,最后先用小数据跑通再看日志。这套流程帮我避开了大部分重复性的翻车,希望帮到你。

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

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

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

立即咨询