Hadoop MapReduce实战:气象数据平均气温统计
2026/9/8 5:30:30 网站建设 项目流程

简介:面向Hadoop初学者与大数据开发者的完整气象数据分析实战资源,覆盖HDFS分布式存储、MapReduce并行计算、SSM框架Web展示全流程,可帮助理解气象数据预处理、统计计算与结果可视化。压缩包共562个文件,约34.88MB,包含Java源码、编译后的class文件、SSM框架jsp/css/js页面、xml配置及可直接部署的war包,目录层次清晰。已有10563人学习下载。资源以TemperatureMapper、TemperatureReducer等具体代码展示气温指标统计的实现思路,同时提供MinTemperatureServiceImpl等业务层示例,便于对照学习MapReduce任务编写、作业提交与Spring整合。对于正在做课程设计或准备大数据岗位面试的读者,这套完整工程具备较高参考价值。 做气象数据处理这个需求,很多人一开始想的都是直接用Python的Pandas一把梭。但在数据量真正上来之后,比如几十年的全球站点观测数据、逐小时级别的自动站数据,单机内存就会成为瓶颈。Hadoop生态里的HDFS分布式存储和MapReduce分布式计算,恰恰就是为了解决这类“数据躺在硬盘上但算不动”的场景而生的。

这篇博文不是讲理论,而是给出一份可以直接跑通的完整方案,覆盖从环境准备、数据格式分析、MapReduce编码到集群提交的完整链路。核心代码是Java版,因为在Hadoop课程设计、面试和真实生产环境里,Java还是主流语言,但文末会补充Python Streaming的替代方案。无论你是正在做课程设计,还是想在公司内网搭一套离线的气象数据统计任务,这篇文章都能作为一份可直接参照的落地手册。

1. 先搞清楚气象数据长什么样,再动手写代码

1.1 源头数据格式与字段语义

气象数据不是只有“温度”一个字段。以全球通用的NCDC(美国国家气候数据中心)数据集为例,每一行是一条独立的站点观测记录,按照固定宽度排列,常见字段包括:

  • 站点ID(如“029070-99999”,前6位是WMO区站号,后5位是扩展编号)
  • 观测日期(格式为“201701010000”,精确到小时)
  • 观测类型(如“TMIN”“TMAX”“PRCP”,分别表示最低温、最高温、降水量)
  • 观测值(气温单位是0.1摄氏度,降水量单位是0.1毫米)
  • 质量控制标志(如“1”表示合理,“9”表示缺失或错误)

这段宽度固定的文本,第一眼看上去非常“丑”,但恰恰是这种固定宽度格式最适合MapReduce处理——按偏移量切片即可,不需要复杂的序列化解析。

提示:如果拿到的是CSV或JSON格式的气象数据,解析方式会更常规,但MapReduce的框架思路完全一样。核心是:map阶段负责提取你关心的字段,reduce阶段负责汇总计算。

1.2 计算需求梳理:我们到底要算什么

在写任何一行Hadoop代码之前,先明确分析目标。最常见的两类需求是:

  • 年度平均气温统计:以“年份”为Key,统计当年所有站点、所有观测时次的平均气温
  • 月度极值统计:以“年份+月份”为Key,统计每月最高气温、最低气温,并追踪是哪个站点创造的

本文以“年度全球平均气温统计”为主线,因为它的逻辑最直观,适合作为第一课。月度极值作为扩展思路附在文末。

1.3 为什么选MapReduce而不是Hive或Spark

这个选择很关键。当数据量在TB级别以内、计算逻辑不复杂时,MapReduce虽然“笨重”,但胜在稳定、易调试、不依赖额外服务。Hive本质是把SQL翻译成MapReduce,适合非程序员;Spark则适合需要迭代计算和实时性要求高的场景。

在本例里,数据格式是固定宽度文本,解析逻辑完全可控,用原生MapReduce写,代码量不过几十行,运行效率反而比Hive的翻译层更可控。

2. 环境准备:Hadoop集群要先能跑起来

2.1 环境清单

执行代码前,先确认Hadoop环境已就位。以下是推荐的版本组合:

  • JDK 1.8+(Hadoop 2.x/3.x均依赖Java运行环境)
  • Hadoop 2.10.x或3.3.x(2.x对新手更友好,3.x对硬件资源占用更小)
  • Linux系统(CentOS 7或Ubuntu 18.04+均测试通过)
  • 至少1台机器即可完成伪分布式,3台以上可搭建真正的集群

