尧图精选

Hadoop MapReduce三个经典编程实例:去重、排序、单表关联

🕒 发布时间:2026/10/2 13:12:02 📁 来源:尧图网络
简介大数据实验五实验报告围绕MapReduce初级编程实践主题面向正在学习大数据处理技术的高校学生与自学者可作为《大数据原理与技术》第三版配套实验的参考材料。报告基于Linux建议Ubuntu 16.04与Hadoop 3.2.2环境重点讲解如何编写MapReduce程序实现两个文件的合并与去重从Mapper读取输入、复制key值到Reducer输出唯一结果并附有A/B输入文件与C输出文件的完整样例便于对照验证。除核心代码外报告还包含了实验环境配置说明、Java代码注解和运行步骤能够帮助初学者理解map与reduce阶段的工作机制并迁移到其他数据处理场景中。该资源包内共1个docx格式文档压缩包大小约1.28MB文件结构紧凑下载后可直接阅读或按需修改使用。目前已有14439人学习使用是一份兼具教学参考与作业参考价值的MapReduce入门实践资料。1. 从一份实验报告到可复用的 MapReduce 代码库三个作业一次讲透手上这份林子雨《大数据原理与应用》第三版的实验5报告表面上是三个入门作业——文件合并去重、整数排序、单表关联找祖孙关系实际上把 MapReduce 最核心的 shuffle、Partitioner、Join 三种机制都过了一遍。我拆完这套代码的感受是它比市面上很多教程都实在代码能直接编译提交坑也都记录在案非常适合刚搭好 Hadoop 3.2.2 集群、正在被作业和面试题两头夹击的人。你不需要再去找零散的 mapreduce 编程实例把这三个作业跑通、跑懂MapReduce 基础编程的地基基本就稳了。这篇笔记会把三个作业的原理、代码、参数和踩坑点全部摊开来讲。2. 文件合并与去重把「按 key 分组」变成一种去重武器2.1 原理为什么去重不需要在 reduce 里写判断第一个作业的目标很简单两个输入文件 A 和 B每行是「日期 空格 字符」合并后剔除完全相同的行。很多第一次写 MapReduce 的人会条件反射地在 reduce 里维护一个 HashSet 来判重这其实是把单机思维带进了分布式模型。MapReduce 框架本身在做的事情就是「按 key 分组」map 输出的所有键值对经过 shuffle 阶段相同的 key 一定会被送到同一个 reduce 调用里。这意味着去重根本不需要额外逻辑——你只要把每一行内容当作 key 输出reduce 阶段对每个 key 只输出一次天然就是去重结果。这个思路第一次想通之后后面很多作业都会用到。这份报告里 Map 类的实现也印证了这一点它没有做任何拆分和过滤直接把整行 value 当成 key 输出value 写一个空字符串占位。Reduce 端更简单不加任何判断直接context.write(key, new Text())。整个作业的逻辑就两句话真正的重活全交给框架的 shuffle 机制了。2.2 代码走读Map、Reduce 与 job 配置的关键点完整代码在报告里我这里只挑核心段落拆解。Map 类的写法如下public static class Map extends MapperObject, Text, Text, Text { private static Text text new Text(); public void map(Object key, Text value, Context content) throws IOException, InterruptedException { text value; content.write(text, new Text()); } }这里有个小细节text被声明为 static 成员变量每次 map 调用直接复用对象避免频繁 new 对象。这是 Hadoop 性能优化里的常见手法因为 map 方法会被调用几万次对象复用能明显减少 GC 压力。在实际生产中如果 value 的内容需要拆分或清洗我一般会在这里用value.toString()生成新字符串处理但纯去重场景直接赋值引用是效率最高的。Reduce 类的实现更短public static class Reduce extends ReducerText, Text, Text, Text { public void reduce(Text key, IterableText values, Context context) throws IOException, InterruptedException { context.write(key, new Text()); } }reduce 方法里根本不需要遍历 values因为每个 key 对应的 value 都是空字符串没有任何聚合需求。这里能这样写的前提是 map 输出阶段已经把 value 统一设成了所以这个 reduce 实际是在「消费」shuffle 的结果而不是在处理数据。主函数里的 Job 配置有几个容易被忽略的参数我重点说一下Configuration conf new Configuration(); conf.set(fs.defaultFS, hdfs://localhost:9000); Job job Job.getInstance(conf, Merge and duplicate removal); job.setJarByClass(Merge.class); job.setMapperClass(Map.class); job.setCombinerClass(Reduce.class); job.setReducerClass(Reduce.class); job.setOutputKeyClass(Text.class); job.setOutputValueClass(Text.class); FileInputFormat.addInputPath(job, new Path(otherArgs[0])); FileOutputFormat.setOutputPath(job, new Path(otherArgs[1])); System.exit(job.waitForCompletion(true) ? 0 : 1);fs.defaultFS指向hdfs://localhost:9000这是伪分布式部署时 NameNode 的默认地址如果你的集群是真正的多节点部署这里的值要改成熟 namenode 的地址。setCombinerClass(Reduce.class)这一行值得注意——combiner 在 map 端先做了一次局部合并对于去重场景能显著减少 map 到 reduce 之间的网络传输量。由于这里的 reduce 函数满足交换律和结合律输出只依赖 key不依赖 value直接复用 Reduce.class 当 combiner 是安全的。setOutputKeyClass和setOutputValueClass都设成 Text 是因为中间结果和最终结果的类型完全一致这是这个作业最简单的地方。真正会翻车的是输入输出路径的配置后面避坑章节会详细展开。2.3 运行验证从 HDFS 到结果核对代码编译通过后常规的提交命令是hadoop jar Merge.jar Merge input output注意这里的 input 和 output 是 HDFS 路径不是本地路径。很多新手在本地跑通 eclipse 后直接照搬本地文件路径去hadoop jar结果报错找不到文件。提交前先确认输入文件已经上传到 HDFShdfs dfs -mkdir -p /user/hadoop/input hdfs dfs -put fileA.txt fileB.txt /user/hadoop/input/ hdfs dfs -ls /user/hadoop/input/跑完后用hdfs dfs -cat output/part-r-00000查看结果对照报告里给的样例20170101 x 20170101 y 20170102 y 20170103 x 20170104 y 20170104 z 20170105 y 20170105 z 20170106 x注意 A 和 B 都有的20170103 x只出现一次这就是去重生效了。检查结果时逐行比对这个文件能快速发现是不是混入了多余数据。3. 全排序作业自定义 Partitioner 与「分区边界」的玄学3.1 为什么全排序不能只靠框架默认行为第二个作业多个输入文件每行一个整数要求把所有整数升序排列输出格式为「位次 原值」。如果数据量小单机排序毫无压力。但 MapReduce 场景下有一个经典问题map 输出的 key 经过 shuffle 后每个分区内部是有序的但分区之间没有全局顺序。框架能保证的是每个 reduce 收到的 key 都是排好序的但如果你有 3 个 reduce第 2 个分区的数据不一定都大于第 1 个分区。要实现全局有序核心思路是让 key 值小的数据永远进入序号小的分区也就是把「分区号」和「key 的取值范围」建立对应关系。这正是自定义 Partitioner 的价值。3.2 代码走读Partitioner 的边界计算与 line_num 计数Map 类的逻辑很简单读入一行解析成整数把整数本身作为 key 输出value 写死为 1。public static class Map extends MapperObject, Text, IntWritable, IntWritable { private static IntWritable data new IntWritable(); public void map(Object key, Text value, Context context) throws IOException, InterruptedException { String text value.toString(); data.set(Integer.parseInt(text)); context.write(data, new IntWritable(1)); } }value 之所以恒为 1是因为 reduce 需要遍历 values 来确定这个 key 出现了几次从而决定输出几行。这里如果改成任何其他常量都行关键在 values 的个数。Reduce 类里有一个全局计数器line_num从 1 开始递增public static class Reduce extends ReducerIntWritable, IntWritable, IntWritable, IntWritable { private static IntWritable line_num new IntWritable(1); public void reduce(IntWritable key, IterableIntWritable values, Context context) throws IOException, InterruptedException { for (IntWritable val : values) { context.write(line_num, key); line_num new IntWritable(line_num.get() 1); } } }这里有个 subtle 的点key 相同的所有 value 在同一个 reduce 调用里被遍历所以重复整数只会在这里输出一次位次递增。注意line_num是 static 变量如果有多个 reduce task 并发跑这个计数器会互相干扰——这也是为什么这个作业必须配合自定义分区保证「全局只有一个 reduce」或者「每个 reduce 处理不相交的 key 区间」否则排序位次会乱。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; } }这个分区的思路是把0 ~ Maxnumber这个值域均匀切成num_Partition份每份的宽度是bound然后判断 key 落在哪个区间里。这里写死Maxnumber 65223是最大安全隐患——如果输入数据里有大于 65223 的整数getPartition会返回 -1作业直接报错。实际生产中正确的做法是从输入数据里算出最大值比如在 map 阶段用计数器统计或者直接设成Integer.MAX_VALUE。另外注意bound Maxnumber/num_Partition1加 1 的目的是保证bound * num_Partition一定大于Maxnumber避免最大值落入「无法匹配任何分区」的情况。这个 1 不写130 这个数在分区数为 2 时就会出问题因为65223/2 32611130 落在分区 0但如果某个数恰好等于32611*2 65222会落不进任何区间。这个细节很多人会忽略属于典型的边界条件翻车点。主函数里比实验一多了这一行job.setPartitionerClass(Partition.class);这行代码告诉框架map 输出的键值对用我们自定义的逻辑决定去哪个分区。如果把每个分区的数据单独拿出来看每个 reduce 内的 key 是升序的又因为分区编号从小到大对应 key 值从低到高最终合并起来就是全局升序。3.3 验证方法查看每个 reduce 的输出文件提交命令和实验一类似hadoop jar MergeSort.jar MergeSort input output如果作业配置了 3 个 reducer输出目录下会出现 3 个 part-r-00000、part-r-00001、part-r-00002。验证全局排序正确的方式是检查第一个文件的最大值小于第二个文件的最小值hdfs dfs -cat output/part-r-00000 | tail -1 hdfs dfs -cat output/part-r-00001 | head -1报告里的样例输出是 11 行排序结果对应 3 个输入文件的 11 个整数。跑通后建议自己构造一批含重复值和大数值的测试数据看看line_num计数和分区边界是否还能正常工作。4. 单表关联作业一张表拆成左右两半做连接4.1 原理把「找祖孙」转换成「两次 join」第三个作业的输入是一张 child-parent 表让你输出 grandchild-grandparent 关系。很多人第一次看到这个题目会懵MapReduce 擅长的是批处理这种需要「跨行推理」的任务怎么下手答案是单表连接(Self-Join)把同一张表当作左右两个关系表来处理。对于每一条记录(child, parent)我们做两次 map 输出第一次以parent作为 key称之为「左表输出」。这样所有以 X 为 parent 的记录会聚到同一个 reducekeyX 时value 里存的是 X 的孩子。第二次以child作为 key称之为「右表输出」。这样所有以 Y 为 child 的记录会聚到同一个 reducekeyY 时value 里存的是 Y 的父母。当某个 key 同时出现在左表集合和右表集合时比如 keyLucy左表里有(Lucy, Steven)和(Lucy, Jone)说明 Lucy 是 Steven 和 Jone 的 parent右表里有(Lucy, Mary)和(Lucy, Frank)说明 Lucy 是 Mary 和 Frank 的 child。于是「Lucy 的孩子」和「Lucy 的父母」两两配对正好得到「(Steven, Mary)、(Steven, Frank)、(Jone, Mary)、(Jone, Frank)」四条祖孙记录。这个思路是所有 MapReduce join 类问题的通用解法理解了它之后做 reduce-side join 就会顺畅很多。4.2 代码走读relation_type 标志位与笛卡尔积Map 类里最核心的就是这个双输出设计String relation_type 1; // 左右表区分标志 context.write(new Text(values[1]), new Text(relation_type child_name parent_name)); // 左表以 parent 为 keyvalue 里带标志 1 relation_type 2; context.write(new Text(values[0]), new Text(relation_type child_name parent_name)); // 右表以 child 为 keyvalue 里带标志 2这里有一个典型的参数坑relation_type用的是 String在 reduce 端用charAt(0)取第一个字符来判断是 1 还是 2这个设计在数据格式稳定时没问题但如果 value 里出现以数字开头的异常行判断就会错位。更稳妥的做法是用Text类型分别加前缀或者直接用两个不同的 key 标记。Reduce 端的处理逻辑是把这个 key 下的所有记录分成两组if (relation_type 1) { grand_child[grand_child_num] child_name; grand_child_num; } else { grand_parent[grand_parent_num] parent_name; grand_parent_num; }然后做笛卡尔积for (int m 0; m grand_child_num; m) { for (int n 0; n grand_parent_num; n) { context.write(new Text(grand_child[m]), new Text(grand_parent[n])); } }这套代码在「一个 key 关联的记录数较少」时没问题但grand_child[]和grand_parent[]数组长度写死为 10一旦某个 key 下的记录超过 10 条就会数组越界。真实生产环境应该用 ArrayList 动态扩容或者直接用字符串拼接输出。这是个典型的「作业能过生产必炸」的写法。4.3 结果验证与数据规模意识的建立跑完后检查输出对照报告里的样例核心验证点是 Steven、Jone、Philip、Mark 四个孙辈是否分别对应到了正确的祖辈。这个作业跑通后建议做一次「数据规模放大测试」把输入表扩到 100 个人随机生成 300 条关系观察 reduce 端内存是否会爆。MapReduce 作业的 reduce 方法是逐 key 调用的如果某个 key 关联的数据特别多也就是所谓的「数据倾斜」这里就是最先卡住的地方。理解了这个问题之后看 hive 或 spark 的 join 优化方案时会更容易。5. 避坑指南三个把作业卡住的典型报错5.1 输出目录已存在FileAlreadyExistsException现象第二次运行同样的作业提交后立刻报org.apache.hadoop.mapreduce.lib.output.FileAlreadyExistsException: Output directory output already exists。原因Hadoop 设计上不允许覆盖已有输出目录这是为了防误删数据。每次跑完作业输出目录都会保留下次运行前必须手动删除。解决报告里给出了自动删除方案但写法有点绕我一般直接在 main 方法里这样加Path in new Path(args[0]); Path out new Path(args[1]); FileSystem fs FileSystem.get(new URI(in.toString()), conf); if (fs.exists(out)) { fs.delete(out, true); }fs.delete(out, true)的第二个参数是递归删除标志必须为 true否则目录非空时删不掉。从那以后我每次提交作业前都强制先跑一遍这段清理逻辑不再手动去 HDFS 删除目录。5.2 输入文件夹里的多余文件LICENSE.txt 混入数据现象第一个实验做文件合并去重输出结果里出现一长串 Apache License 文本和正常的日期字符数据混在一起。原因上传数据时input 目录里遗留了 Hadoop 自带的 LICENSE.txt 文件。MapReduce 的FileInputFormat会读取输入目录下的所有文件不会因为你只放了 A、B 两个文件就只读这两个。多余的 LICENSE.txt 也被当成输入数据一行行进 map输出结果自然脏了。解决运行前用hdfs dfs -ls input检查目录内容确保只有预期的输入文件。这个坑的教训是在 Hadoop 里输入路径是「目录」时目录下的所有文件都会被消费没有「只处理指定文件」的隐含约定。如果想过滤特定文件得用FileInputFormat.setInputPaths指定具体文件路径或者自定义 InputFormat。5.3 空字符串解析报错For input string: 现象第二个实验排序作业删除了输出目录后运行依然报错java.lang.NumberFormatException: For input string: 。原因某个输入文件的末尾多打了一个换行Hadoop 按行读取时把这个空行也读进来了。Integer.parseInt()当然会抛异常。这类问题在手工编辑数据文件时极其常见——文件末尾多个回车、中间多个空行、行尾多个空格都能让解析逻辑崩掉。解决map 阶段加一个空行过滤即可String text value.toString().trim(); if (text.isEmpty()) { return; } data.set(Integer.parseInt(text));trim()顺便把行尾可能残留的\rWindows 换行符也清掉避免parseInt碰见123\r这种脏数据。这是我在处理所有文本类输入时都会写的防御代码。5.4 小集群内存不足Container 被 kill现象第三个实验跑祖孙关系挖掘时数据量稍微放大NodeManager 日志里出现Container killed by the ApplicationMaster。原因reduce 端笛卡尔积在内存里展开某个 key 关联的记录多数组不够或被 OOM。伪分布式环境默认的容器内存很小reduce 端的超大 key 组直接撑爆。解决一方面改造 reduce 逻辑用外部排序或流式输出另一方面调整 yarn 的内存配置yarn-site.xml里的yarn.nodemanager.resource.memory-mb和yarn.scheduler.maximum-allocation-mb适当调大。但根本解法还是控制单 key 的数据倾斜——这也是后来我在真实项目中反复踩的坑单表 join 一旦遇到超级节点内存方案几乎都是死路得换思路。6. 收尾技巧把提交命令固化成一条可重复执行的脚本三个作业都跑通之后我强烈建议做一件事写一个 shell 脚本把编译、清理、提交、查结果四个步骤串起来。这个习惯会让你在后续做更大规模的 mapreduce 编程时省下大量重复劳动。以第二个排序作业为例我的脚本长这样#!/bin/bash # 用法: ./run_sort.sh input_path output_path INPUT$1 OUTPUT$2 # 1. 编译 Java 源码并打 jar 包 rm -rf classes mkdir classes javac -classpath $(hadoop classpath) -d classes MergeSort.java jar cf mergesort.jar -C classes . # 2. 删除可能存在的旧输出目录 hdfs dfs -rm -r -f $OUTPUT # 3. 提交 job hadoop jar mergesort.jar MergeSort $INPUT $OUTPUT # 4. 展示结果 hdfs dfs -cat $OUTPUT/part-r-00000hadoop classpath是动态获取 Hadoop 全部依赖 jar 的快捷方式免去手动拼 classpath 的麻烦。-rm -r -f是幂等的目录清理命令输出目录不存在时也不会报错。这样每次调参只需改输入输出路径整个流程三秒钟跑完。在调分区数的时候我一般会加一个参数-D mapreduce.job.reduces3通过-D动态覆盖mapreduce.job.reduces不用改代码就能测试不同分区数下的排序结果。写完这个脚本后我每次做 MapReduce 作业都强制走一遍「编译 → 清理 → 提交 → 验证」的流程尤其是清理输出目录和检查输入目录这两步已经成了肌肉记忆。这套流程对新手最友好的地方在于它把「作业跑起来」这件事本身变成了确定性行为你只需要专注于调试 map 和 reduce 逻辑本身。最后分享一个从我第一次跑通这三个作业就养成的习惯任何一次 MapReduce 作业运行前我都要花 30 秒确认三件事——输入目录里没有多余的杂文件、输出目录已被清理、输入数据里没有空行。这三件事用一分钟的脚本就能全部自动化但这三件事能帮你省下的排查时间往往是按小时计的。希望这份拆解笔记帮到你至少让你在跑 MapReduce 作业的时候能少掉我当初踩过的那些坑。本文还有配套的精品资源点击获取
上一篇/下一篇内容由系统自动关联 返回资讯列表 →