ARTICLE DETAIL

资讯详情

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

Hadoop实战:气象数据MapReduce完整代码与排错指南

Hadoop实战:气象数据MapReduce完整代码与排错指南 简介这套Hadoop气象数据分析完整代码面向大数据入门与进阶学习者解决气象数据分布式存储、并行计算与结果展示的实战需求。资源以MapReduce为计算核心结合SSM框架搭建Web展示层涵盖数据预处理、温度统计、最低温分析等典型场景。包内共562个文件约34.88MB主要包括Java源码与编译后的class文件、Spring/MyBatis配置、JSP/HTML页面、数据库SQL脚本以及部署所需的jar、war包目录结构贴近真实工程项目其中124个JavaScript、113个PNG、78个JPG、69个CSS和60个XML分别对应前端图表、页面样式、交互逻辑与框架配置。已有10563人学习下载。代码包含TemperatureMapper、TemperatureReducer、MinTemperatureServiceImpl、JobMain等关键类可直接运行获取平均气温、最高最低温度等统计结果完整工程也展示了Hadoop项目从数据清洗、MapReduce并行处理到SSM可视化的全链路实现并可通过自定义POJO返回、Service业务层拆分的写法快速迁移到其他气象指标分析场景适合作为课程设计或毕业设计的参考蓝本。 每年大数据课程设计季都会有一批人被“Hadoop分析气象数据”这个题目卡住。如果你正在搜“Hadoop分析气象数据完整版代码”我先把结论说清楚这套流程我带过好几届学生跑通不是换皮WordCount而是基于NOAA气象站观测记录做真实统计的完整项目——从数据下载、字段格式解析、MapReduce编写到打包提交、结果验证、常见报错排查一条龙全给你。为什么气象数据这么适合做Hadoop实战因为它是个“六边形战士”数据量大、格式规整但又有真实杂质缺失值、符号位、异常数据都齐了、统计结果完全可以人工验证。做完它你不光能交一份能跑的代码面试时被问到MapReduce的shuffle、Combiner、计数器也能从这套项目里掏出真实素材来讲。1. 为什么气象数据是Hadoop课程设计的“默认题”1.1 这个项目到底在解决什么问题气象数据本质上是一堆历史观测记录每天全球各地气象站都在产生新数据。分析需求通常很朴素某一年全球最高温度是多少每个气象站每月的平均气温怎么样某个区域近十年有没有变暖趋势这类统计天然适合“分而治之”Map阶段各管各的文件分片Reduce阶段再汇总。拿“统计每年最高气温”来说数据按行存储每行是一条观测记录包含年份、台站ID、温度、风速、降水量等字段。Map阶段逐行解析出年份和温度Reduce阶段把同一年份的所有温度汇总取最大值。用生活类比解释如果把全球几百年的气象观测记录打印出来堆起来比一栋楼还高。让你一个人从第一页翻到最后一页找“每年最高温度”翻到一半就崩溃了。MapReduce做的事是叫来一千个人每人分一摞纸各自挑出自己手里这摞纸中每一年的最高温度Map阶段再由一个小组长把某个年份的一千个小组结果汇总取最大Reduce阶段。1.2 一个完整方案由哪些部分组成整个项目不是只写一个Java类而是由数据、环境、代码、运行验证四块拼起来的组成说明气象观测数据建议用NOAA/NCEI发布的定长文本记录一行一条观测Hadoop环境伪分布式单机即可课程设计没必要上三台物理机开发环境JDK 1.8、Maven或直接用javac编译MapReduce程序Mapper解析数据Reducer做聚合Driver组装作业整体流程是数据处理前置验证 - 编写Mapper/Reducer/Driver - 编译打包 - 上传HDFS - 提交YARN - 查看结果。这里我要重点提醒很多同学上来就抄代码结果字段偏移量不对跑完全部数据出来一堆垃圾结果还以为是Hadoop算错了。所以第2节我先带你解决数据格式的问题。2. 实验环境与气象数据准备先看清记录长什么样2.1 Hadoop环境的最低要求这个项目用伪分布式模式完全够用。所谓伪分布式就是一台机器同时充当NameNode、DataNode、ResourceManager、NodeManager和真实集群的API调用完全一致只是规模小。准备一台Linux虚拟机或云主机内存建议至少4GB装好JDK 8和Hadoop 3.x。启动之前先格式化NameNodehdfs namenode -format然后启动HDFS和YARNstart-dfs.sh start-yarn.sh jps看到NameNode、DataNode、ResourceManager、NodeManager这4个Java进程都在环境就算OK了。很多同学卡在“Format后启动失败”十有八九是core-site.xml里HDFS地址没配对或者/tmp目录权限问题启动日志里会直接报。2.2 数据获取的两条路径第一路径是用官方数据。NOAA在NCEI官网ncei.noaa.gov提供GHCN-Daily、GSOD等气象数据集这是全球气象数据最权威的来源。但完整数据集很大如果一个课程设计就是从官网慢慢拖数据时间成本太高。第二路径是找课程设计常用的经典样例Hadoop权威指南里使用的那份1901、1902两个年份的气象记录网上很多博客都提供下载文件不大跑起来快格式也是教科书式的NCDC定长格式。如果这两条路都有障碍还有一个保底方法用Python生成一段指定格式的模拟数据。这样能先跑通流程逻辑没问题了再换真实数据。生成时只要保证年份和温度字段的偏移量固定即可import random with open(mock_weather.txt, w) as f: for year in [2001, 2002, 2003]: for _ in range(100): temp random.randint(-50, 350) # 实际温度乘以10 sign if temp 0 else - line f123456789012345{year}010 {sign}{abs(temp):04d}.... f.write(line \n)注意这只是一个格式占位示意真实记录远复杂得多但作为流程测试足够了。2.3 NCDC定长记录格式字段偏移量决定成败NCDC气象记录是定长文本每一列代表什么字段、从第几位开始官方文档里有完整定义。经典样例里最少要关注这三个字段字段起始列从1数结束列从1数Java截取方式0-based说明年份1619line.substring(15, 19)例如 1901温度符号8888line.substring(87, 88)或-温度数值8992line.substring(88, 92)实际温度乘以10例如 0049 表示 49温度为什么要拆成符号和绝对值因为NCDC里不同来源的记录符号位可能缺失或者出现空格。如果直接Integer.parseInt(line.substring(87, 92))遇到符号异常的行会直接抛异常整个Map任务失败。拆开处理的好处是解析不了符号时我们可以调用计数器把它记下来跳过这行不影响整体任务。2.4 一行命令验证字段偏移量写代码之前务必先验证你手里的数据符不符合这套偏移量。用一个Linux命令就能搞定head -5 1901 | awk {print substr($0,16,4), substr($0,88,5)}如果输出类似1901 0049 1901 -0011 1902 0317说明年份在第16到19列、温度在第88到92列偏移量是对的。如果输出乱码、明显错位就得查一查数据文档千万别急着往下写代码。这一步用十分钟后面能省十个小时。3. 完整版代码逐段拆解每年最高气温统计3.1 项目结构与pom.xml我用Maven管理项目工程结构如下weather-analysis/ ├── pom.xml └── src/main/java/ └── com/weather/ └── MaxTemperature.javapom.xml里的核心是引入hadoop-client依赖注意scope设为provided因为运行时Hadoop环境已经自带这些类不需要打进jar包?xml version1.0 encodingUTF-8? project xmlnshttp://maven.apache.org/POM/4.0.0 xmlns:xsihttp://www.w3.org/2001/XMLSchema-instance xsi:schemaLocationhttp://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd modelVersion4.0.0/modelVersion groupIdcom.weather/groupId artifactIdweather-analysis/artifactId version1.0.0/version packagingjar/packaging properties maven.compiler.source8/maven.compiler.source maven.compiler.target8/maven.compiler.target /properties dependencies dependency groupIdorg.apache.hadoop/groupId artifactIdhadoop-client/artifactId version3.3.6/version scopeprovided/scope /dependency /dependencies build plugins plugin groupIdorg.apache.maven.plugins/groupId artifactIdmaven-jar-plugin/artifactId configuration archive manifest mainClasscom.weather.MaxTemperature/mainClass /manifest /archive /configuration /plugin /plugins /build /projectHadoop 3.x的Java版本要求是8所以编译级别填8。3.2 Mapper的解析逻辑符号位比你想的重要下面这套代码是完整可运行的我把它全部塞进一个类方便你在课程设计里直接提交package com.weather; import java.io.IOException; 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; public class MaxTemperature { private static final int MISSING 9999; public static class MaxTemperatureMapper extends MapperLongWritable, Text, Text, IntWritable { private Text yearText new Text(); private IntWritable temperature new IntWritable(); Override protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { String line value.toString(); // 根据验证过的偏移量截取年份和温度 String year line.substring(15, 19); String sign line.substring(87, 88); String rawTemp line.substring(88, 92); int airTemperature; if (.equals(sign)) { airTemperature Integer.parseInt(rawTemp); } else if (-.equals(sign)) { airTemperature Integer.parseInt(rawTemp) * -1; } else { // 符号无法识别跳过并计数 context.getCounter(Weather, InvalidSign).increment(1); return; } // 9999 代表缺失值不能参与统计 if (airTemperature MISSING) { context.getCounter(Weather, MissingTemp).increment(1); return; } yearText.set(year); temperature.set(airTemperature); context.write(yearText, temperature); } } public static class MaxTemperatureReducer extends ReducerText, IntWritable, Text, IntWritable { Override protected void reduce(Text key, IterableIntWritable values, Context context) throws IOException, InterruptedException { int maxValue Integer.MIN_VALUE; for (IntWritable value : values) { maxValue Math.max(maxValue, value.get()); } context.write(key, new IntWritable(maxValue)); } } public static void main(String[] args) throws Exception { if (args.length ! 2) { System.err.println(Usage: MaxTemperature input path output path); System.exit(-1); } Configuration conf new Configuration(); Job job Job.getInstance(conf, Max Temperature); job.setJarByClass(MaxTemperature.class); job.setMapperClass(MaxTemperatureMapper.class); job.setCombinerClass(MaxTemperatureReducer.class); job.setReducerClass(MaxTemperatureReducer.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); } }Mapper里三个细节值得说明第一为什么符号要单独取出来判断因为Integer.parseInt(0049)在Java里能正常解析但遇到符号缺失或异常字符会抛NumberFormatException。把符号单独拿出来等于给数据加了一道质检工序解析不了的记录用计数器记录而不是让整个任务崩掉。第二9999这个值一定要过滤。NCDC规定的缺失温度就是9999如果不过滤哪年有缺失数据当年最高温度就会变成999.9度明显是畸形的。第三计数器Counter是排障利器。任务跑完后在日志里能看到Weather组里InvalidSign和MissingTemp各有多少一眼就知道数据质量靠不靠谱。3.3 Reducer与Combiner最大值可以提前合并Reducer的逻辑非常简单同一个key年份对应的所有温度逐个取最大值。代码里有一行容易忽略但很关键job.setCombinerClass(MaxTemperatureReducer.class);Combiner是Map端的“预聚合器”。最大值运算满足结合律和交换律max(max(a, b), c)和max(a, max(b, c))结果完全一样所以Reducer本身可以直接当作Combiner用。这样一来每个Map任务处理完自己那部分数据后先本地算一次最大值再通过网络shuffle传给Reducer数据量能大幅下降。比如一个Map处理了某个年份的5000条记录没有Combiner就得把这5000个(year, temp)全传给Reducer有了Combiner本地先算出一个最大值shuffle时只传1条。这就是为什么“直接复用Reducer做Combiner”是这类聚合场景的标准做法。3.4 Driver作业参数设置Driver做的事就是把组件拼起来设置Mapper类、Reducer类、Combiner类指定键值输出类型设置输入输出路径。job.waitForCompletion(true)会阻塞等待作业完成并且打印完整的作业进度和计数器信息。跑完后能看到类似这样的输出Weather MissingTemp12345 InvalidSign0这是判断数据质量的重要依据。4. 编译、打包与集群运行从本地到HDFS4.1 两种打包方式有Maven的情况下直接执行mvn clean package -DskipTests在target目录下会生成weather-analysis-1.0.0.jar。有些课程设计环境的网络不方便拉Maven依赖教你一招更皮实的直接用javac编译依赖路径从Hadoop里取mkdir -p weather_classes javac -classpath $(hadoop classpath) -d weather_classes src/main/java/com/weather/*.java jar cf weather.jar -C weather_classes .这种方式的等效性很好因为Hadoop环境已经把需要的jar都装在classpath里了。4.2 上传数据到HDFS并提交作业先把数据放上HDFShdfs dfs -mkdir -p /weather/input hdfs dfs -put 1901 1902 /weather/input/ hdfs dfs -ls /weather/input然后提交作业。这里有个反复被问的细节输出目录必须不存在Hadoop不会覆盖已有目录否则直接报FileAlreadyExistsExceptionhdfs dfs -rm -r /weather/output hadoop jar weather.jar com.weather.MaxTemperature /weather/input /weather/output在伪分布式模式下hadoop jar和yarn jar效果一样底层都是提交到YARN执行。4.3 查看输出与验证结果作业跑完后结果在HDFS的part-r-00000文件里hdfs dfs -cat /weather/output/part-r-00000输出大致是这个样子1901 317 1902 317 2001 312注意这里温度值是317实际摄氏度是31.7度因为气象记录里存的是真实温度乘以10的整数。怎么验证结果正确拿原始数据用grep或awk手工圈出每一年最大值对比MapReduce输出。如果对不上99%是字段偏移量错了回到第2.4节重新验证数据格式而不是怀疑Hadoop算错了。5. 排错记录这些坑我都替你踩过5.1 输出目录已存在错误信息形如org.apache.hadoop.mapred.FileAlreadyExistsException: Output directory hdfs://localhost:9000/weather/output already exists原因和解决方案说过了MapReduce为了防覆盖输出目录必须干净。重跑前先删掉旧目录或者换一个新的输出路径。5.2 NumberFormatException字段截取错位学生问得最多的问题就是这个异常。几乎每次都指向同一个根因substring的起止下标和实际数据不匹配截到了空格或者字母。排查步骤可以这样走先打印一行数据看看长度System.out.println(line.length());再用awk验证偏移量最后在Map里临时输出截取结果System.out.println([ year ] [ rawTemp ]);某个同学的案例他下载的NOAA新版本数据在第16到19列根本不是年份而是操作员ID结果解析出来的“年份”就是四位数字但温度解析全乱。问题不出在代码逻辑而是数据格式变了。5.3 符号位处理错误导致温度集体变负还有一个典型报错输出里某些年份的最高温度是负数明显不符合常识。原因基本是符号位判断写反了或者跳过符号直接把温度绝对值当成最终值。这个坑在气象数据里尤其隐蔽因为温度有正有负。如果判断符号时把和-写反北半球的夏季高温全变成负值。跑完结果先扫一眼发现最高温是负数第一反应就是符号位逻辑错了。5.4 native-hadoop library警告和内存不足伪分布式启动时会看到WARN util.NativeCodeLoader: Unable to load native-hadoop library for your platform...这个警告不影响任务执行原因是Hadoop自带的native库和本机系统不匹配。如果只是做课程设计可以不用管它忽略即可。如果任务运行中Container被kill日志里出现内存相关错误检查yarn-site.xml的两个参数property nameyarn.nodemanager.resource.memory-mb/name value4096/value /property property nameyarn.scheduler.maximum-allocation-mb/name value4096/value /property把内存调大一些或者调小mapreduce.map.memory.mb和mapreduce.reduce.memory.mb让每个任务占用的内存降下来。6. 从60分到90分课程设计的扩展思路6.1 按气象台站分组统计做完“每年最高气温”这个经典案例如果你想在课程设计答辩里多讲几句最简单的扩展是把统计维度换成台站ID、年份的双重分组。做法Map输出key用台站ID 年份拼接成的Text例如029070-99999_1901Reduce逻辑完全不用改。这样就能得到“每个台站每一年的最高气温”比单一“每年最高气温”信息量大得多。6.2 月平均气温统计月平均气温是另一个经典需求。思路是Map阶段的key用年份 月份Reducer里求和再除以记录条数。但这里有个加分知识点平均值不能像最大值那样直接用Reducer当Combiner。原因很简单avg(avg(1, 100), avg(50, 50))不等于avg(1, 100, 50, 50)平均值不满足结合律。如果非要加Combiner就得用自定义的聚合结构把每组的总和和计数一起传而不是只传平均值。这个细节在答辩时讲出来老师会知道你确实理解了Combiner的原理。6.3 用Partitioner控制输出分区默认分区器按key的哈希值散列输出顺序没有规律。如果想把结果按年份分区输出到不同文件可以自定义Partitionerpublic static class YearPartitioner extends PartitionerText, IntWritable { Override public int getPartition(Text key, IntWritable value, int numPartitions) { return Integer.parseInt(key.toString()) % numPartitions; } }这样配置了多个Reducer时每一年份都会固定进入同一个分区便于按年份归档结果。6.4 数据倾斜与Map数量的控制气象数据如果按台站归档某些大城市的台站记录数可能远超其他台站Reduce阶段就会出现一个Reducer处理海量数据、其他Reducer空转的情况。这就是数据倾斜。缓解办法一是加一个随机key做两阶段聚合代价是业务逻辑变复杂二是调整Reduce任务数量或者用Hive、Spark这类更高层框架去处理而不是继续手写MapReduce。最后分享一点个人经验做这个项目最忌讳的就是拿到数据就甩开膀子写代码。先花十分钟看数据长什么样、字段偏移量对不对、有没有缺失值后面所有环节都会顺很多。把这个流程跑通再去研究YARN调度、HDFS读写机制这些底层原理你会发现自己对MapReduce的理解一下子立体起来了。本文还有配套的精品资源点击获取
返回列表