2.2 伪分布式的启动顺序

如果你只有一台电脑,直接使用伪分布式模式,也就是让NameNode、DataNode、ResourceManager等角色都跑在同一台机器上。启动顺序必须严格按照以下命令:

hdfs namenode -format start-dfs.sh start-yarn.sh

启动后使用jps查看进程,确认存在NameNode、DataNode、ResourceManager、NodeManager四个关键进程。很多人的代码本身没问题,却在启动阶段就卡住了——最常见的原因是namenode -format只在第一次启动前执行,重复执行会导致NameNode的clusterID和数据节点的clusterID不一致。

2.3 数据上传到HDFS

假设本机气象数据文件名为weather_data.txt,先创建HDFS目录,然后上传:

hdfs dfs -mkdir -p /input/weather hdfs dfs -put weather_data.txt /input/weather/

上传后验证:

hdfs dfs -ls /input/weather

注意:不要直接读取本地文件路径。MapReduce的默认输入源是HDFS,如果你把本地路径传给FileInputFormat.addInputPath(),运行时会报FileNotFoundException,这是新手最容易踩的坑。

3. MapReduce核心代码:完整版Java实现与逐段解析

3.1 实体类与主类结构

本项目使用Maven工程,引入Hadoop Client依赖后,创建包名com.weather.analysis,下面有三个类:

  • WeatherMapper:负责map阶段解析
  • WeatherReducer:负责reduce阶段聚合
  • WeatherDriver:负责作业的配置与提交

依赖版本按实际集群调整。

3.2 Mapper:把每一行文本变成(年份, 气温)键值对

核心代码如下(完整版):

