ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

MapReduce初级实践:WordCount、Hadoop环境搭建与避坑指南

MapReduce初级实践:WordCount、Hadoop环境搭建与避坑指南 简介基于林子雨《大数据原理与技术》第三版实验5的MapReduce初级编程实践实验报告主要面向学习Hadoop的大数据专业学生和自学者。报告基于Linux建议Ubuntu 16.04环境使用Hadoop 3.2.2围绕文件合并与去重这一典型场景展开。实验中给出了完整的Java MapReduce代码Map阶段将每行文本作为键输出Reduce阶段对相同键只输出一次从而实现去除重复行代码均带有注释便于理解。报告包含输入文件A、B的样例合并去重后的输出文件C以及作业提交参数和运行说明同时展示了MapReduce作业从配置到提交执行的主要步骤能帮助读者理解分布式计算基本流程和去重原理。资源为1个docx文档压缩包大小1.28MB文档结构包括实验环境、实验内容、代码与结果分析便于直接参考或修改。已有14439人学习下载适合课程实验、期末复习或面试准备时查阅。1. 从 WordCount 开始的 MapReduce 初级编程实践到底在练什么如果你打开这份实验指导看到的第一个任务多半是统计文件中单词出现次数。这正是 MapReduce 初级编程实践的核心不是让你一上来就搞倒排索引或矩阵乘法而是通过最经典的 WordCount 把 Map 阶段、Reduce 阶段、输入输出格式、作业提交这一整套流程走通。大数据实验5实验报告要检验的就是你有没有亲手把一份数据切成小块、并行处理、再汇总——而不是只在 IDE 里运行通过就完事。这份实践适合两类人一是刚学完 Hadoop 组件原理、还没写过一行 MR 代码的学生二是想在生产环境摸清作业调度和容错机制的开发新手。初级不代表简单它反而是后面所有复杂实践的骨架你在 WordCount 里理解的 partition、combiner、序列化机制到排序、Join、多输出场景下全是同一套东西。做完这份实验你应该能回答三个问题Mapper 和 Reducer 的输入输出到底是什么驱动类里每个参数对应集群哪一层配置以及数据倾斜或任务卡死时该看哪个日志、哪个计数器。下面从环境搭建开始逐步讲到跑通一个作业的完整链路。2. 跑通实验前的 Hadoop 环境伪分布式比你想的更值得先配好2.1 为什么初级实践不用真实集群很多人在实验前纠结是不是一定要三台服务器起步我每次带学生做 MapReduce 入门都建议先用伪分布式模式。伪分布式Pseudo-Distributed指所有 Hadoop 守护进程跑在同一台机器上DataNode、NameNode、NodeManager 共享一套硬件资源。它和真实集群唯一的区别是节点数而节点数恰恰不是初级实践关注的点。初级实践真正要练的是三件事HDFS 上的文件读写、MapReduce 作业的生命周期、日志和错误排查。这些在伪分布式上全部能复现。数据量小、进程都在本机调试时打一条ssh localhost就能看到全部日志比在远程集群上反复跳板机高效得多。反过来说一上来就买三台云主机却连core-site.xml都还没改过只会让排障半径扩大三倍。2.2 从零安装 Hadoop 最小步骤与核心配置项以 Hadoop 3.3.x 为例安装伪分布式。注意版本选 3.x 是为了用上较新的 YARN 和纠删码但实验配套的旧教程如果写的是 2.x 配置项部分路径名需按 3.x 调整。下载二进制包后解压到/opt/hadoop关键步骤是改五个文件。# 1. 配置 JAVA_HOME 和 Hadoop 环境变量写入 /etc/profile.d/hadoop.sh export JAVA_HOME/usr/lib/jvm/java-8-openjdk-amd64 export HADOOP_HOME/opt/hadoop export PATH$PATH:$HADOOP_HOME/bin:$HADOOP_HOME/sbin # 2. 编辑 /opt/hadoop/etc/hadoop/hadoop-env.sh显式指定 JDK 路径 echo export JAVA_HOME/usr/lib/jvm/java-8-openjdk-amd64 /opt/hadoop/etc/hadoop/hadoop-env.sh # 3. 配置 core-site.xml指定 NameNode 地址 cat /opt/hadoop/etc/hadoop/core-site.xml EOF configuration property namefs.defaultFS/name valuehdfs://localhost:9000/value /property property namehadoop.tmp.dir/name value/opt/hadoop/tmp/value /property /configuration EOF这段配置里fs.defaultFS是 HDFS 的统一入口地址作业里所有相对路径最终都解析到这个地址上。hadoop.tmp.dir存放元数据和本地计算中间结果如果不显式指定默认会落到/tmp目录系统重启后数据被清空NameNode 会因为找不到 fsimage 而报错——这是新手最容易踩的第一个坑。下一步配置 HDFS 副本数和 YARN 资源。伪分布式只有一台机器如果dfs.replication保持默认的 3写入文件时 DataNode 发现副本数不足会一直等待。必须改成 1否则实验里传数据文件进 HDFS 这一步就卡住。# 4. 配置 hdfs-site.xml副本数必须改为 1 cat /opt/hadoop/etc/hadoop/hdfs-site.xml EOF configuration property namedfs.replication/name value1/value /property property namedfs.namenode.name.dir/name value/opt/hadoop/tmp/dfs/name/value /property /configuration EOF # 5. 配置 yarn-site.xml用非 root 用户跑作业时避免本地目录权限问题 cat /opt/hadoop/etc/hadoop/yarn-site.xml EOF configuration property nameyarn.nodemanager.aux-services/name valuemapreduce_shuffle/value /property property nameyarn.nodemanager.local-dirs/name value/opt/hadoop/tmp/yarn-local/value /property property nameyarn.nodemanager.log-dirs/name value/opt/hadoop/tmp/yarn-log/value /property /configuration EOFyarn.nodemanager.aux-services必须配成mapreduce_shuffle少了它 Reduce 任务拉取 Map 输出时会直接失败报错信息是Error: Could not find or load main class org.apache.hadoop.mapreduce.v2.app.MRAppMaster——注意这个报错有迷惑性看起来像类路径缺失实际是 Shuffle 服务没启动。全部配置完成后先格式化 NameNode再启动守护进程。格式化命令只允许执行一次重复执行会清掉原有元数据。启动后立刻用jps检查五个进程是否都在。# 6. 首次启动前必须先格式化之后不能再执行 /opt/hadoop/bin/hdfs namenode -format # 7. 一次性启动所有守护进程 /opt/hadoop/sbin/start-dfs.sh /opt/hadoop/sbin/start-yarn.sh # 8. 验证进程应看到 NameNode、DataNode、SecondaryNameNode、ResourceManager、NodeManager jps我一般会提醒学生把启动日志存一份start-dfs.sh /tmp/hadoop-start.log 21。伪分布式里 NameNode 起不来80% 的原因是格式化时dfs.namenode.name.dir对应的父目录不存在或者 JAVA_HOME 没有正确生效。用jps看到进程后还要访问一次 Web 界面确认活着的进程不是僵尸HDFS 默认端口 9870YARN 默认 8088。这一步虽然不产生代码但决定实验里 HDFS 操作能否执行。2.3 验证环境的最小动作传文件、跑自带示例环境启动后跑 Hadoop 自带的示例程序是成本最低的验证方式。不需要写任何 Java 代码先确认整个数据链路通了再投入写自己的 Mapper 和 Reducer。# 把本地文件放入 HDFS路径以 /input 开头便于识别 /opt/hadoop/bin/hdfs dfs -mkdir -p /input /opt/hadoop/bin/hdfs dfs -put /etc/hadoop/core-site.xml /input/ # 运行官方 WordCount 示例输出目录必须不存在 /opt/hadoop/bin/hadoop jar /opt/hadoop/share/hadoop/mapreduce/hadoop-mapreduce-examples-3.3.6.jar \ wordcount /input/core-site.xml /output/wordcount-test # 查看结果 /opt/hadoop/bin/hdfs dfs -cat /output/wordcount-test/part-r-00000 | head -20如果这个示例能跑完说明 HDFS 读写、YARN 资源分配、MapReduce 框架本身都正常。之后你自己写的作业出错就不需要怀疑环境问题了。注意输出目录output/wordcount-test在运行前必须不存在否则 Job 会抛FileAlreadyExistsException。这个设计是防止上一次运行的结果污染下一次统计我见过有人反复删目录才跑通也有嫌麻烦直接改输出路径的——这是 Hadoop 故意设计的防呆机制不是 bug。跑完示例后还有一个关键习惯去 YARN 的 Web 界面看一次 Application 的日志。ResourceManager 的8088端口点进应用详情能看到 Map 和 Reduce 的进度条、已处理记录数、Shuffle 字节数。这一步很多人忽略但实验报告里的运行结果分析部分这些数字才是最有说服力的素材而不是贴一张终端截图了事。3. 手写 MapReduce 作业WordCount 的分步实现与参数拆解3.1 Mapper、Reducer、Driver 三件套谁负责什么MapReduce 编程模型把一次数据处理拆成三个 Java 类Mapper 负责把输入键值对转换成中间键值对Reducer 负责把相同中间键的值合并Driver 负责组装前两者并提交作业。理解这三者的边界比背代码重要。Mapper 的输入是LongWritable, Text键值对键是行偏移量值是整行文本。它输出的中间键值对经过分组后Reducer 拿到的输入格式是Text, IterableIntWritable同一个单词的所有计数被聚成一个迭代器。这里有个初级实践最常见的概念误区Reducer 的输入值不是一串散装数字而是同一个键对应的值的集合框架保证键相同的记录一定被分发到同一个 Reducer——这就是 Shuffle 阶段做的分区、排序、分组三件事。以 WordCount 为例三个类各司其职Mapper 把一行 hello world hello 切分成(hello,1) (world,1) (hello,1)框架把相同键聚合后交给 Reducer得到(hello, [1,1])Reducer 遍历迭代器求和输出(hello, 2)。Driver 类不处理业务数据它负责设置输入路径、输出路径、键值类型然后调用waitForCompletion(true)提交。3.2 从 IDE 到集群三个类的完整代码与类型匹配逻辑下面这份代码是在 Hadoop 3.x 下可以编译运行的完整 WordCount我把三个类写在一个文件里便于实验提交。注意看泛型类型Mapper 的中间键值类型必须和 Reducer 的输入类型严格对应否则作业提交时直接报类型不匹配。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; import java.io.IOException; import java.util.StringTokenizer; public class WordCount { /** * Mapper 阶段把行文本拆成单词输出 (单词, 1) * LongWritable 是输入键类型行偏移量Text 是输入值类型整行 * 中间键用 Text中间值用 IntWritable必须与 Reducer 输入一致 */ public static class TokenizerMapper extends MapperLongWritable, Text, Text, IntWritable { // 复用对象避免每处理一行就 new 一个减少 GC 压力 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 按空格和制表符切分按行读取由框架保证 StringTokenizer itr new StringTokenizer(value.toString()); while (itr.hasMoreTokens()) { word.set(itr.nextToken()); context.write(word, one); } } } /** * Reducer 阶段对同一个单词的所有 1 求和 * 输入键是 Text 单词输入值是 IterableIntWritable 计数集合 * 输出键是 Text 单词输出值是 IntWritable 总次数 */ public static class IntSumReducer extends ReducerText, IntWritable, Text, IntWritable { private IntWritable result new IntWritable(); Override protected void reduce(Text key, IterableIntWritable 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 { 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); // 以下三行决定 jar 包里的主类是谁实验提交时最容易漏 job.setJarByClass(WordCount.class); job.setMapperClass(TokenizerMapper.class); job.setReducerClass(IntSumReducer.class); // 设置 Mapper 输出类型注意 reducer 输出类型可能不同 job.setMapOutputKeyClass(Text.class); job.setMapOutputValueClass(IntWritable.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); } }这段代码有四个地方需要单独说明。第一setJarByClass(WordCount.class)是告诉 Hadoop 去哪个 jar 里找类。如果通过 IDE 直接运行需要手工指定setJar路径否则会报Class NotFoundException。第二setMapOutputKeyClass与setOutputKeyClass在 WordCount 里虽然一样但设计上允许不同例如 Mapper 输出文本、Reducer 输出自定义序列化对象时这两行必须分别设置。第三Context.write是异步写入的在 map 循环里反复创建新 Text 对象会导致大量小对象我习惯把成员变量定义在类字段上复用这是一个在小数据量下看不出、但大数据量下会影响 GC 的性能细节。第四主函数里通过System.exit把作业成功与否转成进程退出码shell 脚本里可以用$?判断这是方便持续集成用的实验报告里不需要但需要知道。3.3 怎么把代码变成可提交的 jar 包初学者最纠结的是代码写好了怎么跑到集群上。有两条路本地 IDE 运行和打包提交到 Hadoop。实验环境推荐走打包因为你迟早要面对资源管理器上跑作业这个真实场景。打包用 Maven 的标准流程即可。!-- pom.xml 关键依赖版本跟 Hadoop 保持一致 -- dependencies dependency groupIdorg.apache.hadoop/groupId artifactIdhadoop-client/artifactId version3.3.6/version scopeprovided/scope /dependency /dependencies !-- 打包插件生成带 Main-Class 的可执行 jar -- build plugins plugin groupIdorg.apache.maven.plugins/groupId artifactIdmaven-jar-plugin/artifactId version3.3.0/version configuration archive manifest mainClassWordCount/mainClass /manifest /archive /configuration /plugin /plugins /buildscope设为provided是因为 Hadoop 集群自带这些依赖打包时不需要打进 jar否则打出的包会多出十几 MB 的无效体积。maven-jar-plugin的mainClass让hadoop jar命令可以直接找到入口类。有些教程用maven-shade-plugin打成 fat jar对于初级实践完全没有必要反而可能因为 Hadoop 自带类和你的依赖版本冲突出现NoSuchMethodError这种玄学问题。打包命令和提交命令分开写便于实验报告里分别截图说明。# 编译并打包在项目根目录执行 mvn clean package # 查看生成的 jar ls target/*.jar # 提交到 Hadoop 集群运行 WordCount 主类 /opt/hadoop/bin/hadoop jar target/wordcount-1.0.jar /input/core-site.xml /output/wordcount-java提交后终端会持续输出进度信息每个 Map 或 Reduce 完成 100% 时你会看到一行map 100% reduce 100%。如果卡在某一个百分比不动多半不是集群忙而是某个任务在反复重试或等待资源这时去 YARN 页面看日志是最快的排障手段。4. 把 WordCount 改造成生产可用的版本Combiner、自定义序列化与日志验证4.1 Combiner本地预聚合减少集群间数据传输WordCount 虽然能跑但一次完整的作业执行会让每个单词的(hello,1)全部从 Map 节点通过网络传输到 Reduce 节点。假设一个文本里 hello 出现一百万次Reduce 端就要拉取一百万个键值对做累加。Combiner 的作用是在 Map 端先做一次本地累加再把(hello, 1000000)这种合并后的结果传出去。在 WordCount 场景下Combiner 可以直接复用 Reducer 类因为求和操作是幂等可交换的。// Driver 里加一行即可Reducer 类自己可以充当 Combiner job.setCombinerClass(IntSumReducer.class);不要小看这一行。我在生产环境遇到过 Raw 数据量在 10 到 100 GB 的作业加了 Combiner 后 Shuffle 传输量减少 60% 以上作业从 40 分钟压到 15 分钟。但注意 Combiner 的使用条件是多次应用不改变语义WordCount 的求和满足求平均数就不满足——平均数不能本地合并除非你改用(sum, count)的复合值结构。4.2 写出自定义序列化类处理实验里的大概率场景MapReduce 的键值类型必须可序列化内置的Text、IntWritable、LongWritable理解起来简单但实验里常常出现要按某几个字段聚合、同时计算另一个字段的最大最小值这种复合结果需求。这时候需要自定义一个实现Writable接口的类。以下代码是统计每个 IP 的访问次数和平均响应时间的复合值类型这也是实验里会出现的需求变形。import org.apache.hadoop.io.Writable; import java.io.DataInput; import java.io.DataOutput; import java.io.IOException; public class AccessStatWritable implements Writable { private int count; // 访问次数 private long totalTime; // 响应时间总和 // 反射构造时需要无参构造器反序列化才会调用 public AccessStatWritable() {} public AccessStatWritable(int count, long totalTime) { this.count count; this.totalTime totalTime; } Override public void write(DataOutput out) throws IOException { out.writeInt(count); out.writeLong(totalTime); } Override public void readFields(DataInput in) throws IOException { // readFields 的读取顺序必须与 write 的写入顺序完全一致 this.count in.readInt(); this.totalTime in.readLong(); } public int getCount() { return count; } public long getTotalTime() { return totalTime; } }这个类有两个坑需要点出来。第一write和readFields的顺序必须严格配对先写 int 就先读 int。如果你在write里先写count后写totalTime在readFields里顺序换了你得到的就是串位的脏数据。第二必须保留无参构造器。Hadoop 反序列化时通过反射调用无参构造器创建对象如果只定义了带参构造器会出现RuntimeException: No such constructor这个错在代码里看不出来运行到 Shuffle 阶段才爆。调试时我一般会在readFields里打一行临时日志对比读到的值和预期值确认序列化没有错位。4.3 实验报告最需要的运行数据怎么取很多人的实验报告里只有一张终端截图但验证作业正确性需要一些硬指标。在 Driver 里可以主动打印计数器的值这是初级实践中最容易忽略、但最能证明你理解框架运行机制的做法。// 在 main 函数 waitForCompletion 之后增加片段 boolean success job.waitForCompletion(true); if (success) { // 按计数器组名和计数器名取值综合信息在 Map 组 long mapRecords job.getCounters() .findCounter(org.apache.hadoop.mapreduce.TaskCounter, MAP_INPUT_RECORDS) .getValue(); long reduceGroups job.getCounters() .findCounter(org.apache.hadoop.mapreduce.TaskCounter, REDUCE_INPUT_GROUPS) .getValue(); System.out.println(Map 输入记录数: mapRecords); System.out.println(Reduce 输入分组数: reduceGroups); }MAP_INPUT_RECORDS是 Map 阶段实际读入的行数REDUCE_INPUT_GROUPS是经过 Shuffle 分组后的唯一键数量。拿 WordCount 来说如果输入文件有 100 行、5000 个不同单词这两个数字就是 100 和 5000。把它们写进实验报告的分析区再对比输出文件的行数就能证明我的作业没有丢数据也没有重复计算。这是符合工程验证思路的做法而不是只描述程序运行成功。这些计数器信息其实在 YARN 页面里也能看到但通过代码主动打印能让报告更清晰——尤其是初级实践往往要求你在报告中体现对框架的理解这组数字就是最有力的证据。5. 初级实践避坑指南现象、原因与解决的完整排查对照伪分布式环境加基础 WordCount看起来简单实际带过的学生十个里有八个卡在同一批问题上。下面按踩坑频率排序每条都写了现象、原因和解决方式。坑一localhost:9000: Connection refused现象执行 HDFS 命令时直接报连接拒绝随后一条Call From localhost/127.0.0.1 to localhost:9000 failed on connection exception。原因NameNode 进程没有启动。伪分布式里最常见是格式化后没有start-dfs.sh或者启动脚本报错但被忽略次常见是hadoop.tmp.dir目录权限不对NameNode 启动时写不了元数据。解决先jps确认进程是否存在不存在则手动启动并查看日志。tail -f /opt/hadoop/logs/hadoop-xxx-namenode.log看到报错信息再对症处理。如果日志提示NameNode is not formatted重新执行一次hdfs namenode -format再启动。注意格式化命令会清空所有 HDFS 数据只在你确定没有需要保留的文件时执行。坑二输出目录已存在导致作业秒失败现象hadoop jar提交作业后几秒钟就退出日志里有FileAlreadyExistsException: Output directory hdfs://localhost:9000/output/xxx already exists。原因Hadoop 故意设计成不允许覆盖输出目录防止误操作冲掉历史结果。很多同学以为和 Linux 的 file一样可以重复写其实是理解偏差。解决作业运行前先确认输出路径不存在。在 Driver 代码里可以加一段逻辑提交前检查 HDFS 上是否存在输出路径存在就删除。但工程上更推荐保留这个机制因为误删场景比重复跑更危险。实验场景下手动删即可hdfs dfs -rm -r /output/xxx。坑三Map 跑完 Reduce 卡在 33% 不动现象进度条显示map 100% reduce 33%然后长时间无变化最后作业超时失败。原因Reduce 的 33% 代表正在执行 Shuffle 阶段此时会从所有 Map 任务所在节点拉取中间数据。伪分布式里卡住的原因通常是 YARN 日志目录磁盘满了或者yarn.nodemanager.aux-services配置缺失导致 Shuffle 服务未启动。解决检查磁盘df -h清理旧作业的中间文件和日志。确认yarn-site.xml里mapreduce_shuffle已配置且重启过 YARN。如果崩在 33% 之前检查网络和防火墙得更细致一些——本地回环一般不会是这个原因。坑四自定义类实现 Writable 后作业运行报NoSuchMethodException现象作业在 Map 阶段或 Reduce 阶段突然失败日志里有java.lang.RuntimeException: java.lang.NoSuchMethodException: xxx.init()。原因自定义 Writable 类没有写无参构造器。反序列化需要反射创建一个空对象再调用readFields填充字段。Java 里如果你定义了带参构造器无参构造器不会自动生成于是反射创建失败。解决补一个无参构造器。这个坑不仅限于 WritableReducer 的setup方法里实例化自定义对象时也可能碰到类似反射问题。写完自定义类建议习惯性补一个空的构造函数。坑五在 IDE 里能跑通提交到集群就报ClassNotFoundException现象本地运行正常用hadoop jar提交到集群后作业失败错误指向你的类。原因jar 里没有包含主类所在的所有类文件或者setJarByClass指定的类和实际 jar 包路径不匹配。另一个常见因素是 Maven 打包时输出的是不带依赖的瘦包而某些类实际上依赖了第三方库。解决检查target/下的 jar 是否存在用jar tf target/xxx.jar | grep WordCount查看是否有这个类。如果是依赖问题把依赖打进 jar 用maven-shade-plugin但要注意排除org.apache.hadoop相关包优先级比hadoop-client的providedscope 低。坑六运行start-dfs.sh时提示Did not find an YARN configuration现象脚本执行后只启动了 NameNode 和 DataNodeYARN 相关进程没起来日志说没找到配置。原因start-dfs.sh只管 HDFSYARN 需要单独执行start-yarn.sh。两者是独立的启动入口不动脑子的做法是把 yarN-site 配好却不执行第二个脚本。解决补执行start-yarn.sh后jps确认 ResourceManager 和 NodeManager 出现。把两个启动命令写进一行脚本里减少漏启动。这是新手期最容易忽略的分离操作。排查以上问题时我强烈建议保持一个习惯不要只看终端最后的报错把tail -f /opt/hadoop/logs/userlogs/*/syslog打开的日志窗口作为调试桌面的一部分。初级实践的大部分报错在 syslog 里看前五十行就能定位比反复重新提交作业高效得多。6. 从 WordCount 跨出第一步自定义排序与分组把实验报告写成能打的简历素材初级实践完成 WordCount 只是及格线如果想让这份实验在后续找实习、毕业答辩时有说服力我建议在 WordCount 基础上加一个自定义排序需求按单词出现次数降序排列次数相同的按单词字典序排列。这个需求一句话就能说清但完全覆盖了 MapReduce 排序机制、自定义 Comparator 和输出格式三个考点。实现思路是这样的MapReduce 的默认排序只发生在键上而且只按键的自然序。你要求按值排序就得把值和键的位置对调——用IntWritable作为键、Text作为值输出到 Reduce。但框架默认对IntWritable升序排列你又要求降序所以需要自定义一个 Comparator 重写比较逻辑。import org.apache.hadoop.io.IntWritable; // 自定义 IntWritable 降序比较器注意泛型和比较方向 public class IntWritableDecreasingComparator extends IntWritable.Comparator { Override public int compare(byte[] b1, int s1, int l1, byte[] b2, int s2, int l2) { // 调用父类比较后取负实现降序 return -super.compare(b1, s1, l1, b2, s2, l2); } }然后在 Driver 类的main方法里指定排序比较器。setSortComparatorClass作用在 Reduce 阶段的输入键排序上注意它影响的不只是输出顺序而是 Shuffle 分组前的排序规则——如果相同的键因为比较器判定为不同而分不到同一个 Reducer就会出现相同键被分散处理的隐蔽问题。// 在 main 中配置作业 job.setSortComparatorClass(IntWritableDecreasingComparator.class);这个自定义排序做完你就能回答实验必问的一个问题了setSortComparatorClass、setGroupingComparatorClass、setPartitionerClass三者各自管哪个阶段。排序管的是键怎么比大小分组管的是哪些键算一组进同一个 reduce 方法分区管的是哪些键进同一个 Reducer 节点。三个概念各自独立WordCount 里体会不到排序需求一加就清楚了。最后验证结果时对比一下排序前后输出文件和自定义 Comparator 对应的键顺序。在实验报告的收尾部分我会建议写一段待改进项例如你怎么知道这次作业有多少个 Reduce 任务可以在mapred-site.xml里搜mapreduce.job.reduces默认值是 1意味着所有单词都在同一个 Reducer 里处理。输出只有一个文件part-r-00000这个现象可以作为理解分区与并行度的引子。接下来要想让不同字母开头的单词分别进入不同文件你有Partitioner可以玩但那已经是下一个实验的内容了。我在做带教和自学的这些年里反复确认过一件事MapReduce 初级编程实践最难的地方从来不是 API 记不熟而是学会把一份数据处理任务拆成 Map、Shuffle、Reduce 三个环节去思考。能把 WordCount 逐渐改造成排序作业、能看懂 YARN 页面上的 Shuffle 统计、能解释默认 1 个 Reducer 带来的文件数变化你就真正从会写代码跨到了会做工程。希望这份实践梳理能帮你在实验路上少踩几个无谓的坑把时间花在理解机制而不是排错上。本文还有配套的精品资源点击获取
返回列表