ARTICLE DETAIL

资讯详情

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

Hadoop MapReduce实战:气象数据平均气温统计

Hadoop MapReduce实战:气象数据平均气温统计 简介面向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.8Hadoop 2.x/3.x均依赖Java运行环境Hadoop 2.10.x或3.3.x2.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 MapperLongWritable, 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 ReducerText, IntWritable, Text, DoubleWritable { private DoubleWritable result new DoubleWritable(); Override protected void reduce(Text key, IterableIntWritable 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转成实际摄氏度。这里有个隐含细节IterableIntWritable在遍历时每次返回的都是同一个对象引用所以不能把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类。第二setMapOutputKeyClass和setOutputKeyClass的泛型必须按“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-mb和mapreduce.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.mb和mapreduce.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做对比如果你只是想快速看个结果不写JavaHive会是更快的路径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 扩展3Python 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: passReducer脚本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_avgStreaming方案的核心逻辑和Java版完全一致因为底层Shuffle和Sort机制是共用的。区别仅在于Map阶段的输出通过stdin/stdout传递性能比Java略低但胜在开发效率高。写在最后的一点经验从实际项目的角度说气象数据分析这个场景特别适合用来入门Hadoop因为它数据格式规整、计算逻辑清晰、结果指标容易验证你能很直观地看到MapReduce各个阶段在做什么。我在第一次带着完整代码跑通这个任务时最大的感受是Hadoop本身并不难难的是你愿意沉下心去理解数据格式、理解Shuffle过程中每个环节的行为。如果你做的是课程设计建议在报告中重点写清楚“为什么固定宽度文本更适合用substring解析”“为什么要用IntWritable而不是String存温度”这些细节比堆砌大而全的架构图更能体现你的工程素养。照着本文的代码跑一遍再试着改一改需求和参数大概率比自己从零看官方文档有效得多。本文还有配套的精品资源点击获取
返回列表