import java.io.IOException; import org.apache.hadoop.io.IntWritable; import org.apache.hadoop.io.LongWritable; import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.Mapper; public class WeatherMapper extends Mapper<LongWritable, Text, Text, IntWritable> { private Text outKey = new Text(); private IntWritable outValue = new IntWritable(); @Override protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { String line = value.toString(); if (line == null || line.isEmpty()) { return; } try { // 年份:从第0位开始取4位 String year = line.substring(0, 4); if (year.compareTo("1900") < 0 || year.compareTo("2100") > 0) { return; } // 气温整数部分:固定宽度,这里假设数据字段按照NCDC常用长度截取 // 实际开发中需要根据数据源调整偏移量 String tempStr = line.substring(30, 37).trim(); int tempInt = Integer.parseInt(tempStr); outKey.set(year); outValue.set(tempInt); context.write(outKey, outValue); } catch (NumberFormatException | StringIndexOutOfBoundsException e) { // 解析失败的行直接跳过,不中断作业 } } }

逐段解释几个关键点:

  • Mapper的四个泛型参数分别是:输入Key类型(偏移量)、输入Value类型(一行文本)、输出Key类型、输出Value类型
  • 年份判断加上“1900-2100”区间过滤,可以快速排掉文件头部注释行或乱行
  • 气温值用整数表示,因为Hadoop的Writable体系里没有FloatWritable,虽然可以有,但用IntWritable传输Int既省序列化开销又避免浮点比较精度问题。温度单位是0.1摄氏度,最终求平均后除以10即可
  • 异常处理非常重要。真实数据里一定有脏数据,一旦某一行解析失败,如果不捕获异常,整个Mapper任务会直接失败

3.3 Reducer:按年份聚合,计算平均气温

import java.io.IOException; import org.apache.hadoop.io.DoubleWritable; import org.apache.hadoop.io.IntWritable; import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.Reducer; public class WeatherReducer extends Reducer<Text, IntWritable, Text, DoubleWritable> { private DoubleWritable result = new DoubleWritable(); @Override protected void reduce(Text key, Iterable<IntWritable> values, Context context) throws IOException, InterruptedException { int sum = 0; int count = 0; for (IntWritable val : values) { sum += val.get(); count++; } double avg = count == 0 ? 0.0 : (sum * 1.0 / count) / 10.0; result.set(avg); context.write(key, result); } }

Reducer的逻辑很简单:同一个年份的所有气温值会汇聚到同一个Reducer方法中,遍历求和+计数,最后除以10转成实际摄氏度。这里有个隐含细节:Iterable<IntWritable>在遍历时每次返回的都是同一个对象引用,所以不能把val存进List再使用,必须当场计算或通过val.get()取基本类型。

3.4 Driver:配置作业的“胶水层”

import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.fs.Path; import org.apache.hadoop.io.DoubleWritable; import org.apache.hadoop.io.IntWritable; import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.Job; import org.apache.hadoop.mapreduce.lib.input.FileInputFormat; import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat; public class WeatherDriver { public static void main(String[] args) throws Exception { if (args.length != 2) { System.err.println("Usage: WeatherDriver <inputPath> <outputPath>"); System.exit(-1); } Configuration conf = new Configuration(); Job job = Job.getInstance(conf, "Weather Avg Temperature"); job.setJarByClass(WeatherDriver.class); job.setMapperClass(WeatherMapper.class); job.setReducerClass(WeatherReducer.class); job.setMapOutputKeyClass(Text.class); job.setMapOutputValueClass(IntWritable.class); job.setOutputKeyClass(Text.class); job.setOutputValueClass(DoubleWritable.class); FileInputFormat.addInputPath(job, new Path(args[0])); FileOutputFormat.setOutputPath(job, new Path(args[1])); System.exit(job.waitForCompletion(true) ? 0 : 1); } }

两个细节需要特别注意。第一,setJarByClass(WeatherDriver.class)是必须的,否则提交到集群时找不到Mapper和Reducer类。第二,setMapOutputKeyClasssetOutputKeyClass的泛型必须按“map端输出”和“reduce端输出”分别设置,当两者不一致时(本例一个IntWritable一个DoubleWritable),漏设任何一个都会在运行时抛IOException,错误提示还特别不直观。

3.5 用Maven打Jar包

在pom.xml所在目录执行:

mvn clean package -DskipTests

打包后,确认target目录下生成weather-analysis-1.0.jar。提交到集群的命令:

hadoop jar weather-analysis-1.0.jar com.weather.analysis.WeatherDriver /input/weather/weather_data.txt /output/weather_avg

提交后在YARN的ResourceManager界面(默认端口8088)可以看到作业运行进度,也可以在命令行执行yarn application -list查看运行状态。

4. 运行过程中的问题排查与调优心得

4.1 常见异常速查表

本人在多次课程设计和企业开发实践中整理的排查表如下,可以打印出来当备查:

异常信息原因分析解决方案
Invalid Kilobyte Configuration内存参数配置超出系统可用范围yarn.nodemanager.resource.memory-mbmapreduce.map.memory.mb改小一些
Container killed by ApplicationMasterMap或Reduce阶段内存溢出增加mapreduce.map.java.opts,但上限不能超过YARN容器限制
FileAlreadyExistsException输出目录已存在hdfs dfs -rm -r /output/weather_avg先删掉
ClassNotFoundException: com.weather.analysis.WeatherMapper没有设置Jar包或类路径不对检查job.setJarByClass,确认Jar包里有class文件
Input path does not exist传入路径在HDFS上不存在hdfs dfs -ls检查实际路径
GC overhead limit exceeded输入数据行数过多且每行解析有过多的临时对象减少字符串截取次数,避免创建过多String对象

4.2 数据倾斜问题

如果按年份聚合时,某一个年份的数据特别多(比如某一年的站点观测密度突然增大),会导致某个Reducer负载极高,拖慢整体作业。最简单的应对策略是增加Reducer数量:在Driver里加上job.setNumReduceTasks(4),让数据按哈希分发至4个Reducer。代价是同一个Key的结果会被拆分到多个文件,但如果下游是“求总平均”,可以在结果文件之外再聚合一次,问题不大。

4.3 输入小文件过多怎么办

气象数据经常是按年份、按站点切分的多个小文件,比如一个站点一年的数据就是一个几十KB的文本。HDFS上小文件过多会导致NameNode内存被大量Block占用。

临时的解决办法是在Driver里设置:

Configuration conf = new Configuration(); conf.set("mapreduce.input.fileinputformat.split.maxsize", "67108864"); // 64MB

这样会把小文件合并成一个InputSplit,减少MapTask数量。更根本的方案是提前用hdfs dfs -appendToFile或编写一个归并脚本,把一年或一个月的文件合并成一个大文件。

4.4 一个真实的调优案例

有一次在3节点集群跑4年逐小时气象数据,原始文件2.3GB,约3500万行。第一次跑,默认配置,耗时38分钟。后来做了三个调整:

  • 将Reducer数量从默认的1个改为4个,时间降到22分钟
  • 在Mapper里把字符串解析从substring改为复用char[]手动拼接,时间降到17分钟
  • 调整mapreduce.map.memory.mbmapreduce.reduce.memory.mb到1024MB,避免颠簸,总耗时最终稳定在15分钟左右

对于课程设计和大多数中小规模场景,第一项就够用了。后面两项属于锦上添花,但能体现你对Hadoop参数的理解深度。

5. 结果验证与扩展:从“能跑通”到“能应用”

5.1 检查输出

作业跑完后,查看输出文件:

hdfs dfs -cat /output/weather_avg/part-r-00000

输出格式应该是:

2010 16.4 2011 16.2 2012 16.6

这里看到的是“各年份所有站点观测到的气温平均值”,虽然不完全等同于气象学意义上的“全球平均气温”(因为站点分布不均,需要加权),但作为课程设计和日常统计已经足够。

5.2 扩展1:月度极值统计

如果需要统计“每年每月最高气温出现在哪个站点”,Map阶段的Key需要设计为“年份+月份”,Value为“温度”,Reduce阶段需要与当前最大值做比较,同时保留站点ID。一个常用的技巧是使用组合键,即Text类型year + "-" + month,但要注意在这种情况下,Reducer输入会是按组合键整体排序,而不是先按年份再按月份排序。这里更稳妥的做法是自定义WritableComparable,或者退而求其次,接受“月份不连续”的缺陷——大多数统计场景并不要求严格的字典序。

5.3 扩展2:和Hive做对比

如果你只是想快速看个结果,不写Java,Hive会是更快的路径:

CREATE EXTERNAL TABLE weather_data ( station STRING, date STRING, type STRING, value INT, quality STRING ) ROW FORMAT DELIMITED FIELDS TERMINATED BY '\t' LOCATION '/input/weather'; SELECT substr(date, 1, 4), avg(value / 10.0) FROM weather_data WHERE type = 'TMAX' GROUP BY substr(date, 1, 4);

Hive方案的开发速度确实快,但它的调优点和问题排查链路更长。当你已经掌握了MapReduce的原理,再使用Hive会非常顺手;反过来,如果连MapReduce都没跑通过,直接用Hive一旦出问题,会更难定位。

5.4 扩展3:Python Streaming方案

如果你只会Python,也不想学Java,可以直接用Hadoop Streaming,让Mapper和Reducer变成可执行的Python脚本。

Mapper脚本mapper.py

#!/usr/bin/env python import sys for line in sys.stdin: line = line.strip() if not line: continue try: year = line[0:4] temp_str = line[30:37].strip() temp = int(temp_str) print(f"{year}\t{temp}") except ValueError: pass

Reducer脚本reducer.py

#!/usr/bin/env python import sys current_year = None current_sum = 0 current_count = 0 for line in sys.stdin: line = line.strip() if not line: continue year, temp_str = line.split("\t") temp = int(temp_str) if current_year is None: current_year = year if year == current_year: current_sum += temp current_count += 1 else: avg = current_sum / current_count / 10.0 print(f"{current_year}\t{avg:.2f}") current_year = year current_sum = temp current_count = 1 if current_year is not None and current_count > 0: avg = current_sum / current_count / 10.0 print(f"{current_year}\t{avg:.2f}")

提交命令:

hadoop jar $HADOOP_HOME/share/hadoop/tools/lib/hadoop-streaming-*.jar \ -files mapper.py,reducer.py \ -mapper "python mapper.py" \ -reducer "python reducer.py" \ -input /input/weather/weather_data.txt \ -output /output/weather_python_avg

Streaming方案的核心逻辑和Java版完全一致,因为底层Shuffle和Sort机制是共用的。区别仅在于Map阶段的输出通过stdin/stdout传递,性能比Java略低,但胜在开发效率高。

写在最后的一点经验

从实际项目的角度说,气象数据分析这个场景特别适合用来入门Hadoop,因为它数据格式规整、计算逻辑清晰、结果指标容易验证,你能很直观地看到MapReduce各个阶段在做什么。我在第一次带着完整代码跑通这个任务时,最大的感受是:Hadoop本身并不难,难的是你愿意沉下心去理解数据格式、理解Shuffle过程中每个环节的行为。如果你做的是课程设计,建议在报告中重点写清楚“为什么固定宽度文本更适合用substring解析”“为什么要用IntWritable而不是String存温度”,这些细节比堆砌大而全的架构图更能体现你的工程素养。照着本文的代码跑一遍,再试着改一改需求和参数,大概率比自己从零看官方文档有效得多。

本文还有配套的精品资源,点击获取

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询