
简介一套基于Hadoop的气象数据分析完整工程代码面向大数据初学者与MapReduce开发者覆盖从HDFS分布式存储、MapReduce并行预处理到基于SSM框架的可视化展示全链路。资源包共562个文件压缩后约34.88MB主要包含Java源码Mapper/Reducer类、Service实现类、JSP页面、JavaScript与CSS前端样式以及XML配置和JAR依赖目录结构清晰便于直接对照学习。目前已有10577人学习下载。通过研读这套代码可掌握气象数据清洗、聚合统计如均温、极值等典型MapReduce写法并理解如何在Web端整合Spring Boot、Spring MVC与MyBatis呈现分析结果对搭建完整大数据分析项目具有直接参考价值。1. 用 Hadoop 啃气象数据先搞懂这份完整版代码在解决什么问题当手头的气象数据集从几兆涨到几个 G单机 Excel 和 Python 脚本就开始卡壳这正好是 Hadoop MapReduce 入场的时机。这份完整版代码做的事情并不复杂用 MapReduce 对气象观测记录做按年聚合统计出每一年的最高温度/平均温度并把中间结果落到 HDFS 上。适合两类人一是刚学完 Hadoop 基础、想找一份能跑通全流程作业的初学者二是需要在地理信息、农业、环境数据分析里做批处理但不想从零写 Mapper、Reducer、Driver 的开发者。整份代码麻雀虽小五脏俱全从环境配置、数据清洗到提交运行、结果验证都覆盖了值得照着跑一遍再改造成自己的分析逻辑。2. 环境与数据准备把 Hadoop 跑起来并让气象数据进 HDFS2.1 环境选型伪分布式已经足够跑通这份代码很多读者一上来就纠结要不要搭三节点集群我的建议是先别。这份气象数据分析代码的原始输入是某个公开气象观测数据包单机方案完全可以承载用 Hadoop 伪分布式模式即一台机器上同时跑 NameNode、DataNode、ResourceManager、NodeManager就能模拟完整分布式行为包括 HDFS 文件上传、MapReduce 任务调度、Shuffle 排序等核心环节。集群模式只多了一步把/etc/hosts配上多个节点、分发 JDK 和 Hadoop 安装包对于理解这份代码没有实质帮助。等你把伪分布式跑通了再迁移到集群只是改配置文件的事。操作系统上我建议 Linux不管是 Ubuntu 还是 CentOS避免在 Windows 上折腾 Cygwin 那套历史遗留问题。JDK 规格上这份代码里的 pom.xml 和运行脚本是按 Java 8 Hadoop 2.x/3.x 写的我实际跑的时候用的是 Hadoop 3.2.4 JDK 1.8没有遇到兼容性问题。如果你用的是 Hadoop 3.3 以上版本JDK 8 依然能跑但要注意yarn.nodemanager.resource.memory-mb默认配置在低配机器上可能导致 NodeManager 起不来这个在第 5 章避坑里细说。2.2 配置文件与启动流程core-site.xml 和 hdfs-site.xml 两个坑位Hadoop 安装包解压之后配置文件在$HADOOP_HOME/etc/hadoop/目录下。伪分布式只需要改两个文件第一是core-site.xml把默认文件系统指向 HDFS 的 NameNode 地址configuration property namefs.defaultFS/name valuehdfs://localhost:9000/value /property /configuration这个配置的意思是告诉所有 Hadoop 组件文件系统根路径在hdfs://localhost:9000。9000是 NameNode 默认的 RPC 端口如果你在同一台机器上还跑着其他服务占用这个端口需要换成 9001、9020 之类未占用的口子同时后续所有hdfs://路径都要跟着改这个一致性很容易忽略。第二个是hdfs-site.xml伪分布式下必须把副本数从默认的 3 改成 1否则 DataNode 只有一个节点副本数达不到 3HDFS 一直处于Under-Replicated状态虽不影响 MapReduce 跑任务但看着碍眼configuration property namedfs.replication/name value1/value /property /configuration改完这两个文件第一次启动前必须格式化 NameNode之后就不用再格式化了。格式化命令是hdfs namenode -format如果第二次运行还格式化会丢失之前 HDFS 上的所有数据这是新手最容易犯的错。然后依次启动 HDFS 和 YARN常见做法是start-dfs.sh和start-yarn.sh也可以用start-all.sh一把梭。启动完用jps命令确认进程能看到 NameNode、DataNode、ResourceManager、NodeManager 四个进程说明基本健康。提示jps是 JDK 自带工具只显示 Java 进程。看到进程不等于集群可用最好再执行hdfs dfsadmin -report看 DataNode 状态以及访问本机 9870 端口Hadoop 3.x或 50070Hadoop 2.x看 Web UI。2.3 气象数据集格式与上传先知道原始数据长什么样这份代码配套的模拟气象数据集模拟了某标准气象监测站的观测记录每行一条记录用逗号分隔字段结构是站点编号,日期,温度,湿度,风速。日期格式为yyyy-MM-dd HH:mm:ss温度精确到小数点后一位部分观测记录存在缺失字段或非法字符用来说明数据清洗的重要性。比如以下三行ST001,2021-07-15 12:00:00,34.5,62.3,3.6 ST001,2021-07-15 13:00:00,35.8,55.1,2.8 ST002,2021-07-15 14:00:00,,,5.0第三行温度和湿度字段为空是故意留下的脏数据用于验证 Mapper 的过滤逻辑是否有效。拿到数据后先丢到本地文件系统做一次预览用cat或wc -l看总行数然后用hdfs dfs -mkdir -p /weather/input创建 HDFS 上的输入目录再用hdfs dfs -put weather_sample.csv /weather/input/把数据上传到 HDFS 上。检查上传是否完整hdfs dfs -ls /weather/input/ hdfs dfs -du -h /weather/input/第一步列出文件确认文件名和大小是否与本地一致第二步看文件实际占用 HDFS 的物理空间。如果 DataNode 还没有把块复制完成du的结果会比ls里的文件大小小一些稍等几十秒再查一次即可。上传路径里的/weather/input可以任意修改但建议统一用/weather作为父目录把输入和输出分开避免 MapReduce 把输出目录误扫进输入。提示MapReduce 的输入路径默认递归扫描目录下所有文件如果输出目录也在输入路径下会导致任务跑完输出文件又被当成下一轮输入形成数据翻倍或无限递归的隐患。把/weather/input和/weather/output分开是整洁的习惯也是避免这类问题的物理隔离手段。3. MapReduce 核心代码Mapper、Reducer 与 Driver 的实现细节3.1 Mapper从原始行里提取年份和温度整份代码的核心是三个 Java 类第一个是 Mapper。MapReduce 的 Mapper 阶段要做的事就是逐行读取输入数据按业务规则做过滤和映射把结果写成 key-value 形式。气象数据统计每年的最高温度Mapper 要做的就是从日期字段里提取年份再把温度拿出来作为 value。代码里有两种处理路线一种是直接按逗号 split另一种是用自定义的字符串截取逻辑因为原始数据里的日期格式固定substring比 split 后校验更高效public class MaxTempMapper extends MapperLongWritable, Text, Text, IntWritable { Override protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { String line value.toString(); String[] fields line.split(,); // 字段不足 5 个说明是脏数据直接跳过 if (fields.length 5 || fields[1].isEmpty() || fields[2].isEmpty()) { return; } // 取日期时间的年部分例如 2021-07-15 12:00:00 取前 4 位 String year fields[1].substring(0, 4); // 温度转整数MapReduce 里用 IntWritable 而非 DoubleWritable为了排序更稳定 int temp (int) Math.round(Double.parseDouble(fields[2])); context.write(new Text(year), new IntWritable(temp)); } }代码里的LongWritable key是行偏移量MapReduce 框架自动生成不需要业务逻辑去解读Text value是这一行的完整字符串默认以换行符作为行的切分边界。context.write是输出入口key 是年份value 是温度整数。我在这里特意做了两个处理第一字段长度小于 5 的直接丢弃因为原始数据里存在字段缺失的记录第二Math.round把带一位小数的温度转成整数这样用IntWritable作为 value 类型在 Shuffle 阶段的排序比较速度远快于DoubleWritable代价是丢掉了小数精度。如果你的分析需要保留小数把IntWritable换成DoubleWritable即可但第 4 章的 Driver 里对应的类型声明也要同步换这个细节最容易漏。3.2 Reducer按年份聚合出最高温度Mapper 输出的中间数据经过框架的 Shuffle 排序后同一个 key也就是同一年份的所有温度会聚到一个 Reducer 的IterableIntWritable values参数里。Reducer 要做的就很简单了遍历这个迭代器维护一个最大值变量public class MaxTempReducer extends ReducerText, IntWritable, Text, IntWritable { Override protected void reduce(Text key, IterableIntWritable values, Context context) throws IOException, InterruptedException { int max Integer.MIN_VALUE; // 遍历同一年所有温度记录 for (IntWritable value : values) { max Math.max(max, value.get()); } context.write(key, new IntWritable(max)); } }这段代码没有做任何精妙的算法优化因为 MapReduce 的模型决定了大数据量下的复杂度已经被 Shuffle 扛住了Reducer 只需要做纯逻辑的聚合。需要注意到一点Integer.MIN_VALUE作为初始值是为了避免某个年份所有温度都是负数时初始值 0 会错误地把结果变成 0。这个细节在气象数据里尤其重要因为北方地区的冬季温度经常是 -20℃ 以下。如果初始值用 0那 2021 年 1 月的最高温度会被错误统计成 0除非恰好有一天是 0 度这种统计学上的边界条件经常被忽视。你也可以在 Mapper 里过滤掉极端异常值比如温度低于 -100 或者高于 60 的记录这属于气象数据的常识性阈值。3.3 DriverJob 配置与类型声明第三个类 Driver 是 MapReduce 作业的入口负责把 Mapper、Reducer、输入输出格式、路径等全部组装成一个可提交的 Job。代码核心逻辑如下public class MaxTempDriver { public static void main(String[] args) throws Exception { // args[0] 是输入路径args[1] 是输出路径 Configuration conf new Configuration(); Job job Job.getInstance(conf, Max Temp Analysis); job.setJarByClass(MaxTempDriver.class); job.setMapperClass(MaxTempMapper.class); job.setCombinerClass(MaxTempReducer.class); job.setReducerClass(MaxTempReducer.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); } }这里最值得解释的两个点是setCombinerClass和类型声明。setOutputKeyClass与setOutputValueClass指定的是最终输出类型也就是 Reducer 输出的 key 和 value 类型。但有一点需要注意如果 Mapper 的输出类型和 Reducer 的输出类型不一致比如 Mapper 输出Text, IntWritableReducer 输出Text, DoubleWritable那么必须在代码里额外用setMapOutputKeyClass和setMapOutputValueClass声明中间输出类型否则框架默认按最终输出类型处理运行时会报类型转换异常。这份代码里 Mapper 和 Reducer 的输出类型一致所以只声明了最终类型但不代表你改造后不需要补充声明。setCombinerClass用的是MaxTempReducer.class这是合法的因为求最大值这个聚合函数满足结合律先在一部分数据上算最大值再在最大值之间算最大值结果和全局一次求最大完全一致。Combiner 在 Mapper 端本地做一次预聚合减少写入磁盘和网络传输的数据量这是 MapReduce 优化里副作用最小的一步。很多优化手段会引入复杂的分区或排序逻辑Combiner 不会只要你的 Reducer 逻辑满足幂等和结合律直接复用就行。3.4 为什么要用 Text/IntWritable 而不是 Java 的 String/Integer刚接触 Hadoop 的人经常困惑Text和String长度一样IntWritable和Integer数值一样为什么非要绕一圈核心原因是 Writable 接口提供了序列化和反序列化机制MapReduce 框架要把 Mapper 输出的 key-value 写到磁盘和通过网络传输这个过程必须把对象转成字节流。Java 自带的Serializable接口过于重量级序列化结果里会附带很多类信息导致占用的存储空间大、传输效率低。Hadoop 的 Writable 机制只写入每个字段的二进制数据在排序阶段还能直接比较字节不需要反序列化成对象再调用compareTo。这就是为什么这份代码里所有 key 和 value 都用了 Writable 类型。如果你在 Mapper 里直接context.write(new String(year), new Integer(temp))编译能过但运行到 Shuffle 阶段就会报ClassCastException因为框架无法序列化它们。4. 编译打包到运行验证把 jar 提交到 Hadoop 集群4.1 Maven 打包与依赖作用域这份完整版代码配备了 Maven 的pom.xml核心依赖是hadoop-client打包半径很小。关键在于依赖的scope必须设置为provided否则会把 Hadoop 的 class 文件打进 jar 里导致整个 jar 体积膨胀到几十 MB提交到集群后还可能因为版本冲突跑出莫名奇妙的异常dependency groupIdorg.apache.hadoop/groupId artifactIdhadoop-client/artifactId version3.2.4/version scopeprovided/scope /dependencyprovided的含义是这个依赖在编译和测试时提供但运行时由运行环境提供也就是集群里的$HADOOP_HOME/share/hadoop目录。这样打包出来的 jar 只包含你的业务代码也就是三个类文件体积通常只有十几 KB。打包命令用mvn clean package如果本地仓库还没下载过 Hadoop 依赖会花几分钟时间拉依赖属正常现象。打包完成后 target 目录下会出现weather-analysis-1.0.jar之类的文件。提示如果你的 Maven 仓库下载慢可以在仓库镜像配置里换成国内镜像源这个和 Hadoop 本身无关但能省下不少等待时间。4.2 hadoop jar 提交命令与参数解释代码打成 jar 之后提交到 Hadoop 运行的命令如下hadoop jar target/weather-analysis-1.0.jar MaxTempDriver \ /weather/input /weather/output第一个参数是 jar 路径第二个是包含main方法的类名第三、四个是命令行传入的args[0]和args[1]分别对应输入输出路径。这里有几个容易踩坑的参数细节类名必须是MaxTempDriver的全限定名如果你在pom.xml里配置了finalName改变了 jar 名不影响类名。输入路径可以是目录但不能写成/weather/input/weather_sample.csv这样的文件名路径。虽然 Hadoop 允许直接指定单个文件但这份代码是面向批处理的目录路径可以一次性处理多文件。输出路径一定不能预先存在。这是 HDFS 的保护机制防止覆盖掉之前的结果。如果上一次任务跑出过结果第二次提交前先hdfs dfs -rm -r /weather/output。提交之后终端会开始滚动进度条显示 map 和 reduce 任务的百分比。Map 阶段跑完 100% 后Reduce 阶段从 0% 到 100% 通常需要一段时间因为 Reduce 要等待所有 Mapper 完成再拉取属于自己分区的中间数据。如果你在yarn-site.xml里没有设置过yarn.log-aggregation-enable任务跑完再用yarn logs -applicationId id查日志会拿不到内容建议提前配置为true这个参数控制运行日志是否汇总到 HDFS 上任务结束后才方便排查问题。4.3 查看运行结果与日志定位任务跑成功后输出目录下会产生一个part-r-00000文件用命令查看内容hdfs dfs -cat /weather/output/part-r-00000正常输出的每行格式是年份TAB温度比如2008 42 2009 38注意 MapReduce 默认的输出分隔符是 Tab 而不是逗号这个在代码里没有显式指定使用默认的TextOutputFormat。如果你在后续验证里用awk -F,去解析输出文件会发现分割失效这是很隐蔽的一个坑。正确做法有三个要么改代码里conf.set(mapreduce.output.textoutputformat.separator, ,)要么在解析脚本里改用 Tab 分割要么干脆在输出文件上先做一次sed s/\t/,/g再交给下游处理。我在实际项目里一般直接保留 Tab因为下游入库的 ETL 流程里 SQL 的LOAD DATA默认支持 Tab 作为列分隔符。如果任务中途失败优先看 YARN 的 Application 日志查询命令是yarn application -list -appStates FAILED yarn logs -applicationId application_xxx日志会打出每个容器里System.out输出和异常栈重点搜索 Mapper 或 Reducer 类名所在行大多数分析逻辑错误在main里System.exit(job.waitForCompletion(true) ? 0 : 1)不可见只有容器日志里有完整栈信息。5. 运行中的常见问题与避坑指南五条真实的翻车记录5.1 文件与路径类问题输出目录冲突和输入路径不存在现象一第二次运行同一个分析任务时提交后秒报file exists或者Output directory hdfs://localhost:9000/weather/output already exists错误任务直接拒绝启动。原因Hadoop 的FileOutputFormat出于安全考虑拒绝覆盖已有输出目录。这不像本地脚本里重定向能直接覆盖HDFS 上的目录一旦存在Job 初始检查就过不了。解决运行前手动执行hdfs dfs -rm -r /weather/output或者改Driver代码里的输出路径参数/weather/output_2。每次换新路径看起来有点笨但在实际生产任务里是一种好习惯因为有输出路径自带时间戳如/weather/output_20240512天然形成多版本结果。现象二输入路径写成本地文件系统路径如hdfs dfs -cat /weather/output/part-r-00000没问题但在提交 jar 时把输入路径写成file:///root/data/weather.csv报Input path does not exist: file:/root/data/weather.csv。原因数据上传时用了hdfs dfs -put这是往 HDFS 上传对应路径是hdfs://localhost:9000/weather/input/weather_sample.csv。而file:///前缀指向的是 Linux 本地文件系统数据不在本地自然找不到。解决提交 jar 时输入路径写成/weather/input这是最稳妥的写法Hadoop 会依据core-site.xml里的fs.defaultFS自动补全为完整 HDFS 路径。诊断时用hdfs dfs -ls /weather/input确认数据确实在 HDFS 上别用ls /weather/input后者查的是本地目录。5.2 数据与类型类问题脏数据破坏任务和序列化类型不匹配现象三某个年份算出的最高温度是 0而数据里这一年的温度全部是零下的。比如 2015 年冬季记录都是 -15℃、-22℃但 Reducer 输出2015 0。原因Reducer 的初始值代码里写死了int max 0然后逐个和温度取Math.max。由于 0 比所有负数都大循环结束后 max 保持 0。这不是Math.max写错了是初始值选取没有考虑负值域。解决把初始值改成Integer.MIN_VALUE这是求最大值聚合的标准安全初始化。同理如果是求最小温度初始值应该用Integer.MAX_VALUE。这个 bug 在气象数据里特别阴险因为南方数据全年正温度时跑不出来一换到北方数据集就中招属于典型的数据分布驱动代码缺陷。现象四Mapper 或 Reducer 里context.write传入的IntWritable转成了DoubleWritable但Driver里只声明了setOutputValueClass(IntWritable.class)任务跑到 Reduce 阶段报ClassCastException: org.apache.hadoop.io.DoubleWritable cannot be cast to org.apache.hadoop.io.IntWritable。原因MapReduce 在合并 Mapper 和 Reducer 输出时会按Job声明的输出类型做序列化和反序列化。当 Mapper 实际写入的类型与声明不一致Shuffle 阶段的反序列化直接失败。解决要么让 Mapper 和 Reducer 的输出类型保持统一并和setOutputClass保持一致要么在Driver里显式地分开设置中间类型和最终类型job.setMapOutputKeyClass(Text.class); job.setMapOutputValueClass(DoubleWritable.class); job.setOutputKeyClass(Text.class); job.setOutputValueClass(DoubleWritable.class);这个问题的排查思路是先看报错发生在 map 阶段还是 reduce 阶段。ClassCastException在 Map 阶段末尾出现多属setMapOutputClass没设置在 Reduce 处理时出现则优先怀疑setOutputClass和 Reducer 实际context.write类型不一致。现象五Mapper 里做Integer.parseInt(fields[2])时遇到空温度字段直接抛NumberFormatException整个任务失败跑不了几分钟就报Task failed。原因这份气象数据集里存在大量缺测记录比如第三行数据里的,,连续空字段split(,)之后fields[2]得到的是空字符串parseInt直接异常。MapReduce 的特性是一个 Task 失败会重试 4 次重试完了还没成功整个 Job 就失败不像单机脚本能打日志继续跑。解决在 Mapper 解析温度前加校验代码里现在有fields.length 5 || fields[1].isEmpty() || fields[2].isEmpty()的过滤isEmpty()已经把空字符串挡在门外。如果你用的是Double.parseDouble还要把fields[2].trim()放在isEmpty()判断之后因为气象数据里可能存在空格和不可见字符。这属于数据质量的防御性编码宁可多写两个条件也不要依赖数据源一定干净。6. 进阶Combiner 提速与本地结果交叉验证6.1 用 Combiner 减少 Shuffle 数据量伪分布式跑这份代码数据量不大时 Combiner 的效果看不出来。但如果你把数据源换到几年的全球气象观测记录Mapper 输出的中间键值对可能有几百万条Shuffle 阶段的网络传输和磁盘读写会成为瓶颈。这份代码里 Combiner 已经配置上了job.setCombinerClass(MaxTempReducer.class);原因在于求最大值是幂等聚合分片聚合再聚合结果一致。如果换成计算平均值就不能直接复用 Reducer 做 Combiner因为平均值不可结合。此时需要单独写一个 Combiner 类输出部分合计和计数Reducer 再做二次计算代码会复杂一截。这个区分是 MapReduce 优化里最常见的分界线建议在实际改造前先想清楚你的聚合函数是否满足结合律。6.2 本地结果交叉验证防止 HDFS 输出与真实值有偏差任务跑完后我习惯性在本地做一步交叉验证。把 HDFS 上的输入数据下到本地一份用 awk 按相同逻辑算最大温度再和 MapReduce 输出对比。对比命令很简单hdfs dfs -cat /weather/output/part-r-00000 mapreduce_result.txt awk -F, NR1{next} {yearsubstr($2,1,4); temp$3; if (temp0 max[year]) max[year]temp} END{for (y in max) print y, max[y]} weather_sample.csv | sort local_result.txt diff mapreduce_result.txt local_result.txt第一个cat命令重定向到文件第二个awk按逗号切分字段取年份和温度算最大值第三个diff看两端输出是否一致。如果diff输出为空说明 MapReduce 全流程跑通了如果输出有差异差在哪一年就聚焦那年的原始数据看是不是 Mapper 的缺失字段过滤规则和 awk 的temp判断不一致导致。这套交叉验证方法不需要额外引入任何工具Hadoop 命令行和系统自带的 awk、diff 就能完成却能在整合数据源和调整 Mapper 逻辑时提供充分的安全感。从那以后我每次改完 Mapper 的过滤条件或温度类型转换都强制走一遍 Hadoop 交叉验证流程省去了反复手工翻 HDFS 文件和猜答案的过程希望帮到你。本文还有配套的精品资源点击获取