Java操作MapReduce实战:从环境搭建到集群提交
发布时间:2026/9/17 15:15:54 锦皓数字建站

简介Hadoop大数据处理技术实验报告完整版面向大数据初学者与高校学生围绕Java操作MapReduce的完整实践流程讲解如何用MapReduce编程模型分析2020年中国地区气象监测站数据找出每个月出现最高温度的监测站信息。报告包含实验目的、CentOS 7环境配置、Map与Reduce函数编写、自定义分区器将1-6月与7-12月结果分别输出两个文件、Maven打包jar并在Hadoop集群运行等核心步骤并对温度正负值、9999异常值处理等易错点做了说明最后附有实验总结与心得能帮助读者快速掌握Hadoop MapReduce基础编程与集群部署方法。资源为单个doc文档大小765KB内容结构完整适合当作课程实验参考或复习备考资料。目前已有190人学习下载。1. 把 Java 操作 MapReduce 的实验跑通先跟 HDFS 和 YARN 握手实验报告标题里出现“Java 操作 MapReduce”说明这门课的训练目标不是让你在 IDE 里把 WordCount 打印出来而是让你站在客户端视角用 Java 程序向 Hadoop 集群提交一个分布式计算任务。很多人在这一步卡住并不是 Mapper 和 Reducer 写得不对而是“本地跑通”和“集群跑通”之间的鸿沟没填平JAVA_HOME 没配对、core-site.xml 没有被打进 ClassPath、输出目录已存在导致作业直接拒绝。真正把实验做厚的人会在环境里同时看到 NameNode、DataNode、ResourceManager、NodeManager 四个进程再用一段独立的 Java 代码去读写 HDFS最后才动手写 MapReduce。适合的人群是那些“代码会写、但不知道提交后到底发生了什么”的初学者以及想把这套流程固化成工程模板的开发者。这篇从环境、执行模型到参数坑逐个捋一遍。2. 实验环境搭建伪分布式 Hadoop 下 Java 操作 MapReduce 的工程基座2.1 伪分布式是实验报告里最稳妥的部署形态常见做法是搭一个 Hadoop 伪分布式集群也就是在一台物理机上同时启动 HDFS 的 NameNode、DataNode 和 YARN 的 ResourceManager、NodeManager。它比 local 模式多了一层完整的 RPC 和调度逻辑但比真正的多节点集群少了很多网络问题。实验课里如果直接上三节点集群光排查节点间免密和端口连通性就能耗掉你一个下午而这些跟 MapReduce 本身的编程能力没有关系。我一般会明确告诉学生除非题目要求 HA否则不要在这一步引入 Zookeeper。伪分布式环境里HDFS 的副本数必须显式设成 1。默认的dfs.replication3在单 DataNode 上会导致 DataNode 一直尝试复制 block日志里刷There are 0 datanode(s) running之类的信息作业虽然不会失败但报告里留下这种日志很难看。2.1.1 准备 JDK、SSH 与 Hadoop 安装目录先检查三个前置条件JDK 版本、SSH 免密登录、Hadoop 解压目录。JDK 建议用 8 或 11Hadoop 3.3.x 用 JDK 17 在部分场景会报Unsupported Class Files major version。SSH 免密用来支撑 HDFS 和 YARN 的本地进程启动脚本配置完可以用ssh localhost验证。然后编辑etc/hadoop/hadoop-env.sh把JAVA_HOME写成绝对路径。这一步很多人省略但通过ssh localhost拉起脚本时环境变量丢失NameNode 会启动失败错误信息还不直观。# 首次格式化 NameNode只需要执行一次 hdfs namenode -format # 启动 HDFS 与 YARN start-dfs.sh start-yarn.sh # 验证四个关键进程是否齐全 jps启动后jps至少要看到NameNode、DataNode、ResourceManager、NodeManager四个名字。如果只有Jps说明启动脚本因为环境变量或 SSH 问题中断了需要去logs/目录看对应进程的 log 文件。HDFS 的 Web UI 端口是 9870YARN 的资源管理器页面是 8088这两个页面能打开说明集群级别的服务已经正常。2.2 用 core-site.xml 和 mapred-site.xml 锁定客户端连接目标接下来是五个配置文件的职责边界这也是实验报告里“环境参数说明”板块最常扣分的地方。伪分布式的最小配置可以用下面这张表来对照检查文件核心属性作用hadoop-env.shJAVA_HOME让通过 SSH 启动的进程找到 JDKcore-site.xmlfs.defaultFS指定 HDFS 地址客户端 Java 代码靠它找到 NameNodehdfs-site.xmldfs.replication伪分布式必须把副本数改为 1mapred-site.xmlmapreduce.framework.name设为yarn让作业走 YARN 调度yarn-site.xmlyarn.nodemanager.aux-services设为mapreduce_shuffle否则 reduce 阶段拉不到数据core-site.xml里的fs.defaultFS是所有 Java 客户端连接的入口。你在代码里不显式写hdfs://localhost:9000时客户端会读取这个配置来决定文件系统地址。所以拿到一台实验机器第一件事是打开这个文件看fs.defaultFS的值而不是上来就写死 IP。很多“连接拒绝”的问题都是因为代码里写了localhost但实际 NameNode 挂在另一台机器的 IP 上。configuration property namefs.defaultFS/name valuehdfs://localhost:9000/value /property /configurationmapred-site.xml的mapreduce.framework.name这个属性很容易被忽略。IDE 里直接运行 main 方法时默认会退化成 local 模式文件读写走的是本地文件系统如果改成yarn后再在本地 IDE 里跑必须保证core-site.xml里的 NameNode 地址可达否则客户端会反复尝试连接。2.3 用 Java 代码连 HDFS 验证环境put 与 get 都通过再写 Mapper写 MapReduce 之前先写一段 20 行左右的 Java 代码验证 HDFS 连通性。这一步能把环境问题和业务逻辑问题彻底分开如果这段代码跑不通后面 Mapper 写得再对也没有意义。下面这段代码完成了文件上传和目录检查// HdfsProbe.java —— 用 Java 客户端连接 HDFS 做基本读写验证 import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.fs.FileSystem; import org.apache.hadoop.fs.Path; public class HdfsProbe { public static void main(String[] args) throws Exception { Configuration conf new Configuration(); conf.set(fs.defaultFS, hdfs://localhost:9000); FileSystem fs FileSystem.get(conf); Path src new Path(build/input.txt); Path dst new Path(/user/exp/input/input.txt); // 参数含义是否删除源文件、是否覆盖目标文件 fs.copyFromLocalFile(false, true, src, dst); Path dir new Path(/user/exp/input); System.out.println(exists fs.exists(dir) , isDirectory fs.isDirectory(dir)); fs.close(); } }代码里有两个细节值得写进报告copyFromLocalFile的第一个参数是 delSrc设为 false 保证源文件不被删除第二个参数是 overwrite实验中反复测试时设为 true 能省去手动删文件的步骤。在 Windows 上直接跑这段代码如果报Failed to locate the winutils binary说明缺少 Hadoop 的 Windows 本地库常见做法是切到 WSL2 或 Linux 虚拟机里运行。不要试图在 Windows 上硬绕过这个问题后面打 jar 提交时还会有更多权限和路径的坑。3. MapReduce 执行模型Java 代码里的 Mapper、Reducer、Combiner 分别在哪一步生效3.1 InputFormat 决定 map() 的输入边界行数不是你说了算MapReduce 的输入不是一个一个文件行而是由 InputFormat 划分出来的“分片”。FileInputFormat 默认按文件大小除以dfs.blocksize得到 split 数量每个 split 再经过 RecordReader 解析成(key, value)对喂给 map()。也就是说一个 200MB 的文件在 128MB 块大小下会被切成两个分片即使这个文件只有两行文本。这条规则直接影响你“Java 操作 MapReduce”的结果如果实验数据是大量几个 KB 的小文件每个文件都会产生独立的分片和 task启动开销远大于计算开销。常见做法是换用CombineTextInputFormat把多个小文件合并到同一个分片里。配置方式与默认的 FileInputFormat 完全不同需要在 main 方法里显式设置// 在 Job main 方法中加入合并小文件输入分片4MB 以内合并进同一个 map task job.setInputFormatClass(CombineTextInputFormat.class); CombineTextInputFormat.setMaxInputSplitSize(job, 4194304);这条配置在实验报告中非常有用因为实验数据经常是几十个小 JSON 或 CSV 文件拆分存放的。setMaxInputSplitSize的单位是字节4MB 是一个比较稳妥的阈值。另一个容易踩的细节是压缩格式gzip 压缩文件不可切分无论文件多大都只会启动一个 map task而 LZO 和 bzip2 支持切分。实验里如果看到某个 map task 跑得特别久先检查输入文件是什么压缩格式。3.2 shuffle 与二次排序Reducer 拿到的是分组后的迭代器map() 输出的(key, value)不会直接进入 reduce()中间隔着官方文档里最模糊的一段shuffle 和 sort。map 端的输出先写入环形缓冲区达到阈值后溢写到本地磁盘多个溢写文件在 map 端合并时做一次 sort。reduce 端拉取属于自己分区的数据后再做一次 merge 和 sort最终保证到达 reduce() 时相同 key 的 value 是连续排列的并且 key 整体有序。明白了这一点你就知道 reduce() 方法签名里为什么是IterableIntWritable而不是List框架已经完成了按 key 的分组。常见误区是写这样的代码在 reduce() 里先把所有 value 收集到 ArrayList再手动排序求 Top N。这样做功能上没错但内存开销和排序成本都白付了。二次排序的标准做法是构造“key排序字段”的组合键再通过自定义 WritableComparator 控制排序和分组的边界。以“统计同一用户最近一次登录时间”为例key 用(userId, loginTime)组合分区和分组都只按 userId但排序按完整组合键。这样 reduce 拿到的迭代器里第一个 value 就是该用户的最新登录记录。Java 实现里要重写三个方法getPartition按 userId 分区、compareTo定义排序、WritableComparator定义分组。此时不需要在 reduce() 里做任何遍历判断直接取iterator.next()就是结果。把这一步写进报告实验的“完整版”三个字才站得住。3.3 Combiner 和 Partitioner 的 Java 写法与默认值Combiner 是 map 端的局部 reducer它的作用是减少下游拉取的数据量。shuffle 过程要经过磁盘写入和网络拷贝context.write()发出的每一对 key-value 都可能落盘一次。设置 Combiner 后map 端会先对本地输出做一次预聚合下游拿到的数据量急剧下降。Combiner 有一个硬性约束它的输入和输出类型必须与 Reducer 一致并且运算必须满足结合律。以 WordCount 为例可以直接把 Reducer 类本身当作 Combiner 类因为累加操作无论分几步算结果都一样。但“求平均值”就不行map 端先算局部均值再在 reduce 端对均值求平均会得到错误的全局均值。这一点在实验报告里值得单独解释。Partitioner 决定一条记录进入哪个 reducer。默认实现是HashPartitioner源码逻辑非常简单// 默认分区逻辑对 key 的哈希值取模 return (key.hashCode() Integer.MAX_VALUE) % numReduceTasks;注意 Integer.MAX_VALUE这一步Java 的hashCode()可能返回负数直接取模会得到负数分区号进而抛Illegal partition异常。自定义 Partitioner 时最容易犯的错误是返回的分区号不在[0, numReduceTasks)区间内。下面这张表总结了三个组件在 Java API 中的位置组件继承的类调用时机常用配置Mapperorg.apache.hadoop.mapreduce.Mapper每个分片里逐条解析记录job.setMapperClassReducerorg.apache.hadoop.mapreduce.Reducer每组相同 key 调用一次job.setReducerClassPartitionerorg.apache.hadoop.mapreduce.Partitionermap 输出后、写入环形缓冲前job.setPartitionerClass这里还有个容易被忽略的配置job.setNumReduceTasks(0)。当你把 Reducer 数量设为 0作业变成 map-only输出文件命名也会从part-r-变成part-m-。如果你写了个实验逻辑只需要过滤和格式化输出map-only 作业会比强行塞一个空 Reducer 更高效。4. 实验核心用 Java 提交一个可复现的 WordCount 并逐参数解读4.1 完整工程代码Mapper、Reducer、main 放同一个类里实验课的代码结构不用太过设计一个 Java 文件里可以同时放 Mapper、Reducer 和 main 方法。这样打 jar、提交和阅读都方便。下面是一份可以直接编译运行的 WordCount 完整代码// WordCount.java —— Java 操作 MapReduce 的完整可运行示例 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; public class WordCount { // Mapper 类输入 key 是行偏移量value 是整行文本 public static class WordMapper extends MapperLongWritable, Text, Text, IntWritable { private final static IntWritable ONE new IntWritable(1); private final Text word new Text(); Override protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { String[] words value.toString().split(\\s); for (String w : words) { if (w.isEmpty()) continue; word.set(w); context.write(word, ONE); } } } // Reducer 类相同 key 的 value 已经由框架分组归并 public static class WordReducer extends ReducerText, IntWritable, Text, IntWritable { private final 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 output); System.exit(2); } Configuration conf new Configuration(); Job job Job.getInstance(conf, word count); job.setJarByClass(WordCount.class); job.setMapperClass(WordMapper.class); job.setCombinerClass(WordReducer.class); job.setReducerClass(WordReducer.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)是让框架定位包含代码的 jar如果你的 main 方法没有设置它提交到 YARN 后 task 执行时会报找不到类。Combiner 直接复用 Reducer 类前提是 WordCount 的逻辑满足结合律。waitForCompletion(true)的参数决定是否打印运行进度传 true 后客户端会每秒刷新进度条这也是观察作业是否卡住的主要手段。4.2 hadoop jar 提交与输出路径约束参数顺序错一个都不行代码写完后通过 Maven 打包然后提交作业。Hadoop 3.x 里hadoop jar和yarn jar在多数场景下等价但更规范的是用yarn jar。提交命令的参数顺序是固定的jar 包路径、主类全限定名、输入路径、输出路径。# 第一处Maven 打包跳过测试减少干扰 mvn clean package -DskipTests # 第二处把本地 jar 提交到 YARN 集群 yarn jar target/hadoop-exp-1.0.jar \ com.example.exp.WordCount \ /user/exp/input \ /user/exp/output/wordcount-01主类名必须紧跟 jar 包路径之后因为yarn jar会把第一个非选项参数当作主类名后面剩余的参数全部作为字符串传入 main 方法的 args 数组。实验中最常见的启动报错Error: Could not find or load main class一半是因为主类名拼写错误或 jar 内类路径不对另一半是因为把输入输出路径放到了主类名前面。输出路径有一个硬性约束必须不存在。Hadoop 为了避免覆盖历史结果在检测到输出目录已存在时会直接抛出FileAlreadyExistsException。所以每次重新运行前要么删除旧目录要么在命令里加上时间戳。waitForCompletion是同步阻塞的客户端进程会一直挂在终端直到作业结束很多新手看到终端没有输出就以为卡死了实际上打开 YARN 的 8088 页面能看到作业正在 RUNNING。和 Java 多线程里的 join 不同waitForCompletion内部是周期性轮询作业状态并拉取进度日志。4.3 用 Combiner 和自定义 Writable 扩展成“完整版”实验内容如果实验报告要求“完整版”WordCount 往往不够。常见的要求是按品类统计销售额总和与平均值。这时要自定义一个 Writable 类型承载组合结果因为Text和IntWritable无法同时存放金额总和与计数。下面是一个常规写法// SalesWritable.java —— 自定义 Writable承载 sum 与 count 两个字段 import org.apache.hadoop.io.Writable; import java.io.DataInput; import java.io.DataOutput; import java.io.IOException; public class SalesWritable implements Writable { private double sum; private int count; public void write(DataOutput out) throws IOException { out.writeDouble(sum); out.writeInt(count); } public void readFields(DataInput in) throws IOException { this.sum in.readDouble(); this.count in.readInt(); } // getter / setter 省略 }自定义 Writable 最核心的一点是write和readFields的字段顺序必须完全一致。这个类经过序列化后会在 map 端写出、在 shuffle 期间落盘、在 reduce 端反序列化任何一步的字段顺序不一致读出来的数据都是错乱的。Combiner 在这里就不能直接复用 Reducer 了因为“求和”和“求平均”语义不同需要单独写一个只做部分聚合的 Combiner 类。5. 排错与调参Java 操作 MapReduce 提交后最常见的 5 个报错5.1 ClassNotFoundException / NoClassDefFoundError 的根治fat jar 与 -libjars这个报错几乎每个人都遇到过。IDE 里跑得好好的打 jar 提交到集群就报ClassNotFoundException: com.example.exp.WordCount或NoClassDefFoundError: org/apache/hadoop/crypto。前者的根因是setJarByClass没有生效或者 jar 包里没有包含主类。后者是因为本机代码用到了某个第三方库但提交时没有把这个库传上去。常见做法是使用 maven-shade-plugin 把所有依赖打进一个 fat jar。这样本地依赖被合并进一个 jar 包集群端直接运行即可。如果不愿意打 fat jar也可以保留瘦 jar提交时用-libjars参数追加依赖# -libjars 需要配合 mapreduce.job.classloadertrue 使用 yarn jar hadoop-exp-1.0.jar com.example.exp.WordCount \ /input /output \ -libjars /lib/commons-lang3.jar在mapred-site.xml里把mapreduce.job.classloader设为 true 后-libjars的依赖才会生效。这个参数往往被忽略导致-libjars明明写了 jar 路径集群端仍然报类找不到。5.2 输出目录已存在、权限拒绝、虚拟内存超限的处理对照下面这张表收集了实验中出现频率最高的三个错误拿到终端日志后可以直接定位日志关键词根本原因处理方式Output directory ... already existsHadoop 拒绝覆盖已有输出删除旧输出目录或每次运行换新路径Permission denied用户名对目标 HDFS 目录无写权限hdfs dfs -chmod -R 777 /user或切换 HDFS 超级用户Container exited with a non-zero exit code 143容器物理内存或虚拟内存超过限制调大mapreduce.map.memory.mb与yarn.nodemanager.vmem-pmem-ratio虚拟内存超限这个坑最隐蔽。日志里出现 exit code 143 时第一反应往往是去看 Java 堆栈但实际原因可能是机器上还有别的进程占用了大量虚拟内存NodeManager 按比例估算虚拟内存用量时误杀了你的容器。常见解法是调大yarn.nodemanager.vmem-pmem-ratio从默认的 2.1 调到 3.0 以上或者把放大倍数关系反过来用yarn.nodemanager.vmem-check-enabled关闭虚拟内存校验。不过后者只是绕过问题生产环境不建议。5.3 看日志和 8088 页面定位慢任务而不是猜作业运行很慢时第一件事是找到 Application ID。终端会打印application_xxx_0001之类的 ID拿到 YARN Web UI 上搜索可以按 attempt 维度查看每个 task 的启动时间、结束时间和日志链接。单独看某个 task 的日志用命令# 拉取单个 application 的标准输出和错误日志 yarn logs -applicationId application_xxx_0001 -log_files stdout如果某个 map task 的耗时是其他 task 的十倍最常见的原因就是数据倾斜。这时候有两个调整方向第一是检查 Partitioner 是否对 key 做了均衡切分第二是关掉推测执行。推测执行在均衡负载下能加速但在数据倾斜场景下框架会重复启动一个慢任务的副本反而抢占了本该属于正常任务的资源。在mapred-site.xml中把mapreduce.map.speculative和mapreduce.reduce.speculative同时设为 false通常能让倾斜场景下的总耗时更稳定。6. 用 Counter 和输出文件把实验拼成一份可验证的完整报告6.1 统计 Counter 并回写到 HDFS实验报告需要“真实可信”的数据支撑。与其在截图里手写几行文字不如直接用 Hadoop Counter 统计作业运行指标。Counter 是框架内置的全局计数器可以在 Mapper 或 Reducer 里通过 context 自增最终由 Job 对象汇总。自定义枚举类型就可以创建新的计数项// 在 main 方法中定义计数分组 public enum ExpCounter { MAP_INPUT_RECORDS } // 在 Mapper.map() 中累加 context.getCounter(ExpCounter.MAP_INPUT_RECORDS).increment(1); // 作业跑完后在客户端读取 Counters counters job.getCounters(); long total counters.findCounter(ExpCounter.MAP_INPUT_RECORDS).getValue(); System.out.println([exp-report] total records total);Counter 的计数过程发生在 task 内部reduce 端也能读取天然适合作为实验结论数据。比如“统计出现次数超过 10 的单词数量”这种结论不需要在 reduce 里写完后再把结果拼接到外部文件直接定义一个 Counter 自增即可。客户端在waitForCompletion返回 true 后读取 Counter 值再把它写回 HDFS 的指定位置报告的“实验数据”部分就有了机器生成的来源。6.2 把 part-r-00000 和 application 日志变成报告附件实验报告的“完整版”往往需要把结果文件附在末尾。MapReduce 输出目录里的part-r-00000就是 reducer 产出的最终结果文件直接用hdfs dfs -cat拿回本地再按词频降序看 Top 20 结果# 从 HDFS 拉取结果并取数量最高的前 20 行 hdfs dfs -cat /user/exp/output/wordcount-01/part-r-00000 \ | sort -t $\t -k2 -nr | head -20结果文件的 key 和 value 之间是制表符所以排序时要把分隔符指定成\t。除了结果文件环境总览和作业日志也是报告素材用jps输出进程列表、用hdfs dfsadmin -report输出容量与副本状态、用yarn application -status输出作业的最终状态把这些命令的输出原样贴进报告的“运行环境”和“执行过程”章节。最后把这条命令存成脚本放进项目根目录名为run.sh实验数据就能一键复现不用每次手工敲一长串参数。本文还有配套的精品资源点击获取
锦
锦皓数字建站
深耕本土企业品牌数字化升级,专注原创端正雅致商务官网,从视觉设计到稳定运维全程保驾护航。