☰
Hadoop气象数据分析:MapReduce完整链路与踩坑实践
2026/9/26 2:28:17 网站建设 项目流程

简介:一套围绕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日期,YYYYMMDD
2avg_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/year=2020/month=01/site_54308.csv /weather/raw/year=2020/month=02/site_54308.csv /weather/raw/year=2020/month=12/site_54308.csv /weather/raw/year=2021/month=01/site_54308.csv

建目录并上传的常用命令如下:

hadoop fs -mkdir -p /weather/raw/year=2020/month=01 hadoop fs -put site_54308_202001.csv /weather/raw/year=2020/month=01/

后续提交作业时,只把/weather/raw/year=2020作为输入路径,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/year=2020 hadoop fs -put merged/2020_all.csv /weather/raw/year=2020/

上面的脚本用tail -n +2去掉除首个文件外的表头,避免合并后出现大量表头行。如果发现合并后同一站点同一天有重复记录,统计前必须做一次合并去重,否则均值、总量都会被放大。最简单的去重是按“站点编号+日期”去重:

awk -F, 'NR==1 || !seen[$1"_"$2]++' merged/2020_all.csv > merged/2020_dedup.csv

2.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 数据口径与列号约定:先定规则再写代码

写代码前必须把统计口径定死。我的口径是:只统计quality=0的记录;气温字段中-9999、9999视为缺测,整条记录跳过;降水字段缺测当天按0毫米计入,但不参与有效样本计数。所有输出的平均气温保留两位小数,输出格式使用Tab分隔,方便后续用Hive直接加载。

下面按列下标取字段,所以表头和字段顺序必须和2.1节的表格保持一致。如果实际数据列序不同,预处理阶段需要先做一次列重排,而不是在Mapper里到处改下标。

3.2 完整源码:Mapper、Combiner、Reducer与Driver

import 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 Mapper<Object, 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 Reducer<Text, Text, Text, Text> { @Override protected void reduce(Text key, Iterable<Text> 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(Iterable<Text> 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, Iterable<Text> 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.NaN,Mapper只对气温有效的行做统计,这是防止结果出现离谱负数的第一道闸门。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) / (10+2) = 11.67。如果Combiner直接把两个片段平均再平均,结果就是(10+20)/2=15,偏差明显。因此这里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.xml,Hadoop客户端依赖使用provided作用域,避免把整个Hadoop打进业务jar:

<dependencies> <dependency> <groupId>org.apache.hadoop</groupId> <artifactId>hadoop-client</artifactId> <version>3.3.4</version> <scope>provided</scope> </dependency> </dependencies> <build> <plugins> <plugin> <groupId>org.apache.maven.plugins</groupId> <artifactId>maven-jar-plugin</artifactId> <configuration> <archive> <manifest> <mainClass>WeatherAnalysisMR</mainClass> </manifest> </archive> </configuration> </plugin> </plugins> </build>

provided的含义是编译时可以使用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/year=2020 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 $OUTPUT

put -f用于本地文件已存在时强制覆盖,避免重复上传报错。清理输出目录放在提交前而不是代码里,是因为正式环境里作业不应该有删除HDFS目录的权限;课程作业图省事可以把删除逻辑放进Driver,但我不建议养成这个习惯。

命令行里的-D参数必须放在hadoop jar命令名之后、jar包之前,写法如下:

hadoop jar -D mapreduce.job.reduces=8 target/weather-analysis-1.0.jar \ weather.analysis.WeatherAnalysisMR /weather/raw/year=2020 /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 | head

stderr文件里通常是Java异常栈,syslog里能看到容器启动和GC信息。比较隐蔽的情况是作业已经跑完但终端没刷新,这时去YARN ResourceManager界面看Application状态最直接,或者用命令查:

yarn application -list -appStates FINISHED,FAILED,KILLED

如果拿到Application ID,可以拉取完整日志:

yarn logs -applicationId application_1710000000000_0001

4.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,进度纹丝不动

现象:作业提交后显示running,Map进度持续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/year=2020'; 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,两边对上了再写进结论。这样处理出来的分析结果,比贴一张控制台截图有说服力得多。希望帮到你。

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

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

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

立即咨询