简介一套围绕Hadoop生态的气象数据分析完整代码工程面向大数据入门学习者和相关课程设计人员覆盖分布式存储、并行计算与Web可视化的全流程。资源包共包含562个文件大小约34.88MB既有源码、编译后的类文件、依赖库等后端核心也有大量脚本、样式、动态页面和图片等前端展示资源还附带数据库脚本与配置文件结构清晰便于按模块查阅。这套代码目前已有10570人学习或下载热度较高适合作为课程设计、毕业设计或大数据实训的参照素材。代码中的MapReduce作业与业务服务类清晰展示了从原始气象数据清洗、分组聚合到业务封装的完整实现结合前端页面与SSM整合配置能够掌握前后端联调、数据查询以及项目部署的常用方法为独立开展同类大数据分析项目提供可复用的实战参考。1. 先说清楚用Hadoop分析气象数据这条链路到底解决什么问题很多人第一次做Hadoop分析气象数据的课程设计第一反应是赶紧写一个MapReduce类把CSV读进来跑出平均数就算完事。我的建议相反先想清楚统计口径和数据怎么放再写代码。标题里的“完整版代码”指的是从原始气象文件上传HDFS开始到MapReduce作业提交到YARN再到结果取回本地并验证结论的整条闭环而不是单个Mapper类。这套链路适合刚完成Hadoop伪分布式搭建、手里有一批真实气象站点数据、想拿MapReduce练手的同学也适合作为Hive作业、Spark作业的对照基准。一旦跑通数据切分、内存配置、数据倾斜这些问题就都有了一个可以反复复现的试验场。2. 气象数据怎么组织HDFS目录设计与文件格式选择气象数据分析最常翻车的地方不在MapReduce代码而在数据一进门就乱了。这一章先把文件层面的事情定下来后面的代码才能少改。2.1 真实气象数据长什么样字段、量级与脏数据气象站导出的数据通常是CSV或固定列宽文本。我用的报表格式是每行一个站点一天的气象要素列顺序固定。这里给出一个通用结构后面的代码都按这个列序读取列号字段名示例值说明0station_id54308站点编号字符串1date20200105日期YYYYMMDD2avg_temp8.6日平均气温摄氏度3max_temp15.2日最高气温4min_temp1.1日最低气温5precipitation2.5日降水量毫米6quality0质量控制码0表示有效一行样例大致长这样54308,20200105,8.6,15.2,1.1,2.5,0 54308,20200106,9.1,16.0,2.3,0.0,0拿到手的数据通常不会这么干净常见脏数据包括表头行混在数据里、气温字段出现-9999或9999表示缺测、降水字段为空、同一站同一天出现重复记录。这些值如果不提前处理后面算出来的年平均气温可能是零下几百摄氏度一眼假。数据量上一个省级区域的800个自动站一年日值数据大约29万行无压缩CSV约十几MB如果积累十年也就两三个GB。这个量级对HDFS来说很小但恰恰因为小很多人会忽略目录和文件组织把所有数据零零散散put进去最后在NameNode里留下几万个小文件为后续作业埋坑。2.2 HDFS目录设计按年份和月份分区不要让MapReduce扫全表MapReduce对输入路径的处理方式是递归读取目录下所有文件。如果把三十年数据放在同一个扁平目录每次分析都要把全部数据扫描一遍哪怕只需要其中一年。更合理的做法是把原始数据按时间分区存放/weather/raw/year2020/month01/site_54308.csv /weather/raw/year2020/month02/site_54308.csv /weather/raw/year2020/month12/site_54308.csv /weather/raw/year2021/month01/site_54308.csv建目录并上传的常用命令如下hadoop fs -mkdir -p /weather/raw/year2020/month01 hadoop fs -put site_54308_202001.csv /weather/raw/year2020/month01/后续提交作业时只把/weather/raw/year2020作为输入路径MapReduce只会读取2020年的子目录。这个设计在数据量小的时候看不出区别但到了几十GB甚至TB级分区裁剪能让作业执行时间从小时级降到分钟级。HDFS层面还有一个隐性收益分区后的目录天然成为数据访问边界后续用Hive建分区表时可以直接把/weather/raw映射成外部表ALTER TABLE ... ADD PARTITION不需要重新导入数据。提示如果你把每个站点的数据拆成上千个小文件处理它们的时间可能比分析本身还长。小文件问题在气象场景里很常见。每个站点一天的CSV可能只有几百字节但HDFS默认Block大小是128MB一个300字节的文件也要占一个Block的元数据。文件数超过NameNode内存承载能力后集群会变得异常慢。处理方法是在本地先把小文件合并再一次性上传mkdir -p merged for f in data_2020/site_*.csv; do if [ ! -f merged/2020_all.csv ]; then cp $f merged/2020_all.csv else tail -n 2 $f merged/2020_all.csv fi done hadoop fs -mkdir -p /weather/raw/year2020 hadoop fs -put merged/2020_all.csv /weather/raw/year2020/上面的脚本用tail -n 2去掉除首个文件外的表头避免合并后出现大量表头行。如果发现合并后同一站点同一天有重复记录统计前必须做一次合并去重否则均值、总量都会被放大。最简单的去重是按“站点编号日期”去重awk -F, NR1 || !seen[$1_$2] merged/2020_all.csv merged/2020_dedup.csv2.3 文件格式选择CSV、SequenceFile 还是 Parquet处理气象数据时文件格式的选择直接影响作业能否并行读取和压缩效率。简单对比格式可读性压缩率单个文件是否可切分适用场景纯CSV高低是原始数据、排错、课程作业gzip压缩CSV解压后可读中否归档、小集群SequenceFile低中是MapReduce中间结果、小文件合并Parquet低高是数据量大、列式查询、后续接Hive/Spark我的建议是课程设计级项目直接用CSV保留原始可读性如果磁盘吃紧可以给CSV文件单独做gzip压缩。有一个很容易踩的坑值得提前说gzip压缩后的文件在HDFS里是不可切分的一个gzip文件只能被一个Map任务读取。如果合并后生成一个2GB的gzip文件所有数据会被单个Map串行处理分布式并行能力等于没有。规避办法是上传前把数据切分成多个1GB以内的文件每个分别gzip让不同文件跑在不同Map上。SequenceFile适合把大量小文件打包成少数几个顺序文件但用户可读性差Parquet则是Hive和Spark场景下的首选列式存储对按字段聚合的SQL特别友好。如果你预期后续还要用Hive做复核分析Parquet值得投入如果只是跑一次MapReduce出结论用CSV足够不要为了格式而格式。3. 完整版MapReduce代码按“年-站点”统计气象要素这一章给出一个能直接编译运行的完整Java类实现“年-站点”维度的年平均气温、最高/最低气温极值和累计降水量统计。代码里刻意把Combiner和Reducer分开写因为平均值统计的Combiner有边界问题后面单独说明。3.1 数据口径与列号约定先定规则再写代码写代码前必须把统计口径定死。我的口径是只统计quality0的记录气温字段中-9999、9999视为缺测整条记录跳过降水字段缺测当天按0毫米计入但不参与有效样本计数。所有输出的平均气温保留两位小数输出格式使用Tab分隔方便后续用Hive直接加载。下面按列下标取字段所以表头和字段顺序必须和2.1节的表格保持一致。如果实际数据列序不同预处理阶段需要先做一次列重排而不是在Mapper里到处改下标。3.2 完整源码Mapper、Combiner、Reducer与Driverimport org.apache.hadoop.conf.Configuration; import org.apache.hadoop.fs.Path; 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 WeatherAnalysisMR { private static final double MISSING -9999.0; public static class WeatherMapper extends MapperObject, Text, Text, Text { private Text outKey new Text(); private Text outVal new Text(); Override protected void map(Object key, Text value, Context context) throws java.io.IOException, InterruptedException { String line value.toString(); String[] f line.split(,, -1); if (f.length 7) return; if (f[0].trim().equals(station_id)) return; // 跳过表头 String date f[1].trim(); if (date.length() 8) return; double avgTemp parseTemp(f[2]); double maxTemp parseTemp(f[3]); double minTemp parseTemp(f[4]); double precip parseTemp(f[5]); // 缺测值不参与统计任一气温缺测则丢弃整行 if (Double.isNaN(avgTemp) || Double.isNaN(maxTemp) || Double.isNaN(minTemp)) return; if (Double.isNaN(precip)) precip 0.0; String year date.substring(0, 4); outKey.set(year - f[0].trim()); // 中间格式count, avgSum, max, min, precipSum, precipCount outVal.set(1\t avgTemp \t maxTemp \t minTemp \t precip \t (precip 0 ? 1 : 0)); context.write(outKey, outVal); } private double parseTemp(String s) { s s.trim(); if (s.isEmpty()) return Double.NaN; double v Double.parseDouble(s); if (v MISSING || v 9999.0) return Double.NaN; return v; } } public static class WeatherCombiner extends ReducerText, Text, Text, Text { Override protected void reduce(Text key, IterableText values, Context context) throws java.io.IOException, InterruptedException { double[] agg aggregate(values); Text outVal new Text(); outVal.set((int) agg[0] \t agg[1] \t agg[2] \t agg[3] \t agg[4] \t (int) agg[5]); context.write(key, outVal); } protected double[] aggregate(IterableText values) { int count 0; double sum 0.0, max -Double.MAX_VALUE, min Double.MAX_VALUE; double precipSum 0.0; int precipCount 0; for (Text val : values) { String[] p val.toString().split(\t); count Integer.parseInt(p[0]); sum Double.parseDouble(p[1]); max Math.max(max, Double.parseDouble(p[2])); min Math.min(min, Double.parseDouble(p[3])); precipSum Double.parseDouble(p[4]); precipCount Integer.parseInt(p[5]); } return new double[]{count, sum, max, min, precipSum, precipCount}; } } public static class WeatherReducer extends WeatherCombiner { Override protected void reduce(Text key, IterableText values, Context context) throws java.io.IOException, InterruptedException { double[] agg aggregate(values); double avg agg[1] / agg[0]; String result String.format(%.2f\t%.2f\t%.2f\t%.2f, avg, agg[2], agg[3], agg[4]); context.write(key, new Text(result)); } } public static void main(String[] args) throws Exception { Configuration conf new Configuration(); Job job Job.getInstance(conf, weather-year-station); job.setJarByClass(WeatherAnalysisMR.class); job.setMapperClass(WeatherMapper.class); job.setCombinerClass(WeatherCombiner.class); job.setReducerClass(WeatherReducer.class); job.setOutputKeyClass(Text.class); job.setOutputValueClass(Text.class); // 4个Reduce任务输出4个part文件适合后续加载 job.setNumReduceTasks(4); FileInputFormat.addInputPath(job, new Path(args[0])); FileOutputFormat.setOutputPath(job, new Path(args[1])); System.exit(job.waitForCompletion(true) ? 0 : 1); } }代码里的parseTemp把-9999、空字符串、9999统一转成Double.NaNMapper只对气温有效的行做统计这是防止结果出现离谱负数的第一道闸门。Mapper输出value不是最终结果而是count、sum、max、min、precipSum、precipCount六个中间量的拼装串这样Combiner或Reducer拿到后可以继续聚合不会丢失精度。注意Reducer继承了WeatherCombiner直接复用aggregate聚合逻辑但输出时做了一步avg sum / count并格式化字符串。Combiner输出的格式和Mapper输出保持一致都是中间序列格式这保证了Combiner可以在Map端被多次执行而不改变正确性。3.3 平均值的Combiner边界为什么不能对平均值再求平均Combiner是Map端的局部Reducer它的输出会作为Reducer的输入所以必须满足“可结合、可交换”的约束。对求和、求最大值、最小值来说局部汇总结果不影响全局结果但平均值不行——两个片段各自的平均温度不能直接相加除以二因为片段里的样本数不同。举例来说Map任务A产生10条记录的片段平均气温为10度Map任务B产生2条记录的片段平均气温为20度全局正确平均值是(10*10 2*20) / (102) 11.67。如果Combiner直接把两个片段平均再平均结果就是(1020)/215偏差明显。因此这里Combiner传递的必须是count sum而不是平均温度本身最终平均值只在Reducer最后一步计算。这一点在气象统计里特别容易被忽略也是面试和答辩时高频追问的点。3.4 参数设置Reduce数量与序列化类型setNumReduceTasks(4)是刻意设成4的。这个作业的数据量在几十MB量级Reduce任务太多会产生大量小文件太少则并发不够。4个Reduce对应4个part-r-00000到part-r-00003文件后续用getmerge合并结果或按分区加载都很方便。setOutputKeyClass和setOutputValueClass设置的是Reducer最终输出类型也就是输出文件里键值对的序列化格式。如果Mapper输出的键值对和Reducer不同还需要额外调用setMapOutputKeyClass和setMapOutputValueClass声明否则作业在Shuffle阶段会因反序列化类型不匹配直接报错。这个例子中Mapper、Combiner、Reducer的键值对类型都是Text, Text所以可以只设一处。4. 从打包到跑通提交气象分析作业的完整命令序列代码写完不等于能跑出结果。从Maven打包开始到HDFS上传、YARN提交、日志排查、结果取回每一步都有坑。这一章按顺序给出可复制的命令序列。4.1 Maven打包依赖作用域与插件配置IDE里直接运行MapReduce作业经常翻车原因是Eclipse链接配置Hadoop时classpath和Hadoop环境变量没对上跑起来要么报ClassNotFoundException要么访问不到HDFS。命令行提交是最可控的方式。先准备pom.xmlHadoop客户端依赖使用provided作用域避免把整个Hadoop打进业务jardependencies dependency groupIdorg.apache.hadoop/groupId artifactIdhadoop-client/artifactId version3.3.4/version scopeprovided/scope /dependency /dependencies build plugins plugin groupIdorg.apache.maven.plugins/groupId artifactIdmaven-jar-plugin/artifactId configuration archive manifest mainClassWeatherAnalysisMR/mainClass /manifest /archive /configuration /plugin /plugins /buildprovided的含义是编译时可以使用Hadoop API打包时不包含Hadoop依赖运行时由hadoop jar命令提供集群classpath。这样打出来的jar只有几十KB上传快也不会和集群自带依赖冲突。打包命令mvn clean package -DskipTests ls target/weather-analysis-1.0.jar如果项目里用了Commons CSV这类第三方库需要把maven-shade-plugin加进来把第三方类合并进fat jar否则运行时会报NoClassDefFoundError。这个细节等到5.1节真正用到时再处理。4.2 上传数据到HDFS目录检查与输出路径清理数据上传在2.2节已经做过一次。这里强调一个执行顺序问题每次运行作业前必须确认HDFS上输出目录不存在。MapReduce框架禁止覆盖已有输出目录第二次跑同一作业时会直接抛FileAlreadyExistsException。我通常把上传和清理写成一个shell脚本#!/bin/bash INPUT/weather/raw/year2020 OUTPUT/weather/output/result hadoop fs -mkdir -p $INPUT hadoop fs -put -f merged/2020_dedup.csv $INPUT/2020_all.csv # 清理上一次结果目录 hadoop fs -test -d $OUTPUT hadoop fs -rm -r $OUTPUT # 提交作业 hadoop jar target/weather-analysis-1.0.jar \ weather.analysis.WeatherAnalysisMR \ $INPUT $OUTPUTput -f用于本地文件已存在时强制覆盖避免重复上传报错。清理输出目录放在提交前而不是代码里是因为正式环境里作业不应该有删除HDFS目录的权限课程作业图省事可以把删除逻辑放进Driver但我不建议养成这个习惯。命令行里的-D参数必须放在hadoop jar命令名之后、jar包之前写法如下hadoop jar -D mapreduce.job.reduces8 target/weather-analysis-1.0.jar \ weather.analysis.WeatherAnalysisMR /weather/raw/year2020 /weather/output/result这和hadoop jar target/xxx.jar -D的写法效果不同后者会把-D当成main方法的参数Configuration里读不到作业还是按默认值执行。这个位置问题让不少人排查了半天。4.3 观察作业状态与应用日志作业提交后终端会打印Map和Reduce的进度条。伪分布式环境下如果进度一直卡在0%或者容器反复重启优先看YARN日志目录ls $HADOOP_HOME/logs/userlogs/ # 例application_1710000000000_0001 find $HADOOP_HOME/logs/userlogs/ -name stderr -mmin -5 | headstderr文件里通常是Java异常栈syslog里能看到容器启动和GC信息。比较隐蔽的情况是作业已经跑完但终端没刷新这时去YARN ResourceManager界面看Application状态最直接或者用命令查yarn application -list -appStates FINISHED,FAILED,KILLED如果拿到Application ID可以拉取完整日志yarn logs -applicationId application_1710000000000_00014.4 伪分布式常用内存参数4GB机器能跑起来的配置Hadoop默认参数面向生产集群伪分布式单机常因内存不足让容器被ResourceManager杀掉。从零开始Hadoop安装和配置时很多人漏改yarn-site.xml导致作业一提交就失败或超慢。适合4GB内存机器的参数组合如下参数建议值说明mapreduce.framework.nameyarn强制走YARN而不是local模式yarn.nodemanager.resource.memory-mb3072NodeManager可用总内存mapreduce.map.memory.mb512单个Map容器内存mapreduce.map.java.opts-Xmx384mMap JVM堆内存留部分给元空间mapreduce.reduce.memory.mb512单个Reduce容器内存mapreduce.reduce.java.opts-Xmx384mReduce JVM堆yarn.app.mapreduce.am.resource.mb1024ApplicationMaster容器内存mapreduce.job.reduces4本作业Reduce数量这几项在yarn-site.xml里配置后重启YARN才会生效。Map和Reduce堆内存比容器内存小128MB左右是为了给JVM非堆区域留空间否则容器会因超出内存限制被kill翻车现场通常是日志里出现Container killed by ResourceManager。4.5 结果取回getmerge合并多分区输出Reduce数量为4时HDFS输出目录下会有part-r-00000到part-r-00003四个文件。逐个put很多余用getmerge一步到位hadoop fs -getmerge -nl /weather/output/result weather_result.csv wc -l weather_result.csv head -n 10 weather_result.csv-nl参数会在每个part文件之间插入换行避免前一个文件末尾没有换行导致两条记录粘在一起。取回结果后先看一眼head确认字段顺序符合预期再进入下一步验证。5. 气象数据跑MapReduce时我踩过的5个真坑这一章是血泪经验汇总。每条都按照“现象 → 原因 → 解决”来写很多问题不报错只是结果悄悄变错比直接报错更危险。5.1 坑一CSV字段里出现引号和逗号split导致列错位现象跑完统计后某站点的年降水量高达几十万毫米和常识差了好几个数量级。原因line.split(,, -1)只处理简单逗号分割。如果原始数据里某个字段被双引号包裹且引号内部含逗号简单split会把引号内逗号当成字段分隔符导致后续列全部错位。气象导出数据偶尔会出现这种“脏CSV”尤其是备注字段和降水信息混合时。解决不要用split硬切改用Commons CSV库。引入依赖后在Mapper里用CSVRecord解析import org.apache.commons.csv.CSVFormat; import org.apache.commons.csv.CSVParser; import org.apache.commons.csv.CSVRecord; String line value.toString(); CSVParser parser CSVParser.parse(line, CSVFormat.DEFAULT.builder().build()); CSVRecord record parser.getRecords().get(0); String stationId record.get(0); String date record.get(1);这样引号包裹的逗号会被解析为字段内容的一部分列顺序不再错乱。如果不想引第三方库至少要在预处理阶段统一清洗引号并在split时使用split(,, -1)保留尾部空字段防止数组越界。5.2 坑二缺测值参与聚合平均气温变成离谱负数现象某站点的年平均气温统计结果是 -498.2摄氏度max_temp和min_temp也全是 -9999 附近的值。原因气象数据用 -9999、9999、32700 表示观测缺测或仪器故障这些值没有物理意义。直接Double.parseDouble后参与sum和count聚合结果自然被拉偏。解决所有字段解析必须经过缺测过滤逻辑见3.2节的parseTemp。注意过滤要放在统计结构体累加之前而不是在Reducer里过滤否则无效值已经污染了count。课程设计里如果只做了“大于-100”的粗略过滤遇到-9999时依然会把平均值拉低几十度。5.3 坑三伪分布式下作业一直卡在Running进度纹丝不动现象作业提交后显示runningMap进度持续0%或容器反复重启后终止终端没有任何Java异常。原因最常见的是mapreduce.framework.name没有配置为yarn作业实际跑在local模式另一个高频原因是内存不足NodeManager把容器杀了但界面提示不明显。解决先在yarn-site.xml确认配置生效grep -A 3 framework.name $HADOOP_HOME/etc/hadoop/mapred-site.xml没有这个文件或配置缺失就补上再按4.4节的参数表调整内存。调完后重启YARN重新提交作业。如果日志里出现Java heap space单独调大mapreduce.map.java.opts出现Container killed则需要同时调大容器内存。5.4 坑四Combiner把平均值求了平均结果偏差0.3度现象同一份数据开启Combiner后年平均气温和手工Excel计算的结果相差零点几度但极值和降水总量完全正确。原因Combiner被错误地配置成输出“平均温度”Reducer拿到的是多个片段平均温度再次相加求平均。样本数不均匀时平均的平均不等于真实平均。解决用3.2节WeatherCombiner的写法中间值只传count sum max min precipSum precipCount最终平均值只在Reducer里计算。判断Combiner是否安全的标准只有一条输出的数据格式是否和Mapper输出保持同一语义。语义不一致Combiner就会改变最终结果。5.5 坑五输出目录没清理第二次提交作业直接失败现象第一次跑成功后不修改代码和参数第二次提交同一个作业立刻抛org.apache.hadoop.mapred.FileAlreadyExistsException: Output directory hdfs://.../result already exists。原因MapReduce框架为了安全禁止作业输出路径预先存在避免覆盖历史结果。解决在提交命令里显式清理或者每次用一个带时间戳的输出目录OUTPUT/weather/output/result_$(date %Y%m%d_%H%M%S)带时间戳的方式保留了每次运行的结果方便回查。用固定目录则需要在跑前执行hadoop fs -rm -r。两种方式我都用过课程设计阶段固定目录更顺手但答辩演示时如果切换历史结果时间戳目录反而更省事。6. 进阶三步验证让MapReduce结果站得住脚跑出part-r-00000只是第一步给答辩或报告用的数字必须能被复核。我的习惯是三步验证缺一不可。第一步总数对账。在Mapper里对每条有效输入做计数器累加运行结束后用hadoop job -history all查Counter和原始文件行数对比。如果Mapper输出比输入少说明有数据被过滤或解析失败这个差距本身也是数据质量报告的一部分。第二步极值回查。从结果里挑出全年最高温站点回到原始CSV用grep定位该站点该日期的记录确认数值不是缺测码也不是脏数据。极值比均值更容易暴露错误因为极端值往往只有少数几条人工检查成本低。第三步SQL复核。用Hive外部表直接读同一份HDFS原始数据按相同口径聚合对比MapReduce输出CREATE EXTERNAL TABLE weather_daily( station_id STRING, dt STRING, avg_t DOUBLE, max_t DOUBLE, min_t DOUBLE, precip DOUBLE, quality INT) ROW FORMAT DELIMITED FIELDS TERMINATED BY , LOCATION /weather/raw/year2020; SELECT substr(dt,1,4) AS year, station_id, avg(avg_t) AS avg_temp, max(max_t) AS max_temp, min(min_t) AS min_temp, sum(precip) AS total_precip FROM weather_daily WHERE avg_t NOT IN (-9999, 9999) AND max_t NOT IN (-9999, 9999) AND min_t NOT IN (-9999, 9999) GROUP BY substr(dt,1,4), station_id;MapReduce输出与SQL查询结果一致这条链路才算真正可信。我的习惯是永远保留这份验证脚本之后换数据集或改统计口径重跑一遍作业和复核SQL两边对上了再写进结论。这样处理出来的分析结果比贴一张控制台截图有说服力得多。希望帮到你。本文还有配套的精品资源点击获取