ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

MapReduce初级编程实践:三个Java程序详解合并去重与全局排序

MapReduce初级编程实践:三个Java程序详解合并去重与全局排序 简介这是一份面向高校大数据课程学习者的实验报告资源基于林子雨《大数据原理与技术》第三版第五章内容完整呈现MapReduce初级编程实践过程。报告以文件合并与去重为实战目标给出Hadoop 3.2.2环境下的可运行Java代码涵盖Map、Reduce阶段的关键实现、作业配置与提交方式并附有输入输出样例供对照验证适合正在完成同类实验或复习MapReduce原理的读者参考。资源包共1个文件为docx格式文档大小约1.28MB排版清晰、内容完整可直接作为实验报告撰写与代码调试的参照。已有14000余人浏览学习是同类资源中较受关注的实验案例。通过阅读该报告读者能够快速理解如何利用MapReduce解决去重与合并问题掌握编写、配置和运行MapReduce作业的基本思路为后续分布式数据处理实践打下扎实基础。1. 大数据实验5 MapReduce 初级编程实践三个跑通的 Java 程序一次讲清大数据实验五的 MapReduce 初级编程实践给的是三个已经跑通的 Hadoop Java 程序文件合并与去重、多文件全局排序、单表关联挖掘祖孙关系配套信息是林子雨《大数据技术原理与应用》实验5Hadoop 版本 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.sh和start-yarn.sh。上传数据后可以用hdfs dfs -ls /user/ubuntu/input确认文件完整也可以在浏览器打开localhost:9870Hadoop 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_xxxhadoop jar和hdfs dfs是两套命令前者提交 MapReduce 作业后者操作 HDFS 文件新手常把这俩搞混。output 目录必须是 HDFS 上不存在的路径这是 FileOutputFormat 的硬性规定防止覆盖旧数据。跑完后目录里会出现两个文件_SUCCESS标记作业成功part-r-00000是真正的输出。如果作业失败优先看yarn logs里的堆栈信息比在控制台瞎猜有效得多。3. 合并去重与全局排序第一个 Job 和自定义 Partitioner 背后的为什么这一章讲前两个实验文件合并去重和全排序。这两个实验放在一起看特别合适因为它们的 Map 阶段几乎一样简单但 Reduce 和分区策略完全不同对比着看能理解 Shuffle 阶段到底替你做了什么。3.1 合并去重Map 阶段把整行当 keyReduce 阶段只透传一次代码的巧妙之处在于它没有做任何去重操作只是把整行文本作为 key 输出value 全部置空。Shuffle 阶段会按 key 分组相同 key 的所有 value 会合并成一个 list 送到同一个 Reducer。Reducer 拿到 key 后不管 values 里有几条只输出一次 key去重就完成了。public static class Map extends MapperObject, 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 ReducerText, Text, Text, Text { // 同一个 key即同一行文本只输出一次重复内容自然消失 public void reduce(Text key, IterableText 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 分到同一个 keyReducer 只输出一次最终得到的就是合并且去重的 C 文件。3.2 全局排序三个 Reduce 各排各的自定义分区才能全局有序第二个实验的输出要求是每行两个数字排序位次和原始整数。如果只靠默认的 HashPartitioner多个 Reducer 各自有序但把多个 part-r-0000x 文件拼起来看整体是乱序的。实验要求输出到一个文件里的数据是升序排列所以必须自定义 Partitioner让数值按大小区间划分到不同分区每个分区只负责一个数值段。默认情况下 MapReduce 会启动一个 Reducer这个数量不足以体现分区的意义。贴出来的代码里没有显式设置 Reducer 数量走的是默认值 1这种情况下自定义 Partitioner 没有实际效果但它展示了自定义分区的完整写法值得逐行拆解。// 自定义分区根据数值大小划分到不同分区保证分区之间数值范围严格分离 public static class Partition extends PartitionerIntWritable, 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分区数是 3bound 是 21742。那么 0 到 21741 落在分区 021742 到 43483 落在分区 143484 到 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 parentMapper 需要输出两条记录。第一条以 parent 为 keyvalue 里包含 child 信息这构成了找孙子的左表第二条以 child 为 keyvalue 里包含 parent 信息这构成了找爷爷的右表。为了区分这两条记录value 前面加了一个标志位1或2。输入行输出 key输出 value语义Steven LucyLucy1StevenLucy左表Steven 是 Lucy 的孩子Steven LucySteven2StevenLucy右表Lucy 是 Steven 的父辈public static class Map extends MapperObject, 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 的格式是标志位childparent加号是自定义的分隔符这里不像表头那样受空格干扰只要保证 Mapper 输出和 Reducer 解析用的是同一个分隔符就行。4.2 Reducer 里做笛卡尔积拆串、分组、全组合Reducer 收到的是相同 key 的 value-list。以Lucy这个 key 为例输入数据里有Steven Lucy和Lucy Mary两行。Mapper 对第一行输出了左表记录1StevenLucy对第二行输出了右表记录2LucyMary。这两个 value 会同时落到 keyLucy 的 reduce 组里Reducer 把左表记录里的 childSteven放进 grand_child 数组把右表记录里的 parentMary放进 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 exists或FileAlreadyExistsException。原因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 并核对大小最后先用小数据跑通再看日志。这套流程帮我避开了大部分重复性的翻车希望帮到你。本文还有配套的精品资源点击获取
返回